Files
inquiry_robot/inquiry-agent/agent/channel/outbox/sender.py
T

140 lines
4.9 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
企微主动发消息客户端(应用消息)。
本文件职责:取 access_token、发送文本;密钥来自环境变量。
禁止:在发送成功时冒充业务成功;禁止按 userid 伪造业务分支。
出站唯一入口应经 outbox 消费调用本客户端。
"""
from __future__ import annotations
import logging
import threading
import time
from dataclasses import dataclass
from typing import Any, Optional
import httpx
logger = logging.getLogger(__name__)
@dataclass
class WeComSendResult:
"""企微发送结果(通道层)。"""
ok: bool
errcode: int = 0
errmsg: str = ""
raw: dict[str, Any] | None = None
class WeComAppClient:
"""
询价应用(默认 AgentId 1000010)主动消息客户端。
线程安全:token 缓存带锁;HTTP 带超时,禁止无限阻塞。
"""
def __init__(
self,
*,
corp_id: str,
secret: str,
agent_id: int,
api_base: str = "https://qyapi.weixin.qq.com",
timeout_sec: float = 10.0,
) -> None:
self._corp_id = corp_id
self._secret = secret
self._agent_id = agent_id
self._api_base = api_base.rstrip("/")
self._timeout = timeout_sec
self._token: Optional[str] = None
self._token_expire_at: float = 0.0
self._lock = threading.Lock()
def send_text(self, *, touser: str, content: str) -> WeComSendResult:
"""
向指定企微 userid 发送文本应用消息。
参数:touser 必须是入站 sender_id;content 测试环境应由上层加 [test] 前缀。
副作用:出站 HTTP 调企微;失败返回 ok=False,不抛给回调线程(由 outbox 重试)。
"""
if not self._corp_id or not self._secret:
return WeComSendResult(ok=False, errcode=-1, errmsg="WECOM_CORP_ID/SECRET 未配置")
if not touser or not content:
return WeComSendResult(ok=False, errcode=-1, errmsg="touser/content 为空")
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("获取 access_token 失败")
return WeComSendResult(ok=False, errcode=-1, errmsg=str(exc))
url = f"{self._api_base}/cgi-bin/message/send"
body = {
"touser": touser,
"msgtype": "text",
"agentid": self._agent_id,
"text": {"content": content},
"safe": 0,
}
try:
with httpx.Client(timeout=self._timeout) as client:
resp = client.post(url, params={"access_token": token}, json=body)
data = resp.json()
except Exception as exc: # noqa: BLE001
logger.exception("企微 message/send 网络失败")
return WeComSendResult(ok=False, errcode=-1, errmsg=str(exc))
errcode = int(data.get("errcode") or 0)
errmsg = str(data.get("errmsg") or "")
if errcode != 0:
# token 失效则清空缓存,下次重试
if errcode in (40014, 42001, 42007):
with self._lock:
self._token = None
self._token_expire_at = 0.0
logger.warning("企微发送失败 errcode=%s errmsg=%s", errcode, errmsg)
return WeComSendResult(ok=False, errcode=errcode, errmsg=errmsg, raw=data)
preview = content if len(content) <= 200 else (content[:200] + "…")
logger.info("企微发送成功 touser=%s reply=%s", touser, preview)
return WeComSendResult(ok=True, errcode=0, errmsg=errmsg, raw=data)
def _get_token(self) -> str:
with self._lock:
now = time.time()
if self._token and now < self._token_expire_at - 60:
return self._token
# 注意:勿把 corpsecret 打进日志;httpx 默认 INFO 会打印完整 URL
url = f"{self._api_base}/cgi-bin/gettoken"
with httpx.Client(timeout=self._timeout) as client:
resp = client.get(
url,
params={"corpid": self._corp_id, "corpsecret": self._secret},
)
data = resp.json()
errcode = int(data.get("errcode") or 0)
if errcode != 0 or not data.get("access_token"):
raise RuntimeError(f"gettoken 失败 errcode={errcode}")
token = str(data["access_token"])
expires_in = int(data.get("expires_in") or 7200)
with self._lock:
self._token = token
self._token_expire_at = time.time() + expires_in
return token
# 企微侧明确不可重试的错误(IP 白名单、非法 userid 等)
PERMANENT_SEND_ERRORS = frozenset(
{
60020, # not allow to access from your ip
60111, # invalid user
81013, # user not in app visible range
40003, # invalid openid/userid
40013, # invalid corpid
40001, # invalid credential(非 token 过期类时也停)
}
)