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
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +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
18from typing import Optional
20import httpx
22from agentos.channels.base import (
23 BaseChannelAdapter, ChannelConfig, ReplyResult,
24)
25from agentos.channels.message import ChannelMessage, ChannelType, MessageType
28class WeChatAdapter(BaseChannelAdapter):
29 """微信公众号适配器。"""
31 channel_type = ChannelType.WECHAT_MP
33 def __init__(self, config: ChannelConfig):
34 super().__init__(config)
35 self._token: str = ""
36 self._token_expires: float = 0
38 # ── Webhook ──
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
51 def parse_webhook(self, raw_body: bytes, headers: dict) -> ChannelMessage | list[ChannelMessage]:
52 """解析微信 XML 报文。"""
53 root = ET.fromstring(raw_body.decode("utf-8"))
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)
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 "[语音]"
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 )
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 )
110 # ── 主动推送 ──
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')}")
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))
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}")
151 # ── Token ──
153 async def get_access_token(self) -> str:
154 if self._token and time.time() < self._token_expires - 300:
155 return self._token
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}")
172 # ── Helpers ──
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