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

643 lines
25 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_outbound(
self,
*,
touser: str,
content: str,
payload: dict[str, Any] | None = None,
) -> WeComSendResult:
"""
按 outbox 内容发应用消息:有模板卡先发卡,失败再回退文本。
副作用:出站 HTTP 调企微;失败返回 ok=False,由 outbox 重试。
"""
extra = dict(payload or {})
if extra.get("msgtype") == "image":
image_result = self.send_image(
touser=touser,
filename=str(extra.get("filename") or "image.png"),
filepath=str(extra.get("filepath") or ""),
file_b64=str(extra.get("file_b64") or extra.get("fileBase64") or ""),
)
if image_result.ok:
return image_result
logger.warning(
"图片发送失败,回退文本 err=%s:%s",
image_result.errcode,
image_result.errmsg,
)
return self.send_text(touser=touser, content=content)
if extra.get("msgtype") == "file":
file_result = self.send_file(
touser=touser,
filename=str(extra.get("filename") or "quote.xlsx"),
filepath=str(extra.get("filepath") or ""),
file_b64=str(extra.get("file_b64") or extra.get("fileBase64") or ""),
)
# 成交卡绑在同一条出站:文件 HTTP 成功后再发卡。
# 企微文件气泡要等素材落地,卡片是即时的;同一毫秒连发,销售会先看到成交跟进。
if file_result.ok:
gap = float(getattr(self, "_followup_gap_sec", 1.6) or 0)
if gap > 0:
time.sleep(gap)
self._send_followup_card(touser=touser, extra=extra)
return file_result
if extra.get("msgtype") == "update_template_card":
return self.update_template_card(
touser=touser,
response_code=str(extra.get("response_code") or ""),
replace_name=str((extra.get("button") or {}).get("replace_name") or ""),
)
if extra.get("msgtype") == "template_card" and extra.get("template_card"):
card_result = self._post_message(
{
"touser": touser,
"msgtype": "template_card",
"agentid": self._agent_id,
"template_card": extra["template_card"],
"enable_id_trans": 0,
},
preview="[template_card]",
)
if card_result.ok:
return card_result
logger.warning(
"模板卡发送失败,回退文本 err=%s:%s",
card_result.errcode,
card_result.errmsg,
)
return self.send_text(touser=touser, content=content)
def update_template_card(
self,
*,
touser: str,
response_code: str,
replace_name: str,
) -> WeComSendResult:
"""
按回调 ResponseCode 更新原模板卡按钮文案,保持置灰不可再点。
失败不回落成新文本气泡,避免销售看到重复「已提交」。
"""
if not touser or not response_code or not replace_name:
return WeComSendResult(ok=True, errcode=0, errmsg="skip_card_update")
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("更新模板卡取 token 失败")
return WeComSendResult(ok=True, errcode=-1, errmsg=str(exc))
url = f"{self._api_base}/cgi-bin/message/update_template_card"
body = {
"userids": [touser],
"agentid": self._agent_id,
"response_code": response_code,
"enable_id_trans": 0,
"button": {"replace_name": replace_name},
}
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("企微 update_template_card 网络失败")
return WeComSendResult(ok=True, errcode=-1, errmsg=str(exc))
errcode = int(data.get("errcode") or 0)
errmsg = str(data.get("errmsg") or "")
if errcode != 0:
logger.warning("更新模板卡失败 errcode=%s errmsg=%s", errcode, errmsg)
return WeComSendResult(ok=True, errcode=errcode, errmsg=errmsg, raw=data)
logger.info("企微模板卡已置灰 touser=%s", touser)
return WeComSendResult(ok=True, errcode=0, errmsg=errmsg, raw=data)
def _send_followup_card(self, *, touser: str, extra: dict[str, Any]) -> None:
"""
报价文件发出后再发跟进卡。失败不回滚文件、不触发整条重试。
跟进卡单独入队会被另一条出站线程抢走,销售就会先看到成交跟进。
"""
card = extra.get("followup_card")
if not isinstance(card, dict) or not card:
nested = extra.get("followup_payload")
if isinstance(nested, dict):
card = nested.get("template_card")
if not isinstance(card, dict) or not card:
return
result = self._post_message(
{
"touser": touser,
"msgtype": "template_card",
"agentid": self._agent_id,
"template_card": card,
"enable_id_trans": 0,
},
preview="[followup_card]",
)
if not result.ok:
logger.warning(
"报价文件已发出,跟进卡失败 err=%s:%s",
result.errcode,
result.errmsg,
)
def send_file(
self,
*,
touser: str,
filename: str,
filepath: str = "",
file_b64: str = "",
) -> WeComSendResult:
"""
先上传临时素材再发文件消息。
字节来自本机路径或 Base64;失败回 ok=False 由 outbox 重试。
禁止在回调线程调用;本方法只给出站消费线程用。
"""
if not touser:
return WeComSendResult(ok=False, errcode=-1, errmsg="touser 为空")
raw = b""
if filepath:
try:
from pathlib import Path
raw = Path(filepath).read_bytes()
except OSError as exc:
return WeComSendResult(ok=False, errcode=-1, errmsg=f"read_file:{exc}")
elif file_b64:
import base64
try:
raw = base64.b64decode(file_b64)
except Exception as exc: # noqa: BLE001
return WeComSendResult(ok=False, errcode=-1, errmsg=f"b64:{exc}")
if not raw:
return WeComSendResult(ok=False, errcode=-1, errmsg="file_empty")
media_id, up = self._upload_media(
filename=filename or "quote.xlsx", content=raw, media_type="file"
)
if not media_id:
return up
return self._post_message(
{
"touser": touser,
"msgtype": "file",
"agentid": self._agent_id,
"file": {"media_id": media_id},
},
preview=f"[file]{filename}",
)
def send_image(
self,
*,
touser: str,
filename: str,
filepath: str = "",
file_b64: str = "",
) -> WeComSendResult:
"""
先上传 image 素材再发图片消息,聊天里直接看图。
字节来自本机路径或 Base64。禁止在回调线程调用。
"""
if not touser:
return WeComSendResult(ok=False, errcode=-1, errmsg="touser 为空")
raw = b""
if filepath:
try:
from pathlib import Path
raw = Path(filepath).read_bytes()
except OSError as exc:
return WeComSendResult(ok=False, errcode=-1, errmsg=f"read_file:{exc}")
elif file_b64:
import base64
try:
raw = base64.b64decode(file_b64)
except Exception as exc: # noqa: BLE001
return WeComSendResult(ok=False, errcode=-1, errmsg=f"b64:{exc}")
if not raw:
return WeComSendResult(ok=False, errcode=-1, errmsg="image_empty")
name = filename or "image.png"
media_id, up = self._upload_media(
filename=name, content=raw, media_type="image"
)
if not media_id:
return up
return self._post_message(
{
"touser": touser,
"msgtype": "image",
"agentid": self._agent_id,
"image": {"media_id": media_id},
},
preview=f"[image]{name}",
)
def _upload_media(
self,
*,
filename: str,
content: bytes,
media_type: str = "file",
) -> tuple[str, WeComSendResult]:
"""cgi-bin/media/upload;type=file 或 image。密钥不写日志。"""
kind = (media_type or "file").strip() or "file"
mime = "image/png" if kind == "image" else "application/octet-stream"
if not self._corp_id or not self._secret:
return "", WeComSendResult(ok=False, errcode=-1, errmsg="WECOM_CORP_ID/SECRET 未配置")
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("上传素材取 token 失败")
return "", WeComSendResult(ok=False, errcode=-1, errmsg=str(exc))
url = f"{self._api_base}/cgi-bin/media/upload"
try:
with httpx.Client(timeout=max(self._timeout, 30.0)) as client:
resp = client.post(
url,
params={"access_token": token, "type": kind},
files={"media": (filename, content, mime)},
)
data = resp.json()
except Exception as exc: # noqa: BLE001
logger.exception("企微 media/upload 网络失败")
return "", WeComSendResult(ok=False, errcode=-1, errmsg=str(exc))
errcode = int(data.get("errcode") or 0)
media_id = str(data.get("media_id") or "")
if errcode != 0 or not media_id:
logger.warning("企微上传素材失败 type=%s errcode=%s", kind, errcode)
return "", WeComSendResult(
ok=False, errcode=errcode, errmsg=str(data.get("errmsg") or "upload_fail"), raw=data
)
return media_id, WeComSendResult(ok=True, errcode=0, raw=data)
def create_group(self, *, name: str, userids: list[str]) -> dict[str, Any]:
"""
创建企微应用群(appchat/create),群名=工单号。
成员须含销售、产品、会话存档账号(询价机器人);应用自己会进群。失败回 ok=False,不抛给回调。
"""
users = [str(x).strip() for x in userids if str(x).strip()]
if not name or len(users) < 2:
return {"ok": False, "error": "name_or_users", "chat_id": ""}
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("建群取 token 失败")
return {"ok": False, "error": str(exc), "chat_id": ""}
url = f"{self._api_base}/cgi-bin/appchat/create"
body = {"name": name, "userlist": users}
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("企微建群网络失败")
return {"ok": False, "error": str(exc), "chat_id": ""}
errcode = int(data.get("errcode") or 0)
chat_id = str(data.get("chatid") or data.get("chat_id") or "")
if errcode != 0 or not chat_id:
logger.warning("企微建群失败 errcode=%s", errcode)
return {
"ok": False,
"error": str(data.get("errmsg") or f"errcode_{errcode}"),
"chat_id": "",
}
logger.info("企微群已创建 name=%s chat=%s", name, chat_id)
return {"ok": True, "chat_id": chat_id, "error": ""}
def add_group_members(self, *, chat_id: str, userids: list[str]) -> dict[str, Any]:
"""
往已有应用群补人(appchat/update add_user_list)。
用于已建群补拉会话存档账号。已在群里的人企微会忽略。不抛给回调。
"""
users = [str(x).strip() for x in userids if str(x).strip()]
if not chat_id or not users:
return {"ok": False, "error": "chat_or_users"}
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("补拉群成员取 token 失败")
return {"ok": False, "error": str(exc)}
url = f"{self._api_base}/cgi-bin/appchat/update"
body = {"chatid": chat_id, "add_user_list": users}
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("企微补拉群成员网络失败")
return {"ok": False, "error": str(exc)}
errcode = int(data.get("errcode") or 0)
if errcode != 0:
logger.warning("企微补拉群成员失败 errcode=%s", errcode)
return {"ok": False, "error": str(data.get("errmsg") or f"errcode_{errcode}")}
logger.info("企微群已补人 chat=%s", chat_id)
return {"ok": True, "error": ""}
def send_group(
self,
*,
chat_id: str,
content: str,
mention_userids: list[str] | None = None,
) -> dict[str, Any]:
"""
向应用群发文本。必须走 appchat/send,不能当私聊 message/send。
mention_userids 只给调用方记账,不传给企微。
正文开头已经写了 @姓名;再传企微点人名单,结尾会再点一次。
失败回 ok=False,由调用方决定是否重试;不抛给回调线程。
"""
_ = mention_userids
if not chat_id or not content:
return {"ok": False, "error": "chat_or_content"}
if not self._corp_id or not self._secret:
return {"ok": False, "error": "WECOM_CORP_ID/SECRET 未配置"}
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("群消息取 token 失败")
return {"ok": False, "error": str(exc)}
url = f"{self._api_base}/cgi-bin/appchat/send"
body: dict[str, Any] = {"chatid": chat_id, "msgtype": "text", "text": {"content": content}}
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("企微群消息网络失败")
return {"ok": False, "error": str(exc)}
errcode = int(data.get("errcode") or 0)
if errcode != 0:
logger.warning("企微群消息失败 errcode=%s chat=%s", errcode, chat_id)
return {"ok": False, "error": str(data.get("errmsg") or f"errcode_{errcode}")}
logger.info("企微群文本已发送 chat=%s chars=%s", chat_id, len(content))
return {"ok": True, "error": ""}
def send_group_file(
self,
*,
chat_id: str,
filename: str,
filepath: str = "",
file_b64: str = "",
) -> dict[str, Any]:
"""
向应用群发文件:先上传临时素材,再 appchat/send。
与私聊 send_file 相同素材接口,但消息必须走群,不能 message/send。
失败回 ok=False,不抛给回调线程。
"""
if not chat_id:
return {"ok": False, "error": "chat_id"}
raw = b""
if filepath:
try:
from pathlib import Path
raw = Path(filepath).read_bytes()
except OSError as exc:
return {"ok": False, "error": f"read_file:{exc}"}
elif file_b64:
import base64
try:
raw = base64.b64decode(file_b64)
except Exception as exc: # noqa: BLE001
return {"ok": False, "error": f"b64:{exc}"}
if not raw:
return {"ok": False, "error": "file_empty"}
media_id, up = self._upload_media(filename=filename or "quote.xlsx", content=raw)
if not media_id:
return {"ok": False, "error": up.errmsg or "upload_fail"}
return self._post_appchat(
{
"chatid": chat_id,
"msgtype": "file",
"file": {"media_id": media_id},
}
)
def send_group_card(self, *, chat_id: str, template_card: dict[str, Any]) -> dict[str, Any]:
"""
向应用群发卡。企微 appchat 对 template_card 支持不稳定,失败由调用方已发过文本兜底。
"""
if not chat_id or not template_card:
return {"ok": False, "error": "chat_or_card"}
return self._post_appchat(
{
"chatid": chat_id,
"msgtype": "template_card",
"template_card": template_card,
}
)
def dismiss_group(self, *, chat_id: str, userids: list[str]) -> dict[str, Any]:
"""
解散应用群:官方无 dismiss,改为 appchat/update 移出成员。
不抛给回调。移完人不改工单六态。
"""
users = [str(x).strip() for x in (userids or []) if str(x).strip()]
if not chat_id or not users:
return {"ok": False, "error": "chat_or_users"}
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("解散群取 token 失败")
return {"ok": False, "error": str(exc)}
url = f"{self._api_base}/cgi-bin/appchat/update"
body = {"chatid": chat_id, "del_user_list": users}
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("企微解散群网络失败")
return {"ok": False, "error": str(exc)}
errcode = int(data.get("errcode") or 0)
if errcode != 0:
logger.warning("企微解散群失败 errcode=%s", errcode)
return {"ok": False, "error": str(data.get("errmsg") or f"errcode_{errcode}")}
logger.info("企微群已移出成员 chat=%s", chat_id)
return {"ok": True, "error": ""}
def _post_appchat(self, body: dict[str, Any]) -> dict[str, Any]:
"""appchat/send 公共发送。失败不抛。"""
if not self._corp_id or not self._secret:
return {"ok": False, "error": "WECOM_CORP_ID/SECRET 未配置"}
try:
token = self._get_token()
except Exception as exc: # noqa: BLE001
logger.exception("群发送取 token 失败")
return {"ok": False, "error": str(exc)}
url = f"{self._api_base}/cgi-bin/appchat/send"
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("企微 appchat/send 网络失败")
return {"ok": False, "error": str(exc)}
errcode = int(data.get("errcode") or 0)
if errcode != 0:
logger.warning("企微 appchat/send 失败 errcode=%s", errcode)
return {"ok": False, "error": str(data.get("errmsg") or f"errcode_{errcode}")}
return {"ok": True, "error": ""}
def send_text(self, *, touser: str, content: str) -> WeComSendResult:
"""
向指定企微 userid 发送文本应用消息。
参数:touser 必须是入站 sender_id;content 测试环境应由上层加 [test] 前缀。
副作用:出站 HTTP 调企微;失败返回 ok=False,不抛给回调线程(由 outbox 重试)。
"""
if not touser or not content:
return WeComSendResult(ok=False, errcode=-1, errmsg="touser/content 为空")
return self._post_message(
{
"touser": touser,
"msgtype": "text",
"agentid": self._agent_id,
"text": {"content": content},
"safe": 0,
},
preview=content,
)
def _post_message(self, body: dict[str, Any], *, preview: str) -> WeComSendResult:
"""真正打企微 message/send;密钥不写日志。"""
if not self._corp_id or not self._secret:
return WeComSendResult(ok=False, errcode=-1, errmsg="WECOM_CORP_ID/SECRET 未配置")
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"
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:
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)
shown = preview if len(preview) <= 200 else (preview[:200] + "…")
logger.info("企微发送成功 touser=%s reply=%s", body.get("touser"), shown)
return WeComSendResult(ok=True, errcode=0, errmsg=errmsg, raw=data)
def get_access_token(self) -> str:
"""
取出站/JS-SDK 用的 access_token。
只给同进程 H5 签名用;禁止写进页面或日志。
"""
return self._get_token()
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 过期类时也停)
}
)