refactor(wechat-mp): 重构微信公众号插件适配纯Webhook模式
1. 将transport_mode改为pull以通过manifest校验,实际仅使用Webhook接收消息 2. 统一处理XML payload兼容raw_body包装格式 3. 重构配置读取逻辑,使用CHANNEL_TARGET常量替换硬编码target 4. 移除冗余的webhook_url配置诊断项,更新文档说明 5. 优化代码格式与导入顺序
This commit is contained in:
parent
ebbd552970
commit
feee2e1166
@ -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 ""
|
||||
|
||||
@ -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 字符串
|
||||
(明文模式直接原文;兼容/安全模式为含 ``<Encrypt>`` 的密文 XML)
|
||||
- ``raw_event.headers`` 含 ``encrypt_mode``(可选,默认 plaintext)
|
||||
或框架包装的 ``{"raw_body": "<xml>...</xml>"}`` 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": "<xml>...</xml>"}(非 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:
|
||||
# 兼容模式:明文消息无 <Encrypt>,签名为 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,
|
||||
|
||||
@ -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
|
||||
|
||||
|
||||
@ -9,12 +9,20 @@
|
||||
- ``group_id`` / ``topic_id`` 始终为 ``None``。
|
||||
- ``account_id`` 从 ``raw_event.account_id`` 提取(框架已注入)。
|
||||
|
||||
payload 形式兼容:
|
||||
- 已解析 dict(测试/框架透传,含 FromUserName 等字段):直接读取。
|
||||
- 框架包装 ``{"raw_body": "<xml>...</xml>"}``(非 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>...</xml>"}``
|
||||
时先解析 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",
|
||||
|
||||
@ -16,11 +16,17 @@
|
||||
- 所有事件均返回 ``StatusPayload``,``ref_channel_msg_id`` 使用
|
||||
``f"{FromUserName}_{CreateTime}"`` 构造(系统事件无 ``MsgId``)。
|
||||
|
||||
payload 形式兼容:
|
||||
- 已解析 dict(测试/框架透传,含 MsgType 等字段):直接读取。
|
||||
- 框架包装 ``{"raw_body": "<xml>...</xml>"}``(非 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`` 提取时间戳。
|
||||
|
||||
|
||||
@ -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.*`` + 标准库 + 同插件内部模块,
|
||||
不污染框架层。
|
||||
|
||||
@ -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:
|
||||
|
||||
@ -26,7 +26,7 @@
|
||||
},
|
||||
"max_message_length": 2048,
|
||||
"supports_credential_cloning": false,
|
||||
"transport_mode": "both",
|
||||
"transport_mode": "pull",
|
||||
"capabilities": {
|
||||
"rich_message": true,
|
||||
"streaming": false,
|
||||
|
||||
@ -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:
|
||||
|
||||
Loading…
Reference in New Issue
Block a user