dingdanquanliucheng/backend/app/services/event_bus.py

432 lines
17 KiB
Python

"""提醒事件总线模块
统一的提醒触发入口,负责:
1. 创建 system_reminder 记录
2. 发送实时通知到 WebSocket 客户端
3. 调用微信订阅消息 API
使用方式:
from backend.app.services.event_bus import event_bus
# 在订单创建后触发(同步调用)
event_bus.emit("order_created", {
"order_id": order.id,
"order_no": order.order_no,
"salesman_id": order.salesman_id
})
支持的事件类型:
- order_created: 订单创建(通知管理层)
- order_approved: 订单审批通过(通知业务员)
- order_rejected: 订单审批拒绝(通知业务员)
- order_shipped: 订单发厂(通知业务员)
- task_assigned: 物流任务分配(通知司机)
- logistics_updated: 物流状态更新(通知业务员)
- arrears_reminder: 欠款提醒(通知管理层)
- logistics_timeout: 物流超时(通知管理层)
- inactive_customer: 沉默客户(通知业务员)
"""
import asyncio
import logging
import threading
from datetime import datetime
from typing import Any, Dict, List, Optional
from sqlalchemy.orm import Session
from backend.app.core.pubsub import reminder_pubsub
from backend.app.db import SessionLocal
from backend.app.models.business import SystemReminder
from backend.app.models.system import User, Role
from backend.app.repositories.reminder_repository import ReminderRepository
logger = logging.getLogger(__name__)
def _run_async(coro):
"""在后台线程中运行异步任务。
每个线程创建独立的事件循环,避免与主线程的事件循环冲突。
"""
def _wrapper():
try:
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
try:
loop.run_until_complete(coro)
finally:
loop.close()
except Exception as e:
logger.error(f"异步任务执行失败: {e}")
thread = threading.Thread(target=_wrapper, daemon=True)
thread.start()
class EventBus:
"""提醒事件总线"""
# 事件类型 -> (提醒类型, 标题模板, 内容模板, 通知对象)
EVENT_CONFIG = {
"order_created": {
"type": "order_approval_needed",
"title_template": "新订单待审批 - {order_no}",
"content_template": "业务员 {salesman_name} 创建了新订单,客户:{customer_name},金额:{amount}",
"notify_roles": ["admin", "manager"], # 通知管理层
},
"order_approved": {
"type": "order_status_change",
"title_template": "订单审批通过 - {order_no}",
"content_template": "您的订单已通过审批,等待发厂",
"notify_user_ids": "salesman_id", # 通知指定用户
},
"order_rejected": {
"type": "order_status_change",
"title_template": "订单审批拒绝 - {order_no}",
"content_template": "您的订单未通过审批,原因:{reason}",
"notify_user_ids": "salesman_id",
},
"order_shipped": {
"type": "order_status_change",
"title_template": "订单已发厂 - {order_no}",
"content_template": "订单已确认发厂,供应商:{supplier_name}",
"notify_user_ids": "salesman_id",
},
"task_assigned": {
"type": "task_assigned",
"title_template": "新物流任务 - 订单 {order_no}",
"content_template": "您有新的物流任务,请及时处理",
"notify_user_ids": "driver_id",
},
"logistics_updated": {
"type": "order_status_change",
"title_template": "物流状态更新 - {order_no}",
"content_template": "订单物流状态已更新:{status}",
"notify_user_ids": "salesman_id", # 只通知业务员
},
"arrears_reminder": {
"type": "arrears",
"title_template": "客户欠款逾期 - {customer_name}",
"content_template": "客户累计逾期欠款 {amount}元,请及时跟进",
"notify_user_ids": "receiver_user_id",
},
"logistics_timeout": {
"type": "logistics_timeout",
"title_template": "物流超时提醒 - {order_no}",
"content_template": "订单 {order_no} 超时未完成物流流转,请及时跟进",
"notify_user_ids": "receiver_user_id",
},
"inactive_customer": {
"type": "inactive_customer",
"title_template": "沉默客户提醒 - {customer_name}",
"content_template": "客户超过 {days} 天未下单,请关注",
"notify_user_ids": "receiver_user_id",
},
"order_status_changed": {
"type": "order_status_change",
"title_template": "订单状态变更 - {order_no}",
"content_template": "订单 {order_no} 状态已变更",
"notify_user_ids": "salesman_id",
"also_notify": "operator_id",
},
}
def __init__(self):
self._reminder_repo = ReminderRepository()
self.EVENT_CONFIG["order_status_changed"]["title_template"] = "\u8ba2\u5355\u72b6\u6001\u53d8\u66f4 - {order_no}"
self.EVENT_CONFIG["order_status_changed"]["content_template"] = (
"\u8ba2\u5355 {order_no} \u72b6\u6001\u7531\u300c{before_status_text}\u300d"
"\u53d8\u66f4\u4e3a\u300c{after_status_text}\u300d\u3002"
)
def emit(
self,
event_type: str,
data: Dict[str, Any],
session: Optional[Session] = None
) -> List[int]:
"""触发提醒事件(同步方法)
Args:
event_type: 事件类型
data: 事件数据
session: 数据库会话(可选,不传则自动创建)
Returns:
创建的提醒 ID 列表
"""
if event_type not in self.EVENT_CONFIG:
logger.warning(f"未知的事件类型: {event_type}")
return []
config = self.EVENT_CONFIG[event_type]
reminder_ids = []
# 格式化标题和内容
try:
title = config["title_template"].format(**data)
content = config["content_template"].format(**data)
except KeyError as e:
logger.error(f"事件数据缺少必要字段: {e}")
return []
# 确定通知对象
notify_user_ids = self._resolve_notify_targets(config, data, session)
if not notify_user_ids:
logger.warning(f"事件 {event_type} 没有找到通知对象")
return []
# 为每个通知对象创建提醒
use_own_session = session is None
if use_own_session:
session = SessionLocal()
try:
for user_id in notify_user_ids:
# 创建提醒记录
reminder = self._reminder_repo.create_reminder(
session,
{
"reminder_type": config["type"],
"biz_type": data.get("biz_type", ""),
"biz_id": data.get("biz_id", 0),
"receiver_user_id": user_id,
"reminder_title": title,
"reminder_content": content,
"status": "pending",
"sent_at": datetime.now(),
},
skip_wechat=True,
)
if reminder:
reminder_ids.append(reminder.id)
# 构建提醒数据
reminder_data = {
"reminder_id": reminder.id,
"title": title,
"content": content,
"type": config["type"],
"biz_type": data.get("biz_type", ""),
"biz_id": data.get("biz_id", 0),
"status": "pending",
"created_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
}
# 在后台线程中异步发送实时通知(不阻塞主流程)
_run_async(self._send_realtime_notification(user_id, reminder_data))
# 在后台线程中异步发送微信订阅消息
wechat_data = {**data, "title": title, "content": content}
_run_async(self._send_wechat_notification(user_id, event_type, wechat_data))
if use_own_session:
session.commit()
logger.info(f"事件 {event_type} 触发成功,创建 {len(reminder_ids)} 条提醒")
except Exception as e:
logger.error(f"事件 {event_type} 处理失败: {e}")
if use_own_session:
session.rollback()
raise
finally:
if use_own_session:
session.close()
return reminder_ids
def _resolve_notify_targets(
self,
config: Dict[str, Any],
data: Dict[str, Any],
session: Optional[Session]
) -> List[int]:
"""解析通知目标用户 ID 列表"""
user_ids = []
# 直接指定的用户 ID
if "notify_user_ids" in config:
key = config["notify_user_ids"]
if key in data:
uid = data[key]
if uid and uid > 0:
user_ids.append(uid)
# 按角色通知
if "notify_roles" in config:
if session is None:
session = SessionLocal()
should_close = True
else:
should_close = False
try:
for role_code in config["notify_roles"]:
# 查找角色ID
role = session.query(Role).filter(
Role.role_code == role_code,
Role.status == 1
).first()
if role:
# 查找该角色下的所有启用用户
users = session.query(User).filter(
User.role_id == role.id,
User.status == 1
).all()
for user in users:
if user.id not in user_ids:
user_ids.append(user.id)
finally:
if should_close:
session.close()
# 额外通知的用户
if "also_notify" in config:
key = config["also_notify"]
if key in data:
uid = data[key]
if uid and uid > 0 and uid not in user_ids:
user_ids.append(uid)
return user_ids
async def _send_realtime_notification(self, user_id: int, reminder: dict) -> None:
"""发送实时通知到 WebSocket。
使用独立的 Redis 连接,避免跨事件循环复用连接导致错误。
"""
try:
from backend.app.core.pubsub import ReminderPubSub
# 每次创建新实例,确保连接属于当前事件循环
pubsub = ReminderPubSub()
await pubsub.publish(user_id, reminder)
await pubsub.close()
except Exception as e:
logger.error(f"发送实时通知失败: {e}")
async def _send_wechat_notification(
self,
user_id: int,
event_type: str,
data: Dict[str, Any]
) -> None:
# Send WeChat notification. Prefer service account, fallback to mini program.
try:
template_map = {
"order_created": "order_approval_needed",
"order_approved": "order_status_change",
"order_rejected": "order_status_change",
"order_shipped": "order_status_change",
"task_assigned": "task_assigned",
"logistics_updated": "order_status_change",
"arrears_reminder": "arrears",
"logistics_timeout": "logistics_timeout",
"inactive_customer": "inactive_customer",
"order_status_changed": "order_status_change",
}
template_type = template_map.get(event_type)
if not template_type:
logger.warning("Event %s has no template mapping, skipping WeChat", event_type)
return
session = SessionLocal()
try:
user = session.query(User).filter(User.id == user_id).first()
if not user:
return
data = self._build_wechat_data(session, event_type, data)
# Prefer service account template message
if user.service_open_id:
from backend.app.services.service_account_notification_service import service_account_notification_service
sent = await service_account_notification_service.send_reminder(
open_id=user.service_open_id,
reminder_type=template_type,
data=data,
)
if sent:
logger.info("WeChat service account reminder sent to user %s for event %s", user_id, event_type)
return
logger.warning(
"WeChat service account reminder failed for user %s event %s, fallback_to_mini_program=%s",
user_id,
event_type,
bool(user.open_id),
)
else:
logger.info("User %s has no service_open_id for event %s", user_id, event_type)
# Fallback to mini program subscribe message
if user.open_id:
from backend.app.services.wechat_notification_service import wechat_notification_service
sent = await wechat_notification_service.send_reminder(
open_id=user.open_id,
reminder_type=template_type,
data=data,
)
if sent:
logger.info("WeChat mini program reminder sent to user %s for event %s", user_id, event_type)
else:
logger.warning("WeChat mini program reminder failed for user %s event %s", user_id, event_type)
return
logger.warning(
"Skipping WeChat reminder for user %s event %s because no service_open_id/open_id is bound",
user_id,
event_type,
)
finally:
session.close()
except Exception as e:
logger.warning("WeChat notification failed for user %s event %s: %s", user_id, event_type, e)
def _build_wechat_data(self, session: Session, event_type: str, data: Dict[str, Any]) -> Dict[str, Any]:
payload = dict(data)
if event_type == "order_created":
payload.setdefault("approval_type", "订单审批")
payload.setdefault("order_status", "待审批")
payload.setdefault("task_name", payload.get("title", ""))
elif event_type in {"order_approved", "order_rejected", "order_shipped", "order_status_changed"}:
payload.setdefault("status_text", payload.get("after_status_text") or payload.get("title", ""))
payload.setdefault("status_desc", payload.get("after_status_text") or payload.get("status_desc", ""))
payload.setdefault("status", payload.get("after_status_text") or payload.get("status", ""))
payload.setdefault("current_location", payload.get("content", ""))
payload.setdefault("goods_info", payload.get("content", ""))
elif event_type == "task_assigned":
from backend.app.models.business import LogisticsTask, SalesOrder
task = session.get(LogisticsTask, payload.get("task_id") or 0)
order = session.get(SalesOrder, task.order_id) if task else None
driver = session.query(User).filter(User.id == task.driver_id).first() if task else None
payload.setdefault("driver_name", driver.real_name if driver else "")
payload.setdefault("driver_phone", driver.mobile if driver else "")
payload.setdefault("tracking_number", task.tracking_number if task and task.tracking_number else "")
payload.setdefault("order_no", order.order_no if order else payload.get("order_no", ""))
payload.setdefault("delivery_address", task.delivery_address if task else "")
elif event_type == "logistics_updated":
from backend.app.models.business import LogisticsTask, SalesOrder
task = session.get(LogisticsTask, payload.get("task_id") or 0)
order = session.get(SalesOrder, payload.get("order_id") or 0)
payload.setdefault("tracking_number", task.tracking_number if task and task.tracking_number else (order.tracking_number if order else ""))
payload.setdefault("status_text", payload.get("status_desc") or payload.get("status", ""))
payload.setdefault("current_location", payload.get("content", ""))
payload.setdefault("goods_info", payload.get("content", ""))
return payload
# 全局单例
event_bus = EventBus()