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

81 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-06 10:59 +0800

1""" 

2AgentOS Channels — 企业微信适配器。 

3 

4Webhook 规范: https://developer.work.weixin.qq.com/document/path/90238 

5 

6特性: 

7 - XML/JSON 双报文解析 

8 - SHA1 签名验证 

9 - 被动回复 + 主动群机器人 webhook 推送 

10 - access_token 自动续期 

11""" 

12 

13from __future__ import annotations 

14 

15import json 

16import time 

17import xml.etree.ElementTree as ET 

18 

19import httpx 

20 

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

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

23 

24 

25class WeComAdapter(BaseChannelAdapter): 

26 """企业微信适配器。""" 

27 

28 channel_type = ChannelType.WECOM 

29 

30 def __init__(self, config: ChannelConfig): 

31 super().__init__(config) 

32 self._token: str = "" 

33 self._token_expires: float = 0 

34 

35 # ── Webhook ── 

36 

37 def verify_signature(self, raw_body: bytes, headers: dict) -> bool: 

38 """验证企微签名。""" 

39 params = headers.get("x-wx-params", {}) 

40 msg_signature = params.get("msg_signature", "") 

41 timestamp = str(params.get("timestamp", "")) 

42 nonce = str(params.get("nonce", "")) 

43 signature = self.make_signature(self.config.verify_token, timestamp, nonce, "") 

44 return msg_signature == signature 

45 

46 def parse_webhook(self, raw_body: bytes, headers: dict) -> ChannelMessage | list[ChannelMessage]: 

47 text = raw_body.decode("utf-8") 

48 data = json.loads(text) if text.strip().startswith("{") else self._parse_xml(text) 

49 msg_type_str = data.get("MsgType", data.get("msgtype", "text")) 

50 msg_type_map = { 

51 "text": MessageType.TEXT, "image": MessageType.IMAGE, 

52 "voice": MessageType.VOICE, "video": MessageType.VIDEO, 

53 "file": MessageType.FILE, "event": MessageType.EVENT, 

54 } 

55 msg_type = msg_type_map.get(msg_type_str, MessageType.TEXT) 

56 content = "" 

57 if msg_type_str == "text": 

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

59 elif msg_type_str == "image": 

60 content = "[图片]" 

61 

62 return ChannelMessage( 

63 msg_id=data.get("MsgId", "") or "", 

64 channel=ChannelType.WECOM, 

65 msg_type=msg_type, 

66 content=content, 

67 sender_id=data.get("FromUserName", data.get("UserID", "")), 

68 sender_name=data.get("Name", ""), 

69 timestamp=float(data.get("CreateTime", time.time())), 

70 conversation_id=data.get("ChatId", data.get("FromUserName", "")), 

71 media_url=data.get("PicUrl", ""), 

72 media_id=data.get("MediaId", ""), 

73 extra={ 

74 "to_user": data.get("ToUserName"), 

75 "agent_id": data.get("AgentID"), 

76 "msg_type_raw": msg_type_str, 

77 "webhook_url": data.get("WebhookUrl", ""), 

78 "chat_type": data.get("ChatType", "single"), 

79 }, 

80 ) 

81 

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

83 if msg.extra.get("webhook_url"): 

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

85 to_user = msg.extra.get("to_user", msg.sender_id) 

86 create_time = int(time.time()) 

87 return ( 

88 "<xml>" 

89 f"<ToUserName><![CDATA[{to_user}]]></ToUserName>" 

90 f"<FromUserName><![CDATA[{msg.sender_id}]]></FromUserName>" 

91 f"<CreateTime>{create_time}</CreateTime>" 

92 "<MsgType><![CDATA[text]]></MsgType>" 

93 f"<Content><![CDATA[{reply_text}]]></Content>" 

94 "</xml>" 

95 ) 

96 

97 # ── 主动推送(群机器人 webhook 或应用消息)── 

98 

99 async def send_message(self, user_id: str, content: str, msg_type: str = "text") -> ReplyResult: 

100 # 如果有 webhook_url 则走群机器人推送 

101 webhook_url = self.config.extra.get("webhook_url", "") 

102 if webhook_url: 

103 async with httpx.AsyncClient() as client: 

104 resp = await client.post(webhook_url, json={ 

105 "msgtype": "text", 

106 "text": {"content": content}, 

107 }, timeout=10) 

108 return ReplyResult(success=resp.status_code == 200) 

109 

110 token = await self.get_access_token() 

111 url = f"https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token={token}" 

112 payload = { 

113 "touser": user_id, 

114 "msgtype": "text", 

115 "agentid": int(self.config.agent_id or 0), 

116 "text": {"content": content}, 

117 } 

118 async with httpx.AsyncClient() as client: 

119 resp = await client.post(url, json=payload, timeout=10) 

120 data = resp.json() 

121 if data.get("errcode") == 0: 

122 return ReplyResult(success=True, msg_id=data.get("msgid", "")) 

123 return ReplyResult(success=False, error=f"wecom error {data.get('errcode')}: {data.get('errmsg')}") 

124 

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

126 token = await self.get_access_token() 

127 url = f"https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token={token}" 

128 payload = { 

129 "touser": user_id, 

130 "msgtype": "image", 

131 "agentid": int(self.config.agent_id or 0), 

132 "image": {"media_id": image_url}, 

133 } 

134 async with httpx.AsyncClient() as client: 

135 resp = await client.post(url, json=payload, timeout=10) 

136 return ReplyResult(success=resp.json().get("errcode") == 0) 

137 

138 async def send_file(self, user_id: str, file_url: str, filename: str) -> ReplyResult: 

139 token = await self.get_access_token() 

140 url = f"https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token={token}" 

141 payload = { 

142 "touser": user_id, 

143 "msgtype": "file", 

144 "agentid": int(self.config.agent_id or 0), 

145 "file": {"media_id": file_url}, 

146 } 

147 async with httpx.AsyncClient() as client: 

148 resp = await client.post(url, json=payload, timeout=10) 

149 return ReplyResult(success=resp.json().get("errcode") == 0) 

150 

151 # ── Token ── 

152 

153 async def get_access_token(self) -> str: 

154 if self._token and time.time() < self._token_expires - 300: 

155 return self._token 

156 url = ( 

157 "https://qyapi.weixin.qq.com/cgi-bin/gettoken" 

158 f"?corpid={self.config.corp_id}" 

159 f"&corpsecret={self.config.app_secret}" 

160 ) 

161 async with httpx.AsyncClient() as client: 

162 resp = await client.get(url, timeout=10) 

163 data = resp.json() 

164 self._token = data["access_token"] 

165 self._token_expires = time.time() + data.get("expires_in", 7200) 

166 return self._token 

167 

168 @staticmethod 

169 def _parse_xml(xml_str: str) -> dict: 

170 root = ET.fromstring(xml_str) 

171 return {child.tag: child.text for child in root}