"""提醒模块 Repository 层 负责系统提醒、客户欠款、不活跃客户检测、物流超时提醒 等相关数据的数据库查询与持久化操作。 被 ReminderService(提醒服务)调用。 """ 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 wechat_notification_service.send_subscribe_message( user.open_id, reminder.reminder_type, reminder.reminder_title, reminder.reminder_content, ) except Exception: pass 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())