Coverage for agentos/channels/adapters/dingtalk.py: 0%

56 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-08 20:40 +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( 

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) 

60 

61 content = "" 

62 if msg_type_str == "text": 

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

64 elif msg_type_str == "image": 

65 content = "[图片]" 

66 

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 ) 

88 

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

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

91 

92 # ── 主动推送 ── 

93 

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

107 

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

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

110 

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

113 

114 # ── Token ── 

115 

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