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

1""" 

2AgentOS Channels — 钉钉适配器。 

3 

4Webhook 规范: https://open.dingtalk.com/document/orgapp/receive-message 

5 

6特性: 

7 - JSON 报文解析 

8 - 签名验证(timestamp + sign) 

9 - access_token 管理 

10 - 主动推送(工作通知 + 群机器人) 

11""" 

12 

13from __future__ import annotations 

14 

15import json 

16import time 

17 

18import httpx 

19 

20from agentos.channels.base import BaseChannelAdapter, ChannelConfig, ReplyResult 

21from agentos.channels.message import ChannelMessage, ChannelType, MessageType 

22 

23 

24class DingTalkAdapter(BaseChannelAdapter): 

25 """钉钉适配器。""" 

26 

27 channel_type = ChannelType.DINGTALK 

28 

29 def __init__(self, config: ChannelConfig): 

30 super().__init__(config) 

31 self._token: str = "" 

32 self._token_expires: float = 0 

33 

34 # ── Webhook ── 

35 

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 

45 

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) 

55 

56 content = "" 

57 if msg_type_str == "text": 

58 content = data.get("text", {}).get("content", "") 

59 elif msg_type_str == "image": 

60 content = "[图片]" 

61 

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 ) 

83 

84 def build_reply(self, msg: ChannelMessage, reply_text: str) -> str: 

85 return json.dumps({"msgtype": "text", "text": {"content": reply_text}}) 

86 

87 # ── 主动推送 ── 

88 

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())) 

102 

103 async def send_image(self, user_id: str, image_url: str) -> ReplyResult: 

104 return await self.send_message(user_id, f"[图片] {image_url}") 

105 

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}") 

108 

109 # ── Token ── 

110 

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