Coverage for agentos/channels/adapters/dingtalk.py: 0%
56 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-09 09:19 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-09 09:19 +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(
47 self, raw_body: bytes, headers: dict
48 ) -> ChannelMessage | list[ChannelMessage]:
49 data = json.loads(raw_body.decode("utf-8"))
50 msg_type_str = data.get("msgtype", "text")
51 msg_type_map = {
52 "text": MessageType.TEXT,
53 "image": MessageType.IMAGE,
54 "voice": MessageType.VOICE,
55 "video": MessageType.VIDEO,
56 "file": MessageType.FILE,
57 "link": MessageType.LINK,
58 }
59 msg_type = msg_type_map.get(msg_type_str, MessageType.TEXT)
61 content = ""
62 if msg_type_str == "text":
63 content = data.get("text", {}).get("content", "")
64 elif msg_type_str == "image":
65 content = "[图片]"
67 return ChannelMessage(
68 msg_id=data.get("msgId", data.get("msgid", "")),
69 channel=ChannelType.DINGTALK,
70 msg_type=msg_type,
71 content=content,
72 sender_id=data.get("senderStaffId", data.get("senderId", "")),
73 sender_name=data.get("senderNick", ""),
74 timestamp=float(data.get("createAt", time.time() * 1000)) / 1000,
75 conversation_id=data.get("sessionWebhook", ""),
76 reply_token="",
77 media_url=data.get("image", {}).get("picUrl", ""),
78 media_id=data.get("image", {}).get("mediaId", ""),
79 extra={
80 "robot_code": data.get("robotCode"),
81 "chatbot_user_id": data.get("chatbotUserId"),
82 "chat_id": data.get("chatId"),
83 "is_admin": data.get("isAdmin", False),
84 "conversation_type": data.get("conversationType"),
85 "at_users": data.get("atUsers", []),
86 },
87 )
89 def build_reply(self, msg: ChannelMessage, reply_text: str) -> str:
90 return json.dumps({"msgtype": "text", "text": {"content": reply_text}})
92 # ── 主动推送 ──
94 async def send_message(self, user_id: str, content: str, msg_type: str = "text") -> ReplyResult:
95 token = await self.get_access_token()
96 url = "https://api.dingtalk.com/v1.0/robot/oToMessages/batchSend"
97 payload = {
98 "robotCode": self.config.app_id,
99 "userIds": [user_id],
100 "msgKey": "sampleText",
101 "msgParam": json.dumps({"content": content}),
102 }
103 headers = {"x-acs-dingtalk-access-token": token, "Content-Type": "application/json"}
104 async with httpx.AsyncClient() as client:
105 resp = await client.post(url, json=payload, headers=headers, timeout=10)
106 return ReplyResult(success=resp.status_code == 200, msg_id=str(time.time()))
108 async def send_image(self, user_id: str, image_url: str) -> ReplyResult:
109 return await self.send_message(user_id, f"[图片] {image_url}")
111 async def send_file(self, user_id: str, file_url: str, filename: str) -> ReplyResult:
112 return await self.send_message(user_id, f"文件: {filename}\n{file_url}")
114 # ── Token ──
116 async def get_access_token(self) -> str:
117 if self._token and time.time() < self._token_expires - 300:
118 return self._token
119 url = "https://api.dingtalk.com/v1.0/oauth2/accessToken"
120 payload = {"appKey": self.config.app_id, "appSecret": self.config.app_secret}
121 async with httpx.AsyncClient() as client:
122 resp = await client.post(url, json=payload, timeout=10)
123 data = resp.json()
124 self._token = data["accessToken"]
125 self._token_expires = time.time() + data.get("expireIn", 7200)
126 return self._token