Coverage for agentos/channels/adapters/feishu.py: 0%
83 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-08 13:14 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-08 13:14 +0800
1"""
2AgentOS Channels — 飞书适配器。
4Webhook 规范: https://open.feishu.cn/document/server-docs/im-v1/message-content-description
6特性:
7 - JSON 报文解析
8 - 应用 Token + tenant access token 双 token 管理
9 - 卡片消息支持
10 - 消息回复(被动 + 主动)
11"""
13from __future__ import annotations
15import json
16import time
18import httpx
20from agentos.channels.base import BaseChannelAdapter, ChannelConfig, ReplyResult
21from agentos.channels.message import ChannelMessage, ChannelType, MessageType
24class FeishuAdapter(BaseChannelAdapter):
25 """飞书适配器。"""
27 channel_type = ChannelType.FEISHU
29 def __init__(self, config: ChannelConfig):
30 super().__init__(config)
31 self._app_token: str = ""
32 self._tenant_token: str = ""
33 self._token_expires: float = 0
35 # ── Webhook ──
37 def verify_signature(self, raw_body: bytes, headers: dict) -> bool:
38 """验证飞书事件订阅签名。
40 签名算法: Base64Encode(SHA256(timestamp + nonce + encrypt_key))
41 文档: https://open.feishu.cn/document/server-docs/event-subscription-guide/event-subscription-configure-/encrypt-key-encryption-configuration-
42 """
43 import base64
44 import hashlib
46 timestamp = headers.get("X-Lark-Request-Timestamp", "")
47 nonce = headers.get("X-Lark-Request-Nonce", "")
48 signature = headers.get("X-Lark-Signature", "")
50 encrypt_key = self.config.encoding_aes_key or self.config.verify_token
51 if not all([timestamp, nonce, signature, encrypt_key]):
52 return False
54 raw = f"{timestamp}{nonce}{encrypt_key}"
55 computed = base64.b64encode(hashlib.sha256(raw.encode()).digest()).decode()
56 return signature == computed
58 def parse_webhook(
59 self, raw_body: bytes, headers: dict
60 ) -> ChannelMessage | list[ChannelMessage]:
61 data = json.loads(raw_body.decode("utf-8"))
62 # 飞书事件格式: {"schema": "2.0", "header": {...}, "event": {...}}
63 event = data.get("event", data)
64 header = data.get("header", {})
66 # 处理 URL 验证
67 if data.get("type") == "url_verification":
68 return ChannelMessage(
69 msg_id="url_verify",
70 channel=ChannelType.FEISHU,
71 msg_type=MessageType.EVENT,
72 content=data.get("challenge", ""),
73 reply_token=data.get("token", ""),
74 extra={"is_challenge": True, "challenge": data.get("challenge", "")},
75 )
77 msg_type_str = event.get("message", {}).get("message_type", "text")
78 msg_type_map = {
79 "text": MessageType.TEXT,
80 "image": MessageType.IMAGE,
81 "audio": MessageType.VOICE,
82 "media": MessageType.FILE,
83 "file": MessageType.FILE,
84 "post": MessageType.TEXT,
85 }
86 msg_type = msg_type_map.get(msg_type_str, MessageType.TEXT)
88 message = event.get("message", {})
89 content = ""
90 if msg_type_str == "text":
91 content = json.loads(message.get("content", "{}")).get("text", "")
92 elif msg_type_str == "post":
93 content = str(message.get("content", ""))[:200]
95 sender = event.get("sender", {})
96 sender_id = sender.get("sender_id", {}).get("open_id", "")
98 return ChannelMessage(
99 msg_id=header.get("event_id", event.get("message", {}).get("message_id", "")),
100 channel=ChannelType.FEISHU,
101 msg_type=msg_type,
102 content=content,
103 sender_id=sender_id,
104 sender_name="",
105 timestamp=float(header.get("create_time", str(int(time.time() * 1000)))) / 1000,
106 conversation_id=event.get("message", {}).get("chat_id", ""),
107 reply_token=event.get("message", {}).get("message_id", ""),
108 media_url=message.get("image_key", ""),
109 extra={
110 "tenant_key": header.get("tenant_key"),
111 "event_type": header.get("event_type"),
112 "chat_type": event.get("message", {}).get("chat_type", "p2p"),
113 "root_id": event.get("message", {}).get("root_id"),
114 "parent_id": event.get("message", {}).get("parent_id"),
115 },
116 )
118 def build_reply(self, msg: ChannelMessage, reply_text: str) -> str:
119 return json.dumps(
120 {
121 "msg_type": "text",
122 "content": json.dumps({"text": reply_text}),
123 }
124 )
126 # ── 主动推送 ──
128 async def send_message(self, user_id: str, content: str, msg_type: str = "text") -> ReplyResult:
129 token = await self.get_access_token()
130 url = "https://open.feishu.cn/open-apis/im/v1/messages"
131 payload = {
132 "receive_id": user_id,
133 "msg_type": "text",
134 "content": json.dumps({"text": content}),
135 }
136 headers = {"Authorization": f"Bearer {token}", "Content-Type": "application/json"}
137 async with httpx.AsyncClient() as client:
138 resp = await client.post(
139 url,
140 params={"receive_id_type": "open_id"},
141 json=payload,
142 headers=headers,
143 timeout=10,
144 )
145 data = resp.json()
146 if data.get("code") == 0:
147 return ReplyResult(success=True, msg_id=data.get("data", {}).get("message_id", ""))
148 return ReplyResult(
149 success=False, error=f"feishu error {data.get('code')}: {data.get('msg')}"
150 )
152 async def send_image(self, user_id: str, image_url: str) -> ReplyResult:
153 token = await self.get_access_token()
154 url = "https://open.feishu.cn/open-apis/im/v1/messages"
155 payload = {
156 "receive_id": user_id,
157 "msg_type": "image",
158 "content": json.dumps({"image_key": image_url}),
159 }
160 headers = {"Authorization": f"Bearer {token}"}
161 async with httpx.AsyncClient() as client:
162 resp = await client.post(
163 url,
164 params={"receive_id_type": "open_id"},
165 json=payload,
166 headers=headers,
167 timeout=10,
168 )
169 return ReplyResult(success=resp.json().get("code") == 0)
171 async def send_file(self, user_id: str, file_url: str, filename: str) -> ReplyResult:
172 token = await self.get_access_token()
173 url = "https://open.feishu.cn/open-apis/im/v1/messages"
174 payload = {
175 "receive_id": user_id,
176 "msg_type": "file",
177 "content": json.dumps({"file_key": file_url}),
178 }
179 headers = {"Authorization": f"Bearer {token}"}
180 async with httpx.AsyncClient() as client:
181 resp = await client.post(
182 url,
183 params={"receive_id_type": "open_id"},
184 json=payload,
185 headers=headers,
186 timeout=10,
187 )
188 return ReplyResult(success=resp.json().get("code") == 0)
190 # ── Token ──
192 async def get_access_token(self) -> str:
193 if self._tenant_token and time.time() < self._token_expires - 300:
194 return self._tenant_token
195 url = "https://open.feishu.cn/open-apis/auth/v3/tenant_access_token/internal"
196 payload = {"app_id": self.config.app_id, "app_secret": self.config.app_secret}
197 async with httpx.AsyncClient() as client:
198 resp = await client.post(url, json=payload, timeout=10)
199 data = resp.json()
200 self._tenant_token = data["tenant_access_token"]
201 self._token_expires = time.time() + data.get("expire", 7200)
202 return self._tenant_token