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
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-08 01:44 +0800
1"""
2AgentOS Channels — 微信公众号适配器。
4Webhook 规范: https://developers.weixin.qq.com/doc/offiaccount/Message_Management/Receiving_standard_messages.html
6特性:
7 - XML 报文解析
8 - SHA1 签名验证
9 - 被动回复(在 webhook 响应中同步回复)
10 - access_token 管理 + 自动续期
11 - 主动推送(客服消息接口)
12"""
14from __future__ import annotations
16import time
17import xml.etree.ElementTree as ET
19import httpx
21from agentos.channels.base import (
22 BaseChannelAdapter,
23 ChannelConfig,
24 ReplyResult,
25)
26from agentos.channels.message import ChannelMessage, ChannelType, MessageType
29class WeChatAdapter(BaseChannelAdapter):
30 """微信公众号适配器。"""
32 channel_type = ChannelType.WECHAT_MP
34 def __init__(self, config: ChannelConfig):
35 super().__init__(config)
36 self._token: str = ""
37 self._token_expires: float = 0
39 # ── Webhook ──
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
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"))
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)
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 "[语音]"
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 )
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 )
113 # ── 主动推送 ──
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 )
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))
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}")
155 # ── Token ──
157 async def get_access_token(self) -> str:
158 if self._token and time.time() < self._token_expires - 300:
159 return self._token
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}")
176 # ── Helpers ──
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