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
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +0800
1"""
2AgentOS Channels — 企业微信适配器。
4Webhook 规范: https://developer.work.weixin.qq.com/document/path/90238
6特性:
7 - XML/JSON 双报文解析
8 - SHA1 签名验证
9 - 被动回复 + 主动群机器人 webhook 推送
10 - access_token 自动续期
11"""
13from __future__ import annotations
15import json
16import time
17import xml.etree.ElementTree as ET
19import httpx
21from agentos.channels.base import BaseChannelAdapter, ChannelConfig, ReplyResult
22from agentos.channels.message import ChannelMessage, ChannelType, MessageType
25class WeComAdapter(BaseChannelAdapter):
26 """企业微信适配器。"""
28 channel_type = ChannelType.WECOM
30 def __init__(self, config: ChannelConfig):
31 super().__init__(config)
32 self._token: str = ""
33 self._token_expires: float = 0
35 # ── Webhook ──
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
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 = "[图片]"
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 )
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 )
97 # ── 主动推送(群机器人 webhook 或应用消息)──
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)
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')}")
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)
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)
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
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
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}