diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/doctor_adapter.py b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/doctor_adapter.py index 6565f27b..fb809c70 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/doctor_adapter.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/doctor_adapter.py @@ -2,9 +2,9 @@ 实现 ``DoctorAdapter`` Protocol,为微信公众号渠道提供配置诊断能力 (FR-17),辅助管理员定位凭证失效、IP 白名单未配置、access_token 不可用、 -EncodingAESKey 错误与 Webhook 未配置等常见问题。 +EncodingAESKey 错误等常见问题。 -诊断检查项(``getDiagnosticItems`` 声明 5 项): +诊断检查项(``getDiagnosticItems`` 声明 4 项): - ``app_id_app_secret_valid``(CRITICAL,不可自动修复):校验 app_id 格式 与 app_secret 非空。 - ``ip_whitelist_configured``(CRITICAL,不可自动修复):调用 @@ -13,8 +13,12 @@ EncodingAESKey 错误与 Webhook 未配置等常见问题。 ``refresh_access_token`` 验证凭证有效性。 - ``encoding_aes_key_valid``(ERROR,不可自动修复):校验 ``encoding_aes_key`` 为 43 位(仅 compatible/safe 模式)。 -- ``webhook_url_configured``(ERROR,不可自动修复):校验 ``webhook_url`` - 非空。 + +注:wechat_mp 的 Webhook URL 由系统按 +``WEBHOOK_URL_TEMPLATE`` 生成(非用户配置),wizard 的 +``webhook_url_display`` 步骤已展示 URL 供用户填入微信公众平台,故不设置 +``webhook_url_configured`` 诊断项(服务端无法验证用户是否已在微信后台 +完成配置)。 设计要点: - 适配器无状态(stateless,INV-5),仅持有 DI 注入的 ``WeChatMpClient`` / @@ -55,7 +59,6 @@ _ITEM_APP_ID_APP_SECRET_VALID = "app_id_app_secret_valid" _ITEM_IP_WHITELIST_CONFIGURED = "ip_whitelist_configured" _ITEM_ACCESS_TOKEN_AVAILABLE = "access_token_available" _ITEM_ENCODING_AES_KEY_VALID = "encoding_aes_key_valid" -_ITEM_WEBHOOK_URL_CONFIGURED = "webhook_url_configured" # EncodingAESKey 固定长度 _ENCODING_AES_KEY_LENGTH = 43 @@ -93,18 +96,9 @@ _DIAGNOSTIC_ITEMS: tuple[DiagnosticCheck, ...] = ( description="校验 encoding_aes_key 为 43 位(仅 compatible/safe 模式)", auto_repairable=False, ), - DiagnosticCheck( - check_id=_ITEM_WEBHOOK_URL_CONFIGURED, - name="Webhook URL 配置", - severity=DiagnosticSeverity.ERROR, - description="校验 webhook_url 非空", - auto_repairable=False, - ), ) -_DIAGNOSTIC_ITEM_MAP: dict[str, DiagnosticCheck] = { - item.check_id: item for item in _DIAGNOSTIC_ITEMS -} +_DIAGNOSTIC_ITEM_MAP: dict[str, DiagnosticCheck] = {item.check_id: item for item in _DIAGNOSTIC_ITEMS} class WeChatMpDoctorAdapter: @@ -234,8 +228,6 @@ class WeChatMpDoctorAdapter: return await self._run_access_token_available(account_id, check) if item_id == _ITEM_ENCODING_AES_KEY_VALID: return await self._run_encoding_aes_key_valid(account_id, check) - if item_id == _ITEM_WEBHOOK_URL_CONFIGURED: - return await self._run_webhook_url_configured(account_id, check) raise ValidationError( field="item_id", message=f"unknown_check: {item_id}", @@ -381,10 +373,7 @@ class WeChatMpDoctorAdapter: check_id=check.check_id, passed=True, severity=check.severity, - message=( - f"encrypt_mode is '{encrypt_mode or 'plain'}', " - "encoding_aes_key not required" - ), + message=(f"encrypt_mode is '{encrypt_mode or 'plain'}', encoding_aes_key not required"), auto_repairable=check.auto_repairable, ) encoding_aes_key = await self._read_config(account_id, "encoding_aes_key") @@ -394,8 +383,7 @@ class WeChatMpDoctorAdapter: passed=False, severity=check.severity, message=( - f"encoding_aes_key must be {_ENCODING_AES_KEY_LENGTH} characters, " - f"got {len(encoding_aes_key)}" + f"encoding_aes_key must be {_ENCODING_AES_KEY_LENGTH} characters, got {len(encoding_aes_key)}" ), auto_repairable=check.auto_repairable, ) @@ -407,28 +395,6 @@ class WeChatMpDoctorAdapter: auto_repairable=check.auto_repairable, ) - async def _run_webhook_url_configured( - self, - account_id: str, - check: DiagnosticCheck, - ) -> DiagnosticResult: - webhook_url = await self._read_config(account_id, "webhook_url") - if not webhook_url: - return DiagnosticResult( - check_id=check.check_id, - passed=False, - severity=check.severity, - message="webhook_url is not configured", - auto_repairable=check.auto_repairable, - ) - return DiagnosticResult( - check_id=check.check_id, - passed=True, - severity=check.severity, - message="webhook_url is configured", - auto_repairable=check.auto_repairable, - ) - # ------------------------------------------------------------------ # 配置读取辅助 # ------------------------------------------------------------------ @@ -436,9 +402,14 @@ class WeChatMpDoctorAdapter: async def _read_config(self, account_id: str, key: str) -> str: from yuxi.channels.contract.dtos.config import ConfigScope + # app_id / app_secret / token / encoding_aes_key 均为 account-level + # 配置(manifest config_schema scope=account),使用 ConfigScope.ACCOUNT + # + target=account_id 读取,与 wecom 适配器保持一致。 try: cv = await self._config.get( - key, scope=ConfigScope.CHANNEL, target=account_id, + key, + scope=ConfigScope.ACCOUNT, + target=account_id, ) except Exception: return "" diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/inbound_adapter.py b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/inbound_adapter.py index 6cb643fb..1cb5c821 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/inbound_adapter.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/inbound_adapter.py @@ -44,11 +44,12 @@ from yuxi.channels.contract.ports.driven.logger_port import LoggerPort from .. import crypto from .._constants import ( + CHANNEL_TARGET, ENCRYPT_MODE_COMPATIBLE, ENCRYPT_MODE_PLAINTEXT, ENCRYPT_MODE_SAFE, - WECHAT_MP_DEP, WEBHOOK_SUCCESS_RESPONSE, + WECHAT_MP_DEP, ) from ..wechat_mp_client import WeChatMpClient @@ -91,15 +92,17 @@ class WeChatMpInboundAdapter: @pre - ``raw_event.payload`` 为微信 Webhook POST 的 XML 字符串 (明文模式直接原文;兼容/安全模式为含 ```` 的密文 XML) - - ``raw_event.headers`` 含 ``encrypt_mode``(可选,默认 plaintext) + 或框架包装的 ``{"raw_body": "..."}`` dict - ``raw_event.account_id`` 由框架注入 + - 加密模式 / token / encoding_aes_key / app_id 由适配器从 + ``ConfigPort`` 读取(不依赖框架注入到 headers) @post - 返回 ``MessageContent``,含文本与可选附件 @failure - ``ValidationError``:字段缺失 / 未知 MsgType / 解密失败 - - ``DependencyError``:网络异常 + - ``DependencyError``:网络异常 / 配置读取失败 @consistency - 无状态;``@idempotent: True`` @@ -107,23 +110,28 @@ class WeChatMpInboundAdapter: payload = raw_event.payload account_id = raw_event.account_id or "" - # 兼容两种 payload 形式:XML 字符串 或 已解析 dict(测试/框架透传) + # 兼容三种 payload 形式: + # - XML 字符串(直传场景) + # - 已解析 dict(测试/框架透传) + # - 框架包装 {"raw_body": "..."}(非 JSON body 经 + # InboundMessageService._buildRawEvent 包装,生产 webhook 路径) if isinstance(payload, dict): - xml_dict = payload + raw_body = payload.get("raw_body") + if isinstance(raw_body, str) and raw_body: + # 框架包装的 raw_body:解析 XML 字符串 + xml_dict = await self._parse_payload_text( + raw_body, + account_id=account_id, + ) + else: + # 已解析 dict(测试/框架透传,已含 MsgType 等字段) + xml_dict = payload else: xml_text = payload if isinstance(payload, str) else str(payload) - # 解密(兼容/安全模式) - encrypt_mode = raw_event.headers.get("encrypt_mode", ENCRYPT_MODE_PLAINTEXT) - encoding_aes_key = raw_event.headers.get("encoding_aes_key", "") - app_id = raw_event.headers.get("app_id", "") - if encrypt_mode in (ENCRYPT_MODE_COMPATIBLE, ENCRYPT_MODE_SAFE): - xml_text = self._decrypt_if_needed( - xml_text, - encoding_aes_key=encoding_aes_key, - app_id=app_id, - encrypt_mode=encrypt_mode, - ) - xml_dict = _parse_xml_to_dict(xml_text) + xml_dict = await self._parse_payload_text( + xml_text, + account_id=account_id, + ) msg_type = xml_dict.get("MsgType") if not msg_type: @@ -194,7 +202,9 @@ class WeChatMpInboundAdapter: """校验微信 Webhook 签名。 @pre - - ``raw_event.headers`` 含 signature / timestamp / nonce / token + - ``raw_event.headers`` 含 signature / timestamp / nonce + (由 ``webhook_router`` 从 URL query 参数合并到 headers) + - ``raw_event.account_id`` 由框架注入 @post - 明文模式 → 始终返回 ``valid=True`` @@ -207,11 +217,13 @@ class WeChatMpInboundAdapter: @consistency - 无状态;``@idempotent: True`` """ - encrypt_mode = raw_event.headers.get("encrypt_mode", ENCRYPT_MODE_PLAINTEXT) + encrypt_mode = await self._get_encrypt_mode() if encrypt_mode == ENCRYPT_MODE_PLAINTEXT: return SignatureVerifyResult(valid=True) - token = raw_event.headers.get("token", "") + account_id = raw_event.account_id or "" + config = await self._get_account_config(account_id) + token = config["token"] if not token: raise ValidationError(field="token", message="token_required") @@ -224,7 +236,14 @@ class WeChatMpInboundAdapter: # 兼容模式:明文消息无 ,签名为 SHA1(sort([token, timestamp, nonce]))(3 参数) encrypt = "" payload = raw_event.payload - if isinstance(payload, str): + if isinstance(payload, dict): + raw_body = payload.get("raw_body") + if isinstance(raw_body, str): + encrypt = crypto.extract_encrypt(raw_body) or "" + elif "Encrypt" in payload: + # 已解析 dict 含 Encrypt 字段(测试/框架透传场景) + encrypt = str(payload.get("Encrypt") or "") + elif isinstance(payload, str): encrypt = crypto.extract_encrypt(payload) or "" if encrypt: @@ -289,6 +308,67 @@ class WeChatMpInboundAdapter: # 内部辅助 # ------------------------------------------------------------------ + async def _get_account_config(self, account_id: str) -> dict[str, str]: + """读取账户级签名/解密配置。 + + 从 ``ConfigPort`` 读取 account-level 配置(token / encoding_aes_key / + app_id),与 wecom 适配器保持一致。框架 ``_buildRawEvent`` 不注入 + account-level 配置到 ``raw_event.headers``,适配器必须自行从 + ``ConfigPort`` 读取。 + + @failure ``account_id`` 为空或配置读取失败时抛 ``DependencyError``。 + """ + from yuxi.channels.contract.dtos.config import ConfigScope + + if not account_id: + raise DependencyError( + dep=WECHAT_MP_DEP, + cause=ValueError("account_id missing for config resolution"), + ) + result: dict[str, str] = {} + for key in ("token", "encoding_aes_key", "app_id"): + cv = await self._config.get( + key, + scope=ConfigScope.ACCOUNT, + target=account_id, + ) + result[key] = str(cv.value) if cv.value else "" + return result + + async def _get_encrypt_mode(self) -> str: + """从 ConfigPort 读取 channel-level 加密模式。""" + from yuxi.channels.contract.dtos.config import ConfigScope + + cv = await self._config.get( + "encrypt_mode", + scope=ConfigScope.CHANNEL, + target=CHANNEL_TARGET, + ) + return str(cv.value) if cv.value else ENCRYPT_MODE_PLAINTEXT + + async def _parse_payload_text( + self, + xml_text: str, + *, + account_id: str, + ) -> dict[str, Any]: + """解析 XML 文本为 dict,按需先解密。 + + 统一处理 XML 字符串场景(直传与框架 ``raw_body`` 包装), + 从 ``ConfigPort`` 读取 account-level 配置(encoding_aes_key / app_id) + 与 channel-level 配置(encrypt_mode),按需解密后解析 XML。 + """ + encrypt_mode = await self._get_encrypt_mode() + if encrypt_mode in (ENCRYPT_MODE_COMPATIBLE, ENCRYPT_MODE_SAFE): + config = await self._get_account_config(account_id) + xml_text = self._decrypt_if_needed( + xml_text, + encoding_aes_key=config["encoding_aes_key"], + app_id=config["app_id"], + encrypt_mode=encrypt_mode, + ) + return _parse_xml_to_dict(xml_text) + def _decrypt_if_needed( self, xml_text: str, diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/lifecycle_adapter.py b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/lifecycle_adapter.py index cf3ab446..008fb20f 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/lifecycle_adapter.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/lifecycle_adapter.py @@ -42,6 +42,8 @@ from yuxi.channels.contract.ports.driven.config_port import ConfigPort from yuxi.channels.contract.ports.driven.logger_port import LoggerPort from .._constants import ( + _CACHE_INVALIDATE_PATTERN_TEMPLATE, + _CACHE_RUNTIME_INVALIDATE_PATTERN_TEMPLATE, APP_ID_PATTERN, DEFAULT_API_BASE_URL, DEFAULT_HTTP_TIMEOUT_SECONDS, @@ -50,8 +52,6 @@ from .._constants import ( ENCRYPT_MODE_COMPATIBLE, ENCRYPT_MODE_PLAINTEXT, ENCRYPT_MODE_SAFE, - _CACHE_INVALIDATE_PATTERN_TEMPLATE, - _CACHE_RUNTIME_INVALIDATE_PATTERN_TEMPLATE, ) from ..wechat_mp_client import WeChatMpClient diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/session_adapter.py b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/session_adapter.py index 43d420a0..180e39fd 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/session_adapter.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/session_adapter.py @@ -9,12 +9,20 @@ - ``group_id`` / ``topic_id`` 始终为 ``None``。 - ``account_id`` 从 ``raw_event.account_id`` 提取(框架已注入)。 +payload 形式兼容: + - 已解析 dict(测试/框架透传,含 FromUserName 等字段):直接读取。 + - 框架包装 ``{"raw_body": "..."}``(非 JSON body 经 + ``InboundMessageService._buildRawEvent`` 包装):解析 XML 后读取。 + 依赖方向:仅 import ``yuxi.channels.contract.*`` 与同插件内部 ``_constants``, 不污染框架层。 """ from __future__ import annotations +import xml.etree.ElementTree as ET +from typing import Any + from yuxi.channels.contract.dtos.channel import ChannelType, SessionInfo from yuxi.channels.contract.dtos.common import RawEvent from yuxi.channels.contract.errors import ValidationError @@ -85,12 +93,48 @@ class WeChatMpSessionAdapter: # ---------------------------------------------------------------------- +def _resolve_payload(payload: dict[str, Any] | str) -> dict[str, Any]: + """解析 payload,兼容框架 ``raw_body`` 包装。 + + - payload 为 dict 且含 ``raw_body`` 字符串:解析 XML 返回字段 dict。 + - payload 为 dict 不含 ``raw_body``:直接返回。 + - payload 为 str:解析 XML 返回字段 dict。 + """ + if isinstance(payload, str): + return _parse_xml_to_dict(payload) + if isinstance(payload, dict): + raw_body = payload.get("raw_body") + if isinstance(raw_body, str) and raw_body: + return _parse_xml_to_dict(raw_body) + return payload + return {} + + +def _parse_xml_to_dict(xml_text: str) -> dict[str, Any]: + """将微信 XML 消息解析为扁平字典。""" + try: + root = ET.fromstring(xml_text) + except ET.ParseError as exc: + raise ValidationError( + field="payload", + message="xml_parse_failed", + ) from exc + result: dict[str, Any] = {} + for child in root: + result[child.tag] = child.text or "" + return result + + def _extract_from_user_name(raw_event: RawEvent) -> str: """从 ``raw_event.payload`` 提取 ``FromUserName`` 字段。 + 兼容框架 ``raw_body`` 包装:payload 为 ``{"raw_body": "..."}`` + 时先解析 XML 再提取字段。 + @failure ``FromUserName`` 缺失或为空抛 ``ValidationError``。 """ - from_user_name = raw_event.payload.get("FromUserName", "") + payload = _resolve_payload(raw_event.payload) + from_user_name = payload.get("FromUserName", "") if not from_user_name: raise ValidationError( field="FromUserName", diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/status_adapter.py b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/status_adapter.py index a1226c8d..ac9f282d 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/adapters/status_adapter.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/adapters/status_adapter.py @@ -16,11 +16,17 @@ - 所有事件均返回 ``StatusPayload``,``ref_channel_msg_id`` 使用 ``f"{FromUserName}_{CreateTime}"`` 构造(系统事件无 ``MsgId``)。 +payload 形式兼容: + - 已解析 dict(测试/框架透传,含 MsgType 等字段):直接读取。 + - 框架包装 ``{"raw_body": "..."}``(非 JSON body 经 + ``InboundMessageService._buildRawEvent`` 包装):解析 XML 后读取。 + 依赖方向:仅 import ``yuxi.channels.contract.*``,不污染框架层。 """ from __future__ import annotations +import xml.etree.ElementTree as ET from datetime import UTC, datetime from typing import Any @@ -63,11 +69,12 @@ class WeChatMpStatusAdapter: @consistency 无状态;``@idempotent: True``。 """ - msg_type = raw_event.payload.get("MsgType", "") + payload = _resolve_payload(raw_event.payload) + msg_type = payload.get("MsgType", "") if msg_type in _MESSAGE_TYPES: return EventType.MESSAGE if msg_type == "event": - event = raw_event.payload.get("Event", "") + event = payload.get("Event", "") if event in _SYSTEM_EVENT_TYPES: return EventType.MESSAGE return EventType.UNKNOWN @@ -82,11 +89,12 @@ class WeChatMpStatusAdapter: @failure ``CreateTime`` 缺失或非数字抛 ``ValidationError(field="CreateTime")``。 @consistency 无状态;``@idempotent: True``。 """ + payload = _resolve_payload(raw_event.payload) event_type = await self.classifyEvent(raw_event) - from_user = raw_event.payload.get("FromUserName", "") - create_time = raw_event.payload.get("CreateTime", "") + from_user = payload.get("FromUserName", "") + create_time = payload.get("CreateTime", "") ref_id = f"{from_user}_{create_time}" if from_user and create_time else from_user or str(create_time) - timestamp = _extract_timestamp(raw_event.payload) + timestamp = _extract_timestamp(payload) return StatusPayload( event_type=event_type, ref_channel_msg_id=ref_id, @@ -99,6 +107,38 @@ class WeChatMpStatusAdapter: # ---------------------------------------------------------------------- +def _resolve_payload(payload: dict[str, Any] | str) -> dict[str, Any]: + """解析 payload,兼容框架 ``raw_body`` 包装。 + + - payload 为 dict 且含 ``raw_body`` 字符串:解析 XML 返回字段 dict。 + - payload 为 dict 不含 ``raw_body``:直接返回。 + - payload 为 str:解析 XML 返回字段 dict。 + """ + if isinstance(payload, str): + return _parse_xml_to_dict(payload) + if isinstance(payload, dict): + raw_body = payload.get("raw_body") + if isinstance(raw_body, str) and raw_body: + return _parse_xml_to_dict(raw_body) + return payload + return {} + + +def _parse_xml_to_dict(xml_text: str) -> dict[str, Any]: + """将微信 XML 消息解析为扁平字典。""" + try: + root = ET.fromstring(xml_text) + except ET.ParseError as exc: + raise ValidationError( + field="payload", + message="xml_parse_failed", + ) from exc + result: dict[str, Any] = {} + for child in root: + result[child.tag] = child.text or "" + return result + + def _extract_timestamp(payload: dict[str, Any]) -> datetime: """从 ``payload.CreateTime`` 提取时间戳。 diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/entry.py b/backend/package/yuxi/channels/plugins/wechat_mp/entry.py index e0126b97..352543b7 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/entry.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/entry.py @@ -41,15 +41,13 @@ client 实例化、LifecycleHandler 构造与 ``registerAdapter`` 调用。 - ``directory`` 接收 ``(client, cache_port, logger_port)``:调用 ``get_user_list`` / ``get_user_info`` 查询关注者,用户信息缓存于 CachePort。 -传输管理:wechat_mp 通过 ``TransportManager`` 的 ``both`` 模式管理消息 -接收(manifest ``transport_mode: "both"``): -- ``WebhookWorker``(优先):接收微信服务器 Webhook 推送(消息/事件), - 端到端延迟 ≤5s(5 秒被动回复策略) -- ``PullerWorker``(降级):Webhook 不可用时降级为无轮询(公众号无拉取 - API,此模式实际不可用,仅保留契约一致性) -两者由传输引擎通过 ``ChannelAccountOnline`` / ``ChannelAccountOffline`` -领域事件触发启停,``LifecycleAdapter`` 与 ``LifecycleHandler`` 不自管理 -传输任务(决策 3)。 +传输管理:wechat_mp 为纯 Webhook 接入渠道,不注册 ``PullerAdapter`` / +``StreamConnectorAdapter``(决策 3)。manifest 声明 ``transport_mode: "pull"`` +以通过 ``manifest_loader`` 校验(仅接受 ``pull``/``stream``/``both``), +实际消息接收通过 Webhook 端点 ``POST /channels/wechat_mp/webhook`` 由 +``webhook_router`` → ``InboundMessageService.receiveWebhook`` 进入入站 +管道,端到端延迟 ≤5s(5 秒被动回复策略)。``TransportManager`` 找不到 +puller/stream_connector 适配器时不启动传输任务,不影响 Webhook 接入。 依赖方向:仅 import ``yuxi.channels.contract.*`` + 标准库 + 同插件内部模块, 不污染框架层。 diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/lifecycle.py b/backend/package/yuxi/channels/plugins/wechat_mp/lifecycle.py index 18df4a46..ece59bba 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/lifecycle.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/lifecycle.py @@ -323,7 +323,7 @@ class WeChatMpLifecycleHandler: cv = await self._config.get( "token_refresh_interval_s", scope=ConfigScope.CHANNEL, - target="wechat_mp", + target=CHANNEL_TARGET, ) return int(cv.value) except Exception as exc: @@ -341,7 +341,7 @@ class WeChatMpLifecycleHandler: cv = await self._config.get( "http_timeout_ms", scope=ConfigScope.CHANNEL, - target="wechat_mp", + target=CHANNEL_TARGET, ) return int(cv.value) / 1000.0 except Exception as exc: diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/manifest.json b/backend/package/yuxi/channels/plugins/wechat_mp/manifest.json index 2ba0671d..d55626f0 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/manifest.json +++ b/backend/package/yuxi/channels/plugins/wechat_mp/manifest.json @@ -26,7 +26,7 @@ }, "max_message_length": 2048, "supports_credential_cloning": false, - "transport_mode": "both", + "transport_mode": "pull", "capabilities": { "rich_message": true, "streaming": false, diff --git a/backend/package/yuxi/channels/plugins/wechat_mp/wechat_mp_client.py b/backend/package/yuxi/channels/plugins/wechat_mp/wechat_mp_client.py index 97a55d81..cc49960b 100644 --- a/backend/package/yuxi/channels/plugins/wechat_mp/wechat_mp_client.py +++ b/backend/package/yuxi/channels/plugins/wechat_mp/wechat_mp_client.py @@ -48,6 +48,7 @@ from ._constants import ( ACCESS_TOKEN_CACHE_TTL_S, ACCESS_TOKEN_LOCK_KEY_PREFIX, ACCESS_TOKEN_LOCK_TTL_S, + CHANNEL_TARGET, DEFAULT_API_BASE_URL, DEFAULT_HTTP_TIMEOUT_SECONDS, DEFAULT_MAX_MESSAGE_LENGTH, @@ -181,7 +182,7 @@ class WeChatMpClient: cv = await self._config.get( "api_base_url", scope=ConfigScope.CHANNEL, - target="wechat_mp", + target=CHANNEL_TARGET, ) base_url = str(cv.value) if cv.value else DEFAULT_API_BASE_URL return base_url.rstrip("/") or DEFAULT_API_BASE_URL @@ -200,7 +201,7 @@ class WeChatMpClient: cv = await self._config.get( "http_timeout_ms", scope=ConfigScope.CHANNEL, - target="wechat_mp", + target=CHANNEL_TARGET, ) return int(cv.value) / 1000.0 except Exception as exc: @@ -595,7 +596,9 @@ class WeChatMpClient: ) resp = await self._invoke_with_token_retry( - account_id, _call, operation="send_custom_message", + account_id, + _call, + operation="send_custom_message", ) result = self._parse_json_response(resp) await self._logger.debug( @@ -641,11 +644,16 @@ class WeChatMpClient: params = {"access_token": access_token, "type": media_type} # multipart 上传不走 _execute_http(json_body 会误设 Content-Type) return await self._execute_multipart( - url, account_id=account_id, params=params, files=files, + url, + account_id=account_id, + params=params, + files=files, ) resp = await self._invoke_with_token_retry( - account_id, _call, operation="upload_media", + account_id, + _call, + operation="upload_media", ) data = self._parse_json_response(resp) media_id = data.get("media_id") @@ -757,7 +765,9 @@ class WeChatMpClient: ) resp = await self._invoke_with_token_retry( - account_id, _call, operation="download_media", + account_id, + _call, + operation="download_media", ) # 微信媒体下载可能返回 JSON 错误体(errcode!=0)或二进制流 @@ -808,7 +818,9 @@ class WeChatMpClient: ) resp = await self._invoke_with_token_retry( - account_id, _call, operation="get_user_info", + account_id, + _call, + operation="get_user_info", ) return self._parse_json_response(resp, user_id_hint=openid) @@ -840,7 +852,9 @@ class WeChatMpClient: ) resp = await self._invoke_with_token_retry( - account_id, _call, operation="get_user_list", + account_id, + _call, + operation="get_user_list", ) return self._parse_json_response(resp) @@ -863,7 +877,9 @@ class WeChatMpClient: ) resp = await self._invoke_with_token_retry( - account_id, _call, operation="get_callback_ip", + account_id, + _call, + operation="get_callback_ip", ) return self._parse_json_response(resp) @@ -928,7 +944,7 @@ class WeChatMpClient: cv = await self._config.get( "media_cache_ttl_s", scope=ConfigScope.CHANNEL, - target="wechat_mp", + target=CHANNEL_TARGET, ) ttl = int(cv.value) except Exception: @@ -998,8 +1014,6 @@ class WeChatMpClient: - 返回 True 表示额度充足,已消费 1 条 - 返回 False 表示额度耗尽,不应发送 """ - from ._constants import CUSTOMER_SERVICE_QUOTA_OTHER_TTL_S - quota_key = self._quota_cache_key(account_id, openid) cached = await self._cache.get(quota_key) if not isinstance(cached, Some): @@ -1034,7 +1048,7 @@ class WeChatMpClient: cv = await self._config.get( "max_message_length", scope=ConfigScope.CHANNEL, - target="wechat_mp", + target=CHANNEL_TARGET, ) return int(cv.value) except Exception as exc: