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

79 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-08 01:44 +0800

1""" 

2AgentOS Channels — 微信公众号适配器。 

3 

4Webhook 规范: https://developers.weixin.qq.com/doc/offiaccount/Message_Management/Receiving_standard_messages.html 

5 

6特性: 

7 - XML 报文解析 

8 - SHA1 签名验证 

9 - 被动回复(在 webhook 响应中同步回复) 

10 - access_token 管理 + 自动续期 

11 - 主动推送(客服消息接口) 

12""" 

13 

14from __future__ import annotations 

15 

16import time 

17import xml.etree.ElementTree as ET 

18 

19import httpx 

20 

21from agentos.channels.base import ( 

22 BaseChannelAdapter, 

23 ChannelConfig, 

24 ReplyResult, 

25) 

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

27 

28 

29class WeChatAdapter(BaseChannelAdapter): 

30 """微信公众号适配器。""" 

31 

32 channel_type = ChannelType.WECHAT_MP 

33 

34 def __init__(self, config: ChannelConfig): 

35 super().__init__(config) 

36 self._token: str = "" 

37 self._token_expires: float = 0 

38 

39 # ── Webhook ── 

40 

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

42 """验证微信签名(SHA1)。""" 

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

44 if not params: 

45 return True # 无签名时放行(开发模式) 

46 signature = params.get("signature", "") 

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

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

49 expected = self.make_signature(self.config.verify_token, timestamp, nonce) 

50 return signature == expected 

51 

52 def parse_webhook( 

53 self, raw_body: bytes, headers: dict 

54 ) -> ChannelMessage | list[ChannelMessage]: 

55 """解析微信 XML 报文。""" 

56 root = ET.fromstring(raw_body.decode("utf-8")) 

57 

58 msg_type_str = self._xml_text(root, "MsgType") or "text" 

59 msg_type_map = { 

60 "text": MessageType.TEXT, 

61 "image": MessageType.IMAGE, 

62 "voice": MessageType.VOICE, 

63 "video": MessageType.VIDEO, 

64 "location": MessageType.LOCATION, 

65 "link": MessageType.LINK, 

66 "event": MessageType.EVENT, 

67 } 

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

69 

70 content = "" 

71 if msg_type_str == "text": 

72 content = self._xml_text(root, "Content") or "" 

73 elif msg_type_str == "image": 

74 content = "[图片]" 

75 elif msg_type_str == "voice": 

76 content = self._xml_text(root, "Recognition") or "[语音]" 

77 

78 return ChannelMessage( 

79 msg_id=self._xml_text(root, "MsgId") or "", 

80 channel=ChannelType.WECHAT_MP, 

81 msg_type=msg_type, 

82 content=content, 

83 sender_id=self._xml_text(root, "FromUserName") or "", 

84 sender_name="", 

85 timestamp=float(self._xml_text(root, "CreateTime") or time.time()), 

86 conversation_id=self._xml_text(root, "FromUserName") or "", 

87 reply_token="", 

88 media_url=self._xml_text(root, "PicUrl") or self._xml_text(root, "MediaId") or "", 

89 media_id=self._xml_text(root, "MediaId") or "", 

90 extra={ 

91 "to_user": self._xml_text(root, "ToUserName"), 

92 "msg_type_raw": msg_type_str, 

93 "event": self._xml_text(root, "Event"), 

94 "event_key": self._xml_text(root, "EventKey"), 

95 }, 

96 ) 

97 

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

99 """构建微信被动回复 XML。""" 

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

101 from_user = msg.sender_id 

102 create_time = int(time.time()) 

103 return ( 

104 "<xml>" 

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

106 f"<FromUserName><![CDATA[{from_user}]]></FromUserName>" 

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

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

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

110 "</xml>" 

111 ) 

112 

113 # ── 主动推送 ── 

114 

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

116 """发送客服消息。""" 

117 token = await self.get_access_token() 

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

119 payload = { 

120 "touser": user_id, 

121 "msgtype": "text", 

122 "text": {"content": content}, 

123 } 

124 async with httpx.AsyncClient() as client: 

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

126 data = resp.json() 

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

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

129 return ReplyResult( 

130 success=False, error=f"wechat error {data.get('errcode')}: {data.get('errmsg')}" 

131 ) 

132 

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

134 token = await self.get_access_token() 

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

136 payload = { 

137 "touser": user_id, 

138 "msgtype": "image", 

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

140 } 

141 async with httpx.AsyncClient() as client: 

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

143 data = resp.json() 

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

145 return ReplyResult(success=True) 

146 return ReplyResult(success=False, error=str(data)) 

147 

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

149 await self.get_access_token() 

150 # 先上传临时素材 

151 async with httpx.AsyncClient(): 

152 # 简化实现:发文本链接 

153 return await self.send_message(user_id, f"文件: {filename}\n{file_url}") 

154 

155 # ── Token ── 

156 

157 async def get_access_token(self) -> str: 

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

159 return self._token 

160 

161 url = ( 

162 "https://api.weixin.qq.com/cgi-bin/token" 

163 f"?grant_type=client_credential" 

164 f"&appid={self.config.app_id}" 

165 f"&secret={self.config.app_secret}" 

166 ) 

167 async with httpx.AsyncClient() as client: 

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

169 data = resp.json() 

170 if "access_token" in data: 

171 self._token = data["access_token"] 

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

173 return self._token 

174 raise RuntimeError(f"wechat token error: {data}") 

175 

176 # ── Helpers ── 

177 

178 @staticmethod 

179 def _xml_text(element: ET.Element, tag: str) -> str | None: 

180 child = element.find(tag) 

181 return child.text if child is not None else None