Coverage for agentos/channels/adapters/wechat.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://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 

18from typing import Optional 

19 

20import httpx 

21 

22from agentos.channels.base import ( 

23 BaseChannelAdapter, ChannelConfig, ReplyResult, 

24) 

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

26 

27 

28class WeChatAdapter(BaseChannelAdapter): 

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

30 

31 channel_type = ChannelType.WECHAT_MP 

32 

33 def __init__(self, config: ChannelConfig): 

34 super().__init__(config) 

35 self._token: str = "" 

36 self._token_expires: float = 0 

37 

38 # ── Webhook ── 

39 

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

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

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

43 if not params: 

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

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

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

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

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

49 return signature == expected 

50 

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

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

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

54 

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

56 msg_type_map = { 

57 "text": MessageType.TEXT, 

58 "image": MessageType.IMAGE, 

59 "voice": MessageType.VOICE, 

60 "video": MessageType.VIDEO, 

61 "location": MessageType.LOCATION, 

62 "link": MessageType.LINK, 

63 "event": MessageType.EVENT, 

64 } 

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

66 

67 content = "" 

68 if msg_type_str == "text": 

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

70 elif msg_type_str == "image": 

71 content = "[图片]" 

72 elif msg_type_str == "voice": 

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

74 

75 return ChannelMessage( 

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

77 channel=ChannelType.WECHAT_MP, 

78 msg_type=msg_type, 

79 content=content, 

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

81 sender_name="", 

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

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

84 reply_token="", 

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

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

87 extra={ 

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

89 "msg_type_raw": msg_type_str, 

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

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

92 }, 

93 ) 

94 

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

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

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

98 from_user = msg.sender_id 

99 create_time = int(time.time()) 

100 return ( 

101 "<xml>" 

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

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

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

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

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

107 "</xml>" 

108 ) 

109 

110 # ── 主动推送 ── 

111 

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

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

114 token = await self.get_access_token() 

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

116 payload = { 

117 "touser": user_id, 

118 "msgtype": "text", 

119 "text": {"content": content}, 

120 } 

121 async with httpx.AsyncClient() as client: 

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

123 data = resp.json() 

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

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

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

127 

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

129 token = await self.get_access_token() 

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

131 payload = { 

132 "touser": user_id, 

133 "msgtype": "image", 

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

135 } 

136 async with httpx.AsyncClient() as client: 

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

138 data = resp.json() 

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

140 return ReplyResult(success=True) 

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

142 

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

144 token = await self.get_access_token() 

145 # 先上传临时素材 

146 upload_url = f"https://api.weixin.qq.com/cgi-bin/media/upload?access_token={token}&type=file" 

147 async with httpx.AsyncClient() as client: 

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

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

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 

157 url = ( 

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

159 f"?grant_type=client_credential" 

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

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

162 ) 

163 async with httpx.AsyncClient() as client: 

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

165 data = resp.json() 

166 if "access_token" in data: 

167 self._token = data["access_token"] 

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

169 return self._token 

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

171 

172 # ── Helpers ── 

173 

174 @staticmethod 

175 def _xml_text(element: ET.Element, tag: str) -> Optional[str]: 

176 child = element.find(tag) 

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