dingdanquanliucheng/backend/app/repositories/reminder_repository.py
taiyi 1a3b67c423 feat: 微信订阅消息模板字段适配,填入 6 个真实模板 ID
- bootstrap.py 写入 6 个微信订阅消息模板 ID
- wechat_notification_service 新增 TEMPLATE_FIELD_MAP 按模板字段名发送
- reminder_repository 新增 _build_extra_data 从 title/content 提取结构化数据

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-06-06 13:25:10 +08:00

346 lines
13 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.

"""提醒模块 Repository 层
负责系统提醒、客户欠款、不活跃客户检测、物流超时提醒
等相关数据的数据库查询与持久化操作。
被 ReminderService提醒服务调用。
"""
import re
from datetime import datetime
from sqlalchemy import and_, func, select
from sqlalchemy.orm import Session
from backend.app.models.business import Customer, CustomerArrears, LogisticsTask, SalesOrder, SystemReminder
class ReminderRepository:
"""提醒数据仓储
封装系统提醒、客户欠款、不活跃客户检测、物流超时等查询。
被 ReminderService提醒服务调用。
"""
def list_reminders(self, session: Session, filters: dict) -> list[SystemReminder]:
"""查询系统提醒列表
根据提醒类型、状态、接收人等条件过滤提醒记录,
按 id 倒序返回。
参数:
session: 数据库会话
filters: 过滤条件字典,支持 reminder_type / status / receiver_user_id
返回:
符合条件的系统提醒列表
"""
stmt = select(SystemReminder)
if filters.get("reminder_type"):
stmt = stmt.where(SystemReminder.reminder_type == filters["reminder_type"])
if filters.get("status"):
stmt = stmt.where(SystemReminder.status == filters["status"])
if filters.get("receiver_user_id") is not None:
stmt = stmt.where(SystemReminder.receiver_user_id == filters["receiver_user_id"])
stmt = stmt.order_by(SystemReminder.id.desc())
return list(session.execute(stmt).scalars())
def get_reminder(self, session: Session, reminder_id: int) -> SystemReminder | None:
"""根据 ID 查询单条系统提醒
参数:
session: 数据库会话
reminder_id: 提醒记录 ID
返回:
对应的系统提醒对象,不存在则返回 None
"""
stmt = select(SystemReminder).where(SystemReminder.id == reminder_id)
return session.execute(stmt).scalar_one_or_none()
def get_customer(self, session: Session, customer_id: int) -> Customer | None:
"""根据 ID 查询未删除的客户信息
参数:
session: 数据库会话
customer_id: 客户 ID
返回:
对应的客户对象,不存在或已删除则返回 None
"""
stmt = select(Customer).where(Customer.id == customer_id, Customer.deleted == 0)
return session.execute(stmt).scalar_one_or_none()
def create_reminder(self, session: Session, payload: dict) -> SystemReminder:
"""创建系统提醒记录
创建提醒后,自动尝试发送微信订阅消息通知。
参数:
session: 数据库会话
payload: 提醒数据字典,包含 reminder_type / receiver_user_id 等字段
返回:
新创建的系统提醒对象
"""
reminder = SystemReminder(**payload)
session.add(reminder)
session.flush()
self._try_send_wechat_notification(session, reminder)
return reminder
def _try_send_wechat_notification(self, session: Session, reminder: SystemReminder) -> None:
"""尝试通过微信订阅消息发送提醒通知
根据接收人 ID 查询用户 open_id若存在则调用微信通知服务发送订阅消息。
出错时静默忽略,不影响提醒创建流程。
参数:
session: 数据库会话
reminder: 刚创建的系统提醒对象
"""
try:
from backend.app.services.wechat_notification_service import wechat_notification_service
from backend.app.models.system import User
from sqlalchemy import select
receiver_id = reminder.receiver_user_id
if receiver_id <= 0:
return
user = session.execute(select(User).where(User.id == receiver_id)).scalar_one_or_none()
if user is None or not user.open_id:
return
extra_data = self._build_extra_data(reminder)
wechat_notification_service.send_subscribe_message(
user.open_id,
reminder.reminder_type,
extra_data=extra_data,
)
except Exception:
pass
def _build_extra_data(self, reminder: SystemReminder) -> dict:
"""根据提醒类型从 title/content 中提取结构化数据,用于填充微信模板字段。
参数:
reminder: 系统提醒对象
返回:
包含模板所需字段数据的字典
"""
title = reminder.reminder_title or ""
content = reminder.reminder_content or ""
now_str = datetime.now().strftime("%Y-%m-%d %H:%M")
order_no_match = re.search(r'[\d-]{10,}', title + content)
order_no = order_no_match.group() if order_no_match else title
base: dict[str, str] = {
"order_no": order_no,
"change_time": now_str,
"assign_time": now_str,
"timeout_time": now_str,
"check_time": now_str,
"submit_time": now_str,
}
rtype = reminder.reminder_type
if rtype == "order_status_change":
base["status_text"] = content[:50] if content else title
elif rtype == "order_approval_needed":
base["submitter"] = "业务员"
base["order_info"] = title
base["remark"] = "请尽快审批"
elif rtype == "task_assigned":
base["task_no"] = order_no
base["task_type"] = "运输任务"
base["assign_date"] = now_str.split(" ")[0]
elif rtype == "logistics_timeout":
hours_match = re.search(r'(\d+)\s*小时', content)
base["timeout_info"] = f"超过{hours_match.group(1) if hours_match else '48'}小时未流转"
base["remark"] = "请及时跟进"
elif rtype == "arrears":
amount_match = re.search(r'([\d.]+)\s*元', content)
customer_match = re.search(r'客户[:](.+?)[,]', content)
base["customer_info"] = (customer_match.group(1) if customer_match else "") + "欠款"
base["arrears_amount"] = amount_match.group(1) if amount_match else "0"
base["order_date"] = now_str
base["due_date"] = ""
elif rtype == "inactive_customer":
days_match = re.search(r'(\d+)\s*天', content)
last_time_match = re.search(r'最近订单时间[:](.+?)[。,]', content)
base["customer_name"] = title.replace("沉默客户提醒 - ", "")
base["inactive_info"] = f"超过{days_match.group(1) if days_match else '90'}天未下单"
base["last_order_time"] = last_time_match.group(1).strip() if last_time_match else "暂无"
return base
def find_active_reminder(
self,
session: Session,
reminder_type: str,
biz_type: str,
biz_id: int,
receiver_user_id: int,
) -> SystemReminder | None:
"""查找针对同一业务对象的活跃提醒
根据提醒类型、业务类型、业务ID和接收人查找状态为
pending/sent/read 的现有提醒,用于防止重复创建提醒。
参数:
session: 数据库会话
reminder_type: 提醒类型
biz_type: 业务类型(如 order / arrears 等)
biz_id: 业务对象 ID
receiver_user_id: 接收人用户 ID
返回:
匹配的活跃提醒对象,不存在则返回 None
"""
stmt = select(SystemReminder).where(
SystemReminder.reminder_type == reminder_type,
SystemReminder.biz_type == biz_type,
SystemReminder.biz_id == biz_id,
SystemReminder.receiver_user_id == receiver_user_id,
SystemReminder.status.in_(["pending", "sent", "read"]),
)
return session.execute(stmt).scalar_one_or_none()
def list_overdue_arrears(self, session: Session) -> list[tuple[CustomerArrears, Customer | None, SalesOrder | None]]:
"""查询逾期欠款列表
联表查询状态为 pending 或 overdue 的欠款记录,
关联客户和销售订单信息。
参数:
session: 数据库会话
返回:
(欠款记录, 客户对象, 销售订单对象) 的元组列表
"""
stmt = (
select(CustomerArrears, Customer, SalesOrder)
.outerjoin(Customer, Customer.id == CustomerArrears.customer_id)
.outerjoin(SalesOrder, SalesOrder.id == CustomerArrears.order_id)
.where(CustomerArrears.status.in_(["pending", "overdue"]))
.order_by(CustomerArrears.id.asc())
)
return list(session.execute(stmt).all())
def get_arrears_by_order_id(self, session: Session, order_id: int) -> CustomerArrears | None:
"""根据订单 ID 查询欠款记录
参数:
session: 数据库会话
order_id: 关联的销售订单 ID
返回:
对应的欠款记录,不存在则返回 None
"""
stmt = select(CustomerArrears).where(CustomerArrears.order_id == order_id)
return session.execute(stmt).scalar_one_or_none()
def create_arrears(self, session: Session, payload: dict) -> CustomerArrears:
"""创建客户欠款记录
参数:
session: 数据库会话
payload: 欠款数据字典,包含 order_id / customer_id / amount 等字段
返回:
新创建的欠款记录对象
"""
arrears = CustomerArrears(**payload)
session.add(arrears)
session.flush()
return arrears
def list_inactive_customers(
self,
session: Session,
inactive_before: datetime,
) -> list[tuple[Customer, SalesOrder | None, float]]:
"""查询不活跃客户列表
遍历所有未删除客户,计算其累计订单金额并查找最近一笔订单。
若最近订单时间早于或等于指定时间阈值,则视为不活跃客户。
参数:
session: 数据库会话
inactive_before: 不活跃截止时间,最近订单晚于此时间的客户视为活跃
返回:
(客户对象, 最近订单对象, 累计订单金额) 的元组列表
"""
customers = list(
session.execute(
select(Customer)
.where(Customer.deleted == 0)
.order_by(Customer.id.asc())
).scalars()
)
results: list[tuple[Customer, SalesOrder | None, float]] = []
for customer in customers:
total_amount_stmt = (
select(func.coalesce(func.sum(SalesOrder.sale_price_total), 0))
.where(
SalesOrder.deleted == 0,
SalesOrder.customer_id == customer.id,
SalesOrder.order_status != "canceled",
)
)
total_amount = float(session.execute(total_amount_stmt).scalar() or 0)
latest_order_stmt = (
select(SalesOrder)
.where(
SalesOrder.deleted == 0,
SalesOrder.customer_id == customer.id,
SalesOrder.order_status != "canceled",
)
.order_by(SalesOrder.created_at.desc(), SalesOrder.id.desc())
)
latest_order = session.execute(latest_order_stmt).scalars().first()
if latest_order is None or (latest_order.created_at and latest_order.created_at <= inactive_before):
results.append((customer, latest_order, total_amount))
return results
def list_logistics_timeout_candidates(
self,
session: Session,
timeout_before: datetime,
) -> list[SalesOrder]:
"""查询物流超时待提醒的销售订单
联表查询物流任务状态为 pending 或 accepted
且创建时间早于指定超时阈值的销售订单,用于生成物流超时提醒。
参数:
session: 数据库会话
timeout_before: 超时截止时间,物流任务创建时间早于此值的订单视为超时
返回:
超时候选销售订单列表(去重)
"""
stmt = (
select(SalesOrder)
.join(LogisticsTask, LogisticsTask.order_id == SalesOrder.id)
.where(
SalesOrder.deleted == 0,
LogisticsTask.status.in_(["pending", "accepted"]),
LogisticsTask.created_at <= timeout_before,
)
.distinct()
.order_by(SalesOrder.id.asc())
)
return list(session.execute(stmt).scalars())