252 lines
11 KiB
Python
252 lines
11 KiB
Python
"""微信服务号模板消息通知服务。
|
|
|
|
通过微信 templateMessage.send 接口向关注服务号的用户推送服务提醒。
|
|
与小程序订阅消息不同,服务号模板消息无单次授权限制,用户关注后可持续接收。
|
|
"""
|
|
|
|
import json
|
|
import logging
|
|
import time
|
|
from urllib import error, request
|
|
|
|
from backend.app.core.config import get_settings
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class ServiceAccountNotificationService:
|
|
"""微信服务号模板消息通知服务。"""
|
|
|
|
TOKEN_CACHE: dict[str, str] = {"token": "", "expires_at": 0}
|
|
|
|
REMINDER_TYPE_TEMPLATE_MAP: dict[str, str] = {
|
|
"order_status_change": "wechat_sa_template_order_status",
|
|
"order_approval_needed": "wechat_sa_template_approval",
|
|
"task_assigned": "wechat_sa_template_task_assigned",
|
|
"logistics_timeout": "wechat_sa_template_logistics_timeout",
|
|
"arrears": "wechat_sa_template_arrears",
|
|
"inactive_customer": "wechat_sa_template_inactive_customer",
|
|
}
|
|
|
|
# 服务号模板字段映射(与小程序订阅消息字段名完全不同)
|
|
TEMPLATE_FIELD_MAP: dict[str, dict[str, str]] = {
|
|
"order_status_change": {
|
|
"thing2": "status_text",
|
|
"character_string3": "order_no",
|
|
"phrase6": "logistics_status",
|
|
"time7": "update_time",
|
|
"thing13": "current_location",
|
|
},
|
|
"order_approval_needed": {
|
|
"character_string2": "order_no",
|
|
"thing8": "approval_type",
|
|
"thing7": "submitter",
|
|
"time15": "submit_time",
|
|
"phrase3": "order_status",
|
|
},
|
|
"task_assigned": {
|
|
"thing1": "driver_name",
|
|
"phone_number2": "driver_phone",
|
|
"character_string3": "tracking_number",
|
|
"character_string4": "order_no",
|
|
},
|
|
"logistics_timeout": {
|
|
"character_string2": "tracking_number",
|
|
"thing11": "latest_location",
|
|
"time7": "update_time",
|
|
"character_string9": "order_no",
|
|
"thing5": "destination",
|
|
},
|
|
"arrears": {
|
|
"amount2": "pending_amount",
|
|
"time4": "due_time",
|
|
"thing6": "debtor_name",
|
|
},
|
|
"inactive_customer": {
|
|
"thing2": "customer_name",
|
|
"time3": "visit_time",
|
|
"phone_number8": "contact_phone",
|
|
},
|
|
}
|
|
|
|
def send_template_message(self, open_id: str, reminder_type: str, title: str = "", extra_data: dict | None = None, biz_url: str = "") -> bool:
|
|
"""发送服务号模板消息。"""
|
|
settings = get_settings()
|
|
if not settings.wechat_sa_app_id or not settings.wechat_sa_app_secret:
|
|
return False
|
|
|
|
template_id = self._get_template_id(settings, reminder_type)
|
|
if not template_id:
|
|
return False
|
|
|
|
access_token = self._get_access_token(settings)
|
|
if not access_token:
|
|
return False
|
|
|
|
field_map = self.TEMPLATE_FIELD_MAP.get(reminder_type)
|
|
if not field_map:
|
|
return False
|
|
|
|
payload_data = self._build_payload_data(field_map, title, extra_data or {})
|
|
|
|
payload = {
|
|
"touser": open_id,
|
|
"template_id": template_id,
|
|
"url": biz_url,
|
|
"data": payload_data,
|
|
}
|
|
|
|
url = f"https://api.weixin.qq.com/cgi-bin/message/template/send?access_token={access_token}"
|
|
req = request.Request(
|
|
url=url,
|
|
data=json.dumps(payload, ensure_ascii=False).encode("utf-8"),
|
|
headers={"Content-Type": "application/json; charset=UTF-8"},
|
|
method="POST",
|
|
)
|
|
try:
|
|
with request.urlopen(req, timeout=10) as response:
|
|
result = json.loads(response.read().decode("utf-8") or "{}")
|
|
errcode = result.get("errcode", -1)
|
|
if errcode != 0:
|
|
logger.warning("服务号模板消息发送失败: errcode=%s, errmsg=%s, open_id=%s, type=%s",
|
|
errcode, result.get("errmsg", ""), open_id, reminder_type)
|
|
return errcode == 0
|
|
except (error.HTTPError, error.URLError, json.JSONDecodeError) as e:
|
|
logger.warning("服务号模板消息发送异常: %s, open_id=%s", e, open_id)
|
|
return False
|
|
|
|
async def send_reminder(self, open_id: str, reminder_type: str, data: dict | None = None) -> bool:
|
|
"""发送提醒消息(兼容 event_bus 调用)。"""
|
|
extra_data = {}
|
|
if data:
|
|
extra_data = {
|
|
"title": data.get("title", ""),
|
|
"order_no": data.get("order_no", ""),
|
|
"status_text": data.get("status_text", data.get("status_desc", data.get("status", ""))),
|
|
"status_desc": data.get("status_desc", data.get("status", "")),
|
|
"status": data.get("status", data.get("status_desc", "")),
|
|
"change_time": data.get("change_time", time.strftime("%Y-%m-%d %H:%M")),
|
|
"remark": data.get("remark", data.get("content", "")),
|
|
"submitter": data.get("salesman_name", data.get("submitter", "")),
|
|
"task_name": data.get("task_name", data.get("task_no", data.get("title", ""))),
|
|
"submit_time": data.get("submit_time", time.strftime("%Y-%m-%d %H:%M")),
|
|
"assign_time": data.get("assign_time", time.strftime("%Y-%m-%d %H:%M")),
|
|
"order_status": data.get("order_status", data.get("status_desc", data.get("status", ""))),
|
|
"deadline": data.get("deadline", ""),
|
|
"pending_amount": data.get("amount", ""),
|
|
"order_time": data.get("order_time", ""),
|
|
"activity_name": data.get("customer_name", ""),
|
|
"activity_time": data.get("activity_time", time.strftime("%Y-%m-%d %H:%M")),
|
|
}
|
|
|
|
title = data.get("title", "") if data else ""
|
|
biz_url = data.get("biz_url", "") if data else ""
|
|
|
|
return self.send_template_message(
|
|
open_id=open_id,
|
|
reminder_type=reminder_type,
|
|
title=title,
|
|
extra_data=self._build_extra_data(reminder_type, extra_data, title),
|
|
biz_url=biz_url,
|
|
)
|
|
|
|
def _build_payload_data(self, field_map: dict[str, str], title: str, data: dict) -> dict[str, dict[str, str]]:
|
|
now_str = time.strftime("%Y-%m-%d %H:%M")
|
|
payload_data: dict[str, dict[str, str]] = {}
|
|
for wx_field, data_key in field_map.items():
|
|
value = data.get(data_key, "")
|
|
if not value:
|
|
value = self._default_value(data_key, title, now_str)
|
|
payload_data[wx_field] = {"value": str(value)[:100], "color": "#173177"}
|
|
return payload_data
|
|
|
|
def _build_extra_data(self, reminder_type: str, data: dict, title: str) -> dict:
|
|
now_str = time.strftime("%Y-%m-%d %H:%M")
|
|
extra_data = dict(data)
|
|
extra_data.setdefault("order_no", data.get("task_no", ""))
|
|
|
|
if reminder_type == "order_approval_needed":
|
|
extra_data.setdefault("approval_type", "订单审批")
|
|
extra_data.setdefault("submitter", data.get("salesman_name", data.get("submitter", "")))
|
|
extra_data.setdefault("submit_time", data.get("submit_time", now_str))
|
|
extra_data.setdefault("order_status", data.get("order_status", "待审批"))
|
|
|
|
elif reminder_type == "task_assigned":
|
|
extra_data.setdefault("driver_name", data.get("driver_name", ""))
|
|
extra_data.setdefault("driver_phone", data.get("driver_phone", ""))
|
|
extra_data.setdefault("tracking_number", data.get("tracking_number", ""))
|
|
|
|
elif reminder_type == "order_status_change":
|
|
extra_data.setdefault("status_text", data.get("status_text", title or data.get("status_desc", "")))
|
|
extra_data.setdefault(
|
|
"logistics_status",
|
|
data.get("logistics_status")
|
|
or data.get("status_text")
|
|
or data.get("status_desc")
|
|
or data.get("status", ""),
|
|
)
|
|
extra_data.setdefault("update_time", data.get("change_time", data.get("update_time", now_str)))
|
|
extra_data.setdefault("current_location", data.get("current_location", data.get("remark", data.get("content", ""))))
|
|
|
|
elif reminder_type == "logistics_timeout":
|
|
extra_data.setdefault("tracking_number", data.get("tracking_number", data.get("order_no", "")))
|
|
extra_data.setdefault("latest_location", data.get("latest_location", data.get("location", data.get("remark", data.get("content", "")))))
|
|
extra_data.setdefault("update_time", data.get("update_time", data.get("deadline", now_str)))
|
|
extra_data.setdefault("destination", data.get("destination", data.get("delivery_address", "")))
|
|
|
|
elif reminder_type == "arrears":
|
|
extra_data.setdefault("pending_amount", data.get("amount", "0"))
|
|
extra_data.setdefault("due_time", data.get("due_time", data.get("deadline", now_str)))
|
|
extra_data.setdefault("debtor_name", data.get("debtor_name", data.get("customer_name", "")))
|
|
|
|
elif reminder_type == "inactive_customer":
|
|
extra_data.setdefault("customer_name", data.get("customer_name", ""))
|
|
extra_data.setdefault("visit_time", data.get("visit_time", data.get("activity_time", now_str)))
|
|
extra_data.setdefault("contact_phone", data.get("contact_phone", data.get("customer_mobile", "")))
|
|
|
|
return extra_data
|
|
|
|
@staticmethod
|
|
def _default_value(data_key: str, title: str, now_str: str) -> str:
|
|
if "time" in data_key or "date" in data_key:
|
|
return now_str
|
|
if "amount" in data_key:
|
|
return "0.00"
|
|
if "phone" in data_key:
|
|
return "-"
|
|
return title or ""
|
|
|
|
def _get_template_id(self, settings, reminder_type: str) -> str:
|
|
config_key = self.REMINDER_TYPE_TEMPLATE_MAP.get(reminder_type, "")
|
|
if not config_key:
|
|
return ""
|
|
return getattr(settings, config_key, "")
|
|
|
|
def _get_access_token(self, settings) -> str:
|
|
now = time.time()
|
|
if self.TOKEN_CACHE["token"] and self.TOKEN_CACHE["expires_at"] > now + 60:
|
|
return self.TOKEN_CACHE["token"]
|
|
|
|
url = (
|
|
f"https://api.weixin.qq.com/cgi-bin/token"
|
|
f"?grant_type=client_credential"
|
|
f"&appid={settings.wechat_sa_app_id}"
|
|
f"&secret={settings.wechat_sa_app_secret}"
|
|
)
|
|
req = request.Request(url=url, method="GET")
|
|
try:
|
|
with request.urlopen(req, timeout=10) as response:
|
|
result = json.loads(response.read().decode("utf-8") or "{}")
|
|
token = result.get("access_token", "")
|
|
expires_in = int(result.get("expires_in", 0))
|
|
if token and expires_in > 0:
|
|
self.TOKEN_CACHE["token"] = token
|
|
self.TOKEN_CACHE["expires_at"] = now + expires_in
|
|
return token
|
|
except (error.HTTPError, error.URLError, json.JSONDecodeError) as e:
|
|
logger.warning("服务号 access_token 获取失败: %s", e)
|
|
return ""
|
|
|
|
|
|
service_account_notification_service = ServiceAccountNotificationService()
|