Coverage for agentos/channels/adapters/dingtalk.py: 0%
56 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +0800
1"""
2AgentOS Channels — 钉钉适配器。
4Webhook 规范: https://open.dingtalk.com/document/orgapp/receive-message
6特性:
7 - JSON 报文解析
8 - 签名验证(timestamp + sign)
9 - access_token 管理
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 DingTalkAdapter(BaseChannelAdapter):
25 """钉钉适配器。"""
27 channel_type = ChannelType.DINGTALK
29 def __init__(self, config: ChannelConfig):
30 super().__init__(config)
31 self._token: str = ""
32 self._token_expires: float = 0
34 # ── Webhook ──
36 def verify_signature(self, raw_body: bytes, headers: dict) -> bool:
37 """验证钉钉签名(timestamp + sign SHA256)。"""
38 params = headers.get("x-ding-params", {})
39 timestamp = params.get("timestamp", "")
40 sign = params.get("sign", "")
41 if not sign:
42 return True
43 computed = self.hmac_sha256(timestamp, self.config.app_secret)
44 return computed == sign
46 def parse_webhook(self, raw_body: bytes, headers: dict) -> ChannelMessage | list[ChannelMessage]:
47 data = json.loads(raw_body.decode("utf-8"))
48 msg_type_str = data.get("msgtype", "text")
49 msg_type_map = {
50 "text": MessageType.TEXT, "image": MessageType.IMAGE,
51 "voice": MessageType.VOICE, "video": MessageType.VIDEO,
52 "file": MessageType.FILE, "link": MessageType.LINK,
53 }
54 msg_type = msg_type_map.get(msg_type_str, MessageType.TEXT)
56 content = ""
57 if msg_type_str == "text":
58 content = data.get("text", {}).get("content", "")
59 elif msg_type_str == "image":
60 content = "[图片]"
62 return ChannelMessage(
63 msg_id=data.get("msgId", data.get("msgid", "")),
64 channel=ChannelType.DINGTALK,
65 msg_type=msg_type,
66 content=content,
67 sender_id=data.get("senderStaffId", data.get("senderId", "")),
68 sender_name=data.get("senderNick", ""),
69 timestamp=float(data.get("createAt", time.time() * 1000)) / 1000,
70 conversation_id=data.get("sessionWebhook", ""),
71 reply_token="",
72 media_url=data.get("image", {}).get("picUrl", ""),
73 media_id=data.get("image", {}).get("mediaId", ""),
74 extra={
75 "robot_code": data.get("robotCode"),
76 "chatbot_user_id": data.get("chatbotUserId"),
77 "chat_id": data.get("chatId"),
78 "is_admin": data.get("isAdmin", False),
79 "conversation_type": data.get("conversationType"),
80 "at_users": data.get("atUsers", []),
81 },
82 )
84 def build_reply(self, msg: ChannelMessage, reply_text: str) -> str:
85 return json.dumps({"msgtype": "text", "text": {"content": reply_text}})
87 # ── 主动推送 ──
89 async def send_message(self, user_id: str, content: str, msg_type: str = "text") -> ReplyResult:
90 token = await self.get_access_token()
91 url = "https://api.dingtalk.com/v1.0/robot/oToMessages/batchSend"
92 payload = {
93 "robotCode": self.config.app_id,
94 "userIds": [user_id],
95 "msgKey": "sampleText",
96 "msgParam": json.dumps({"content": content}),
97 }
98 headers = {"x-acs-dingtalk-access-token": token, "Content-Type": "application/json"}
99 async with httpx.AsyncClient() as client:
100 resp = await client.post(url, json=payload, headers=headers, timeout=10)
101 return ReplyResult(success=resp.status_code == 200, msg_id=str(time.time()))
103 async def send_image(self, user_id: str, image_url: str) -> ReplyResult:
104 return await self.send_message(user_id, f"[图片] {image_url}")
106 async def send_file(self, user_id: str, file_url: str, filename: str) -> ReplyResult:
107 return await self.send_message(user_id, f"文件: {filename}\n{file_url}")
109 # ── Token ──
111 async def get_access_token(self) -> str:
112 if self._token and time.time() < self._token_expires - 300:
113 return self._token
114 url = "https://api.dingtalk.com/v1.0/oauth2/accessToken"
115 payload = {"appKey": self.config.app_id, "appSecret": self.config.app_secret}
116 async with httpx.AsyncClient() as client:
117 resp = await client.post(url, json=payload, timeout=10)
118 data = resp.json()
119 self._token = data["accessToken"]
120 self._token_expires = time.time() + data.get("expireIn", 7200)
121 return self._token