"""物流管理服务层。 负责司机运输任务的创建、状态流转、物流轨迹查询等业务逻辑。 支持内部轨迹和第三方物流(快递100、通用 APPCODE)轨迹查询。 依赖 LogisticsRepository、OrderRepository 进行数据持久化, 依赖 audit_service 记录操作审计日志,依赖 arrears_service 同步欠款数据, 依赖 file_service 管理司机任务附件。 """ from datetime import datetime import json import logging from urllib import error, request from sqlalchemy.exc import SQLAlchemyError from sqlalchemy.orm import Session from backend.app.core.config import get_settings from backend.app.core.cache import cache_delete_pattern, cache_get, cache_set, make_cache_key from backend.app.core.error_codes import ErrorCode from backend.app.core.exceptions import AppException from backend.app.repositories.logistics_repository import LogisticsRepository from backend.app.repositories.order_repository import OrderRepository from backend.app.repositories.reminder_repository import ReminderRepository from backend.app.services.arrears_service import arrears_service from backend.app.services.audit_service import audit_service from backend.app.services.event_bus import event_bus from backend.app.services.file_service import file_service logger = logging.getLogger(__name__) class LogisticsService: """物流服务,封装司机任务和物流轨迹的业务逻辑。 依赖: LogisticsRepository - 物流任务和轨迹数据访问 OrderRepository - 订单数据访问(任务创建/取消时联动订单状态) ReminderRepository - 任务分配提醒通知 arrears_service - 欠款同步 audit_service - 操作审计日志 file_service - 司机任务附件管理 settings - 第三方物流 API 配置(快递100等) """ def __init__(self) -> None: self.settings = get_settings() self.logistics_repository = LogisticsRepository() self.order_repository = OrderRepository() self.reminder_repository = ReminderRepository() self.trace_records: dict[int, list] = {} def list_tasks( self, session: Session | None = None, filters: dict | None = None, current_user: dict | None = None, ) -> dict: """查询物流任务列表。 Args: session: 数据库会话 filters: 查询过滤条件 current_user: 当前登录用户信息,司机角色自动过滤为本人任务 Returns: 包含 total、page_no、page_size、list 的分页字典 被调用路由: logistics.py - GET /logistics/tasks """ if session is not None: cache_key = make_cache_key("logistics:list", user_id=(current_user or {}).get("user_id"), **(filters or {})) cached = cache_get(cache_key) if cached is not None: return cached try: tasks = self.logistics_repository.list_tasks_by_filters( session, self.normalize_task_filters(filters, current_user), ) result = { "total": len(tasks), "page_no": 1, "page_size": 20, "list": [self._map_task_summary(task, session) for task in tasks], } cache_set(cache_key, result, ttl=30) return result except SQLAlchemyError as exc: raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc return {"total": 0, "page_no": 1, "page_size": 20, "list": []} def create_task( self, payload: dict, session: Session | None = None, current_user: dict | None = None, ) -> dict: """创建司机运输任务。 校验订单状态和任务唯一性,创建任务后订单状态变更为 pending_driver。 向司机发送任务分配提醒通知。自动记录审计日志。 Args: payload: 包含 order_id、driver_id、pickup_address 等的请求体 session: 数据库会话 current_user: 当前登录用户信息 Returns: 任务详情字典 被调用路由: logistics.py - POST /logistics/tasks """ if session is not None: try: order = self.order_repository.get_order(session, payload["order_id"]) if order is None: raise AppException(code=ErrorCode.NOT_FOUND, message="订单不存在", status_code=404) if order.order_status not in {"approved", "pending_logistics"}: raise AppException(code=ErrorCode.INVALID_STATUS, message="当前状态不允许创建司机任务", status_code=400) existing_task = self.logistics_repository.get_task_by_order_id(session, order.id) if existing_task is not None and existing_task.status != "canceled": raise AppException(code=ErrorCode.PARAM_ERROR, message="该订单已存在有效司机任务", status_code=400) task = self.logistics_repository.create_task( session, { "task_no": f"LT{datetime.now().strftime('%Y%m%d%H%M%S')}", "order_id": order.id, "driver_id": payload["driver_id"], "factory_id": payload.get("factory_id") or order.factory_id, "pickup_address": payload["pickup_address"], "delivery_address": payload["delivery_address"], "pickup_content": payload["pickup_content"], "quantity": payload["quantity"], "status": "pending", "tracking_number": payload.get("tracking_number"), "express_company": payload.get("express_company"), "created_by": current_user.get("user_id") if current_user else None, "remark": payload.get("remark"), }, ) # 司机任务创建后,订单进入待司机接单阶段。 self.order_repository.update_order_status(session, order, "pending_driver") cache_delete_pattern(f"order:detail:{order.id}") # 触发任务分配事件,通知司机 event_bus.emit("task_assigned", { "task_id": task.id, "task_no": task.task_no, "order_id": order.id, "order_no": order.order_no, "driver_id": task.driver_id, "biz_type": "logistics_task", "biz_id": task.id, }, session) audit_service.write_log( session, { "operate_type": "logistics_task_create", "biz_type": "logistics_task", "biz_id": task.id, "before_value": None, "after_value": self._map_task_detail(task, session), "remark": f"创建司机任务 {task.task_no}", }, ) session.commit() cache_delete_pattern("logistics:*") return self._map_task_detail(task, session) except AppException: session.rollback() raise except SQLAlchemyError as exc: session.rollback() raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库连接不可用", status_code=500) def reassign_task( self, task_id: int, payload: dict, session: Session | None = None, current_user: dict | None = None, ) -> dict: """重新分配待接单司机任务,复用原任务内容。""" if session is not None: try: task = self.logistics_repository.get_task(session, task_id) if task is None: raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404) self._ensure_task_access(task, current_user) if task.status != "pending": raise AppException(code=ErrorCode.INVALID_STATUS, message="当前任务状态不允许重新分配司机", status_code=400) order = self.order_repository.get_order(session, task.order_id) if order is None: raise AppException(code=ErrorCode.NOT_FOUND, message="订单不存在", status_code=404) if order.order_status != "pending_driver": raise AppException(code=ErrorCode.INVALID_STATUS, message="当前订单状态不允许重新分配司机", status_code=400) new_driver_id = payload["driver_id"] if new_driver_id == task.driver_id: raise AppException(code=ErrorCode.PARAM_ERROR, message="新司机不能与当前司机相同", status_code=400) before_value = self._map_task_detail(task, session) self.logistics_repository.update_task_status(session, task, "canceled") task.canceled_at = datetime.now() task.canceled_by = current_user.get("user_id") if current_user else None task.cancel_reason = "管理员重新分配司机" session.add(task) new_task = self.logistics_repository.create_task( session, { "task_no": f"LT{datetime.now().strftime('%Y%m%d%H%M%S')}", "order_id": task.order_id, "driver_id": new_driver_id, "factory_id": task.factory_id or order.factory_id, "pickup_address": task.pickup_address, "delivery_address": task.delivery_address, "pickup_content": task.pickup_content, "quantity": task.quantity, "status": "pending", "tracking_number": None, "express_company": None, "created_by": current_user.get("user_id") if current_user else None, "remark": task.remark, }, ) cache_delete_pattern(f"order:detail:{order.id}") event_bus.emit( "task_assigned", { "task_id": new_task.id, "task_no": new_task.task_no, "order_id": order.id, "order_no": order.order_no, "driver_id": new_task.driver_id, "biz_type": "logistics_task", "biz_id": new_task.id, }, session, ) audit_service.write_log( session, { "operate_type": "logistics_task_reassign", "biz_type": "logistics_task", "biz_id": new_task.id, "before_value": before_value, "after_value": self._map_task_detail(new_task, session), "remark": f"重新分配司机:{task.task_no} -> {new_task.task_no}", }, ) session.commit() cache_delete_pattern("logistics:*") return { "task_id": new_task.id, "task_no": new_task.task_no, "old_task_id": task.id, "order_id": order.id, "status": new_task.status, } except AppException: session.rollback() raise except SQLAlchemyError as exc: session.rollback() raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库连接不可用", status_code=500) def get_task( self, task_id: int, session: Session | None = None, current_user: dict | None = None, ) -> dict: """查询单个物流任务详情。 Args: task_id: 任务 ID session: 数据库会话 current_user: 当前登录用户信息 Returns: 任务详情字典 被调用路由: logistics.py - GET /logistics/tasks/{task_id} """ if session is not None: try: task = self.logistics_repository.get_task(session, task_id) if task is None: raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404) self._ensure_task_access(task, current_user) return self._map_task_detail(task, session) except SQLAlchemyError as exc: raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404) def list_driver_tasks(self, session: Session | None = None, current_user: dict | None = None) -> dict: """查询当前司机的任务列表。 自动以当前司机 user_id 作为过滤条件。 Args: session: 数据库会话 current_user: 当前登录用户信息 Returns: 任务分页列表字典 被调用路由: logistics.py - GET /logistics/driver/tasks """ return self.list_tasks(session=session, filters={}, current_user=current_user) def cancel_task( self, task_id: int, payload: dict | None = None, session: Session | None = None, current_user: dict | None = None, ) -> dict: """取消司机任务。 仅待接单和已接单状态的任务允许取消。取消后若订单处于 pending_driver 状态, 则回退为 approved。自动记录审计日志。 Args: task_id: 任务 ID payload: 可选的请求体,包含 cancel_reason、remark session: 数据库会话 current_user: 当前登录用户信息 Returns: 取消结果字典 被调用路由: logistics.py - POST /logistics/tasks/{task_id}/cancel """ operate_payload = payload or {} if session is not None: try: task = self.logistics_repository.get_task(session, task_id) if task is None: raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404) self._ensure_task_access(task, current_user) if task.status not in {"pending", "accepted"}: raise AppException(code=ErrorCode.INVALID_STATUS, message="当前任务状态不允许取消", status_code=400) before_status = task.status self.logistics_repository.update_task_status(session, task, "canceled") task.canceled_at = datetime.now() task.canceled_by = current_user.get("user_id") if current_user else None task.cancel_reason = (operate_payload.get("cancel_reason") or "").strip() or None session.add(task) order = self.order_repository.get_order(session, task.order_id) if order is not None and order.order_status == "pending_driver": self.order_repository.update_order_status(session, order, "approved") cache_delete_pattern(f"order:detail:{order.id}") audit_service.write_log( session, { "operate_type": "logistics_task_cancel", "biz_type": "logistics_task", "biz_id": task.id, "before_value": {"task_status": before_status}, "after_value": { "task_status": task.status, "canceled_at": task.canceled_at.strftime("%Y-%m-%d %H:%M:%S") if task.canceled_at else None, "cancel_reason": task.cancel_reason, }, "remark": operate_payload.get("remark") or f"取消司机任务 {task.task_no}", }, ) session.commit() cache_delete_pattern("logistics:*") return {"task_id": task.id, "status": task.status, "order_id": task.order_id} except AppException: session.rollback() raise except SQLAlchemyError as exc: session.rollback() raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库连接不可用", status_code=500) def get_driver_task(self, task_id: int, session: Session | None = None, current_user: dict | None = None) -> dict: """查询单个司机任务详情(司机视角)。 委托给 get_task 实现,通过 _ensure_task_access 控制权限。 Args: task_id: 任务 ID session: 数据库会话 current_user: 当前登录用户信息 Returns: 任务详情字典 被调用路由: logistics.py - GET /logistics/driver/tasks/{task_id} """ return self.get_task(task_id, session=session, current_user=current_user) def update_tracking( self, task_id: int, payload: dict, session: Session | None = None, current_user: dict | None = None, ) -> dict: """修改物流单号。 业务员/管理员修改已有物流任务的快递单号和快递公司。 业务员只能修改自己订单关联的任务。 Args: task_id: 物流任务 ID payload: 包含 tracking_number、express_company、remark 的字典 session: 数据库会话 current_user: 当前登录用户信息 Returns: 更新后的任务信息字典 被调用路由: logistics.py - PUT /logistics/tasks/{task_id}/tracking """ if session is None: raise AppException(code=ErrorCode.PARAM_ERROR, message="数据库会话未初始化", status_code=500) task = self.logistics_repository.get_task(session, task_id) if task is None: raise AppException(code=ErrorCode.NOT_FOUND, message="物流任务不存在", status_code=404) # 权限校验:salesman 只能改自己订单的物流单号 if current_user and current_user.get("role_code") == "salesman": order = self.order_repository.get_order(session, task.order_id) if order and order.salesman_id != current_user.get("user_id"): raise AppException(code=ErrorCode.FORBIDDEN, message="无权修改他人订单的物流单号", status_code=403) new_tracking = (payload.get("tracking_number") or "").strip() if not new_tracking: raise AppException(code=ErrorCode.PARAM_ERROR, message="物流单号不能为空", status_code=400) old_tracking = task.tracking_number or "" old_express = task.express_company or "" task.tracking_number = new_tracking # 快递公司:优先使用传入值,为空时根据单号自动识别 express_company = (payload.get("express_company") or "").strip() if not express_company: auto_code, _ = self._guess_express_by_tracking(new_tracking) if auto_code: express_company = auto_code task.express_company = express_company or None # 如果订单处于已揽货或待绑定物流单号状态,绑定快递单号后自动变为运输中 order = self.order_repository.get_order(session, task.order_id) if order and order.order_status in ("picked_up", "pending_logistics"): before_status = order.order_status self.order_repository.update_order_status(session, order, "in_transit") self._notify_status_change(session, order, before_status, "in_transit", current_user.get("user_id") if current_user else None) audit_service.write_log( session, { "operate_type": "tracking_update", "biz_type": "logistics_task", "biz_id": task.id, "before_value": {"tracking_number": old_tracking, "express_company": old_express}, "after_value": {"tracking_number": new_tracking, "express_company": task.express_company}, "remark": payload.get("remark") or f"修改物流单号:{old_tracking} → {new_tracking}", }, ) session.commit() cache_delete_pattern("logistics:*") # 清除订单相关缓存,确保详情立即更新 cache_delete_pattern(f"order:detail:{task.order_id}") cache_delete_pattern("order:list:*") cache_delete_pattern("dashboard:*") return { "task_id": task.id, "order_id": task.order_id, "tracking_number": task.tracking_number, "express_company": task.express_company or "", } def accept_task( self, task_id: int, payload: dict | None = None, session: Session | None = None, current_user: dict | None = None, ) -> dict: """司机接单。 任务状态从 pending 变更为 accepted,订单状态变更为 accepted。 追加物流轨迹节点"司机已接单"。 Args: task_id: 任务 ID payload: 可选的请求体,包含 remark session: 数据库会话 current_user: 当前登录用户信息 Returns: 任务状态变更结果字典 被调用路由: logistics.py - POST /logistics/driver/tasks/{task_id}/accept """ result = self._change_task_status(task_id, "pending", "accepted", "accepted", session, current_user, payload or {}) self._append_trace( result["order_id"], { "task_id": task_id, "node_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "node_desc": "司机已接单", "node_type": "accepted", "source_platform": "driver", "remark": (payload or {}).get("remark"), }, ) return result def pickup_task( self, task_id: int, payload: dict | None = None, session: Session | None = None, current_user: dict | None = None, ) -> dict: """司机揽货确认。 必须至少上传一张照片。任务状态从 accepted 变更为 picked_up, 订单状态同步变更。追加物流轨迹节点"司机已揽货"。 Args: task_id: 任务 ID payload: 包含 photo_files、video_files 的请求体 session: 数据库会话 current_user: 当前登录用户信息 Returns: 任务状态变更结果字典 被调用路由: logistics.py - POST /logistics/driver/tasks/{task_id}/pickup """ operate_payload = payload or {} if not operate_payload.get("photo_files"): raise AppException(code=ErrorCode.PARAM_ERROR, message="确认揽货至少上传一张照片", status_code=400) # 获取运单号信息(可选) waybill_desc = "" if session is not None: waybills = self.logistics_repository.list_waybills_by_task_id(session, task_id) if waybills: # 同步第一条运单号到 logistics_task(兼容老逻辑) task = self.logistics_repository.get_task(session, task_id) if task and not task.tracking_number: task.tracking_number = waybills[0].tracking_number task.express_company = waybills[0].express_company waybill_desc = ",".join(w.tracking_number for w in waybills) result = self._change_task_status( task_id, "accepted", "picked_up", "picked_up", session, current_user, operate_payload, ) node_desc = f"司机已揽货" + (f",运单号:{waybill_desc}" if waybill_desc else "") self._append_trace( result["order_id"], { "task_id": task_id, "node_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "node_desc": node_desc, "node_type": "picked_up", "source_platform": "driver", "remark": operate_payload.get("remark"), "photo_files": operate_payload.get("photo_files", []), "video_files": operate_payload.get("video_files", []), }, ) return result def deliver_task( self, task_id: int, payload: dict | None = None, session: Session | None = None, current_user: dict | None = None, ) -> dict: """司机送达确认(已禁用)。 送达状态现由快递100物流追踪自动判断,无需手动确认。 此方法保留以兼容旧接口,调用时会抛出异常。 Args: task_id: 任务 ID payload: 可选的请求体 session: 数据库会话 current_user: 当前登录用户信息 Returns: 抛出异常,不返回结果 被调用路由: logistics.py - POST /logistics/driver/tasks/{task_id}/deliver """ raise AppException( code=ErrorCode.PARAM_ERROR, message="送达状态现由系统自动判断(根据物流单号查询快递100状态),无需手动确认", status_code=400, ) def get_trace( self, order_id: int, session: Session | None = None, current_user: dict | None = None, ) -> dict: """查询订单的物流轨迹。 合并内部轨迹和第三方物流轨迹(快递100等),按时间排序。 Args: order_id: 订单 ID session: 数据库会话 current_user: 当前登录用户信息 Returns: 包含 order_id 和 trace_list 的字典 被调用路由: logistics.py - GET /logistics/traces/{order_id} """ if session is not None: try: self._ensure_trace_access(order_id, session, current_user) traces = self.logistics_repository.list_traces_by_order_id(session, order_id) trace_list = [ { "task_id": item.task_id, "node_time": item.node_time.strftime("%Y-%m-%d %H:%M:%S") if item.node_time else "", "node_desc": item.node_desc, "node_type": item.node_type, "source_platform": item.source_platform, "remark": item.remark, } for item in traces ] try: third_party_trace_list = self._load_third_party_traces(order_id, session) except AppException as exc: logger.warning("[get_trace] 第三方轨迹查询失败,降级返回内部轨迹: %s", exc) third_party_trace_list = [] if trace_list or third_party_trace_list: return { "order_id": order_id, "trace_list": self._merge_trace_list(trace_list, third_party_trace_list), } except SQLAlchemyError as exc: raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc return {"order_id": order_id, "trace_list": []} def create_trace( self, order_id: int, payload: dict, session: Session | None = None, current_user: dict | None = None, ) -> dict: """手动创建物流轨迹节点。 Args: order_id: 订单 ID payload: 包含 node_time、node_desc、node_type 等的请求体 session: 数据库会话 current_user: 当前登录用户信息 Returns: 创建结果字典 被调用路由: logistics.py - POST /logistics/traces/{order_id} """ if session is not None: try: self._ensure_trace_access(order_id, session, current_user) node_time = datetime.strptime(payload["node_time"], "%Y-%m-%d %H:%M:%S") self.logistics_repository.create_trace( session, { "order_id": order_id, "task_id": payload.get("task_id"), "node_time": node_time, "node_desc": payload["node_desc"], "node_type": payload.get("node_type"), "source_platform": payload.get("source_platform"), "remark": payload.get("remark"), }, ) session.commit() cache_delete_pattern("logistics:*") return {"order_id": order_id, "created": True} except ValueError as exc: session.rollback() raise AppException(code=ErrorCode.PARAM_ERROR, message="物流节点时间格式错误", status_code=400) from exc except SQLAlchemyError as exc: session.rollback() raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库连接不可用", status_code=500) def _change_task_status( self, task_id: int, current_status: str, target_status: str, order_status: str, session: Session | None = None, current_user: dict | None = None, payload: dict | None = None, ) -> dict: """通用的任务状态变更方法。 校验当前状态、变更任务和订单状态、创建轨迹节点、持久化附件、记录审计日志。 送达时同步欠款数据。 Args: task_id: 任务 ID current_status: 期望的当前状态 target_status: 目标状态 order_status: 对应的订单状态 session: 数据库会话 current_user: 当前登录用户信息 payload: 可选的请求体(照片、视频等) Returns: 任务状态变更结果字典 """ operate_payload = payload or {} if session is not None: try: task = self.logistics_repository.get_task(session, task_id) if task is None: raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404) self._ensure_task_access(task, current_user) if task.status != current_status: raise AppException(code=ErrorCode.INVALID_STATUS, message="当前任务状态不允许执行该操作", status_code=400) self.logistics_repository.update_task_status(session, task, target_status) order = self.order_repository.get_order(session, task.order_id) if order is not None: self.order_repository.update_order_status(session, order, order_status) # 清除订单详情缓存,确保下次查询返回最新状态 cache_delete_pattern(f"order:detail:{order.id}") self.logistics_repository.create_trace( session, { "order_id": task.order_id, "task_id": task.id, "node_time": datetime.now(), "node_desc": self._status_desc(target_status), "node_type": target_status, "source_platform": "driver", "remark": operate_payload.get("remark"), }, ) self._persist_task_attachments(session, task, target_status, operate_payload, current_user) audit_service.write_log( session, { "operate_type": f"driver_task_{target_status}", "biz_type": "logistics_task", "biz_id": task.id, "before_value": {"task_status": current_status, "order_status": order.order_status}, "after_value": { "task_status": task.status, "order_status": order_status, "photo_count": len(operate_payload.get("photo_files", [])), "video_count": len(operate_payload.get("video_files", [])), }, "remark": operate_payload.get("remark") or f"司机任务状态变更为 {target_status}", }, ) if order_status == "pending_logistics": arrears_service.sync_order_arrears(session, order) # 触发物流状态更新事件 event_bus.emit("logistics_updated", { "task_id": task.id, "task_no": task.task_no, "order_id": order.id, "order_no": order.order_no, "salesman_id": order.salesman_id, "status": target_status, "status_desc": self._status_desc(target_status), "biz_type": "sales_order", "biz_id": order.id, }, session) session.commit() cache_delete_pattern("logistics:*") return {"task_id": task.id, "status": task.status, "order_id": task.order_id} except AppException: session.rollback() raise except SQLAlchemyError as exc: session.rollback() raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库连接不可用", status_code=500) def recognize_waybill(self, payload: dict, session: Session | None = None, current_user: dict | None = None) -> dict: """调用快递100面单OCR识别运单号(/elec/detocr 接口)。 使用 application/x-www-form-urlencoded 格式,MD5 签名认证。 失败时返回空结果让前端引导手动输入。 """ import base64 as _b64 import hashlib as _md5 from urllib.parse import quote, urlencode, urlparse, urlunparse image_url = (payload.get("image_url") or "").strip() if not image_url: raise AppException(code=ErrorCode.PARAM_ERROR, message="图片地址不能为空", status_code=400) settings = self.settings tracking_number = "" express_company = "" express_name = "" raw_result = {} if settings.kuaidi100_key and settings.kuaidi100_customer: try: # 1. 下载图片并 base64 编码(去掉 data:image 头) parsed = urlparse(image_url) safe_url = urlunparse(parsed._replace(path=quote(parsed.path))) with request.urlopen(safe_url, timeout=15) as resp: image_bytes = resp.read() image_b64 = _b64.b64encode(image_bytes).decode("utf-8") # 2. 构建 param JSON param_dict = { "image": image_b64, "enableTilt": True, "include": ["barcode", "receiver", "sender"], } param_str = json.dumps(param_dict, ensure_ascii=False, separators=(",", ":")) # 3. 签名: MD5(param + key + customer) 转大写 sign_str = param_str + settings.kuaidi100_key + settings.kuaidi100_customer sign = _md5.md5(sign_str.encode("utf-8")).hexdigest().upper() # 4. 构建 form-urlencoded 请求体 form_data = urlencode({ "key": settings.kuaidi100_key, "param": param_str, "sign": sign, }).encode("utf-8") req = request.Request( url="https://api.kuaidi100.com/elec/detocr", data=form_data, headers={"Content-Type": "application/x-www-form-urlencoded"}, method="POST", ) with request.urlopen(req, timeout=30) as resp: raw = resp.read().decode("utf-8") or "{}" result = json.loads(raw) raw_result = result # 5. 解析返回: data.waybill 为运单号, data.barcode 为条码 data = result.get("data") or {} tracking_number = (data.get("waybill") or "").strip() barcodes = data.get("barcode") or [] if barcodes and not tracking_number: tracking_number = str(barcodes[0]).strip() # 6. 根据运单号前缀推断快递公司 if tracking_number: express_company, express_name = self._guess_express_by_tracking(tracking_number) logger.info("[recognize_waybill] 快递100识别成功: %s, company=%s", tracking_number, express_name) except Exception as exc: logger.warning("[recognize_waybill] 快递100 OCR 失败,将返回空结果: %s", exc) raw_result = {"error": str(exc)} return { "tracking_number": tracking_number, "express_company": express_company, "express_name": express_name, "confidence": 0.9 if tracking_number else 0, "raw_result": raw_result, } # 常见快递公司运单号前缀映射 _EXPRESS_PREFIX_MAP = [ ("610", "ane", "安能物流"), ("773", "ane", "安能物流"), ("YT", "yuantong", "圆通速递"), ("SF", "shunfeng", "顺丰速运"), ("77", "shunfeng", "顺丰速运"), ("75", "shentong", "申通快递"), ("76", "shentong", "申通快递"), ("78", "zhongtong", "中通快递"), ("268", "zhongtong", "中通快递"), ("310", "yunda", "韵达快递"), ("468", "yunda", "韵达快递"), ("55", "yunda", "韵达快递"), ("228", "yunda", "韵达快递"), ("57", "jd", "京东物流"), ("JD", "jd", "京东物流"), ("D", "debang", "德邦快递"), ("dpk", "debang", "德邦快递"), ("33", "zhaijisong", "宅急送"), ("EM", "ems", "EMS"), ("E", "ems", "EMS"), ("CN", "guoji", "国际包裹"), ] @staticmethod def _guess_express_by_tracking(tracking_number: str) -> tuple[str, str]: """根据运单号前缀推断快递公司编码和名称。 Returns: (express_company, express_name) 元组,未匹配时返回空字符串。 """ tn = tracking_number.upper().strip() for prefix, code, name in LogisticsService._EXPRESS_PREFIX_MAP: if tn.startswith(prefix): return code, name return "", "" def submit_waybills(self, task_id: int, payload: dict, session: Session | None = None, current_user: dict | None = None) -> dict: """批量提交运单号,校验任务状态。""" if session is None: raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库连接不可用", status_code=500) task = self.logistics_repository.get_task(session, task_id) if task is None: raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404) self._ensure_task_access(task, current_user) if task.status not in {"accepted", "pending"}: raise AppException(code=ErrorCode.INVALID_STATUS, message="当前任务状态不允许录入运单号", status_code=400) waybills_data = payload.get("waybills", []) created = [] for item in waybills_data: tn = (item.get("tracking_number") or "").strip() if not tn: continue waybill = self.logistics_repository.create_waybill(session, { "task_id": task_id, "order_id": task.order_id, "tracking_number": tn, "express_company": (item.get("express_company") or "").strip() or None, "express_name": (item.get("express_name") or "").strip() or None, "image_url": (item.get("image_url") or "").strip() or None, "status": "pending", "created_by": current_user.get("user_id") if current_user else None, }) created.append({"id": waybill.id, "tracking_number": waybill.tracking_number, "express_company": waybill.express_company, "express_name": waybill.express_name}) audit_service.write_log(session, { "operate_type": "waybill_submit", "biz_type": "logistics_task", "biz_id": task_id, "before_value": None, "after_value": {"count": len(created)}, "remark": f"司机录入{len(created)}个运单号", }) session.commit() # 回填任务主表单号:仅当任务表单号为空时回填第一条运单 if created and not task.tracking_number: first = created[0] task.tracking_number = first["tracking_number"] task.express_company = first.get("express_company") session.commit() cache_delete_pattern("logistics:*") # 清除订单相关缓存,确保详情立即更新 cache_delete_pattern(f"order:detail:{task.order_id}") cache_delete_pattern("order:list:*") cache_delete_pattern("dashboard:*") return {"task_id": task_id, "waybill_count": len(created), "waybills": created} def list_waybills(self, task_id: int, session: Session | None = None, current_user: dict | None = None) -> list[dict]: """查询任务关联的运单号列表。""" if session is None: return [] task = self.logistics_repository.get_task(session, task_id) if task is None: return [] self._ensure_task_access(task, current_user) waybills = self.logistics_repository.list_waybills_by_task_id(session, task_id) return [{"id": w.id, "tracking_number": w.tracking_number, "express_company": w.express_company, "express_name": w.express_name, "image_url": w.image_url, "status": w.status} for w in waybills] def delete_waybill(self, task_id: int, waybill_id: int, session: Session | None = None, current_user: dict | None = None) -> dict: """删除单条运单号。""" if session is None: raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库连接不可用", status_code=500) task = self.logistics_repository.get_task(session, task_id) if task is None: raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404) self._ensure_task_access(task, current_user) if task.status not in {"accepted", "pending"}: raise AppException(code=ErrorCode.INVALID_STATUS, message="当前任务状态不允许删除运单号", status_code=400) ok = self.logistics_repository.delete_waybill(session, waybill_id, task_id) if not ok: raise AppException(code=ErrorCode.NOT_FOUND, message="运单号不存在", status_code=404) session.commit() # 重新查询剩余运单,重算任务主表摘要单号 remaining = self.logistics_repository.list_waybills_by_task_id(session, task_id) if remaining: first = remaining[0] task.tracking_number = first.tracking_number task.express_company = first.express_company else: task.tracking_number = None task.express_company = None session.commit() cache_delete_pattern("logistics:*") # 清除订单相关缓存,确保详情立即更新 cache_delete_pattern(f"order:detail:{task.order_id}") cache_delete_pattern("order:list:*") cache_delete_pattern("dashboard:*") return {"deleted": True} def _map_waybills(self, task_id: int, session: Session) -> list[dict]: """获取任务关联的运单号列表字典。""" waybills = self.logistics_repository.list_waybills_by_task_id(session, task_id) return [{"id": w.id, "tracking_number": w.tracking_number, "express_company": w.express_company, "express_name": w.express_name, "status": w.status} for w in waybills] def normalize_task_filters(self, filters: dict | None, current_user: dict | None = None) -> dict: """标准化任务查询过滤条件。 对司机角色自动注入 driver_id 过滤条件。 Args: filters: 原始过滤条件 current_user: 当前登录用户信息 Returns: 标准化后的过滤条件字典 """ normalized = dict(filters or {}) if current_user and current_user.get("role_code") == "driver": normalized["driver_id"] = current_user.get("user_id") return normalized def _ensure_task_access(self, task: object, current_user: dict | None) -> None: """校验当前用户是否有权访问指定物流任务。 管理员和经理可访问所有任务,司机只能访问自己的任务。 Args: task: 任务 ORM 对象 current_user: 当前登录用户信息 Raises: AppException: 无权限时抛出 FORBIDDEN """ if not current_user: return role_code = current_user.get("role_code") if role_code in {"admin", "manager"}: return if role_code == "driver": if task.driver_id != current_user.get("user_id"): raise AppException(code=ErrorCode.FORBIDDEN, message="仅可访问本人司机任务", status_code=403) return raise AppException(code=ErrorCode.FORBIDDEN, message="无权限访问", status_code=403) def _ensure_trace_access(self, order_id: int, session: Session, current_user: dict | None) -> None: """校验当前用户是否有权访问指定订单的物流轨迹。 Args: order_id: 订单 ID session: 数据库会话 current_user: 当前登录用户信息 Raises: AppException: 无权限时抛出 FORBIDDEN """ if not current_user: return role_code = current_user.get("role_code") if role_code in {"admin", "manager"}: return if role_code == "driver": task = self.logistics_repository.get_task_by_order_id(session, order_id) if task is None or task.driver_id != current_user.get("user_id"): raise AppException(code=ErrorCode.FORBIDDEN, message="仅可访问本人订单轨迹", status_code=403) return if role_code == "salesman": order = self.order_repository.get_order(session, order_id) if order is None or order.salesman_id != current_user.get("user_id"): raise AppException(code=ErrorCode.FORBIDDEN, message="仅可访问本人订单轨迹", status_code=403) return raise AppException(code=ErrorCode.FORBIDDEN, message="无权限访问", status_code=403) def _notify_driver_task_assigned(self, session: Session, task, order) -> None: """向司机发送任务分配提醒通知。 Args: session: 数据库会话 task: 任务 ORM 对象 order: 订单 ORM 对象 """ receiver = task.driver_id or 0 if receiver <= 0: return reminder_repo = ReminderRepository() reminder_repo.create_reminder( session, { "reminder_type": "task_assigned", "biz_type": "logistics_task", "biz_id": task.id, "receiver_user_id": receiver, "reminder_title": f"新任务分配 - {task.task_no}", "reminder_content": f"您有一个新运输任务 {task.task_no}(订单 {order.order_no}),请及时接单。", "status": "pending", "sent_at": datetime.now(), }, ) def _map_task_summary(self, task, session: Session) -> dict: """构建任务摘要字典,包含关联的订单号和工厂名称。 Args: task: 任务 ORM 对象 session: 数据库会话 Returns: 任务摘要字典 """ order = self.order_repository.get_order(session, task.order_id) supplier = self.order_repository.get_supplier(session, task.factory_id) return { "task_id": task.id, "task_no": task.task_no, "order_id": task.order_id, "order_no": order.order_no if order else "", "order_status": order.order_status if order else "", "driver_id": task.driver_id, "driver_name": f"司机{task.driver_id}", "factory_id": task.factory_id, "factory_name": supplier.supplier_name if supplier else "", "pickup_address": task.pickup_address, "delivery_address": task.delivery_address, "pickup_content": task.pickup_content, "quantity": float(task.quantity or 0), "status": task.status, "tracking_number": task.tracking_number or "", "express_company": task.express_company or "", "remark": task.remark, "created_at": task.created_at.strftime("%Y-%m-%d %H:%M:%S") if task.created_at else "", } def _map_task_detail(self, task, session: Session) -> dict: """构建任务详情字典,在摘要基础上补充产品规格和业务员信息。 仅暴露司机工作所需信息:产品规格、业务员。 不暴露客户姓名、手机号、价格等敏感信息。 """ summary = self._map_task_summary(task, session) order = self.order_repository.get_order(session, task.order_id) # 查询订单明细(产品规格) items = self.order_repository.list_order_items(session, task.order_id) if order else [] product_list = [ { "product_name": item.product_name, "specification": item.specification, "unit": item.unit, "quantity": float(item.quantity or 0), } for item in items ] summary["products"] = product_list # 查询业务员姓名 salesman_name = "" if order and order.salesman_id: from backend.app.models.system import User from sqlalchemy import select as sa_select user = session.execute( sa_select(User).where(User.id == order.salesman_id) ).scalar_one_or_none() if user: salesman_name = user.real_name summary["salesman_name"] = salesman_name summary["waybills"] = self._map_waybills(task.id, session) return summary def _append_trace(self, order_id: int, payload: dict) -> None: """将轨迹节点追加到内存中的轨迹记录列表。 Args: order_id: 订单 ID payload: 轨迹节点数据字典 """ trace_list = self.trace_records.setdefault(order_id, []) trace_list.append(payload) def _status_desc(self, status: str) -> str: """获取任务状态的中文描述。 Args: status: 任务状态编码 Returns: 中文描述字符串 """ if status == "accepted": return "司机已接单" if status == "picked_up": return "司机已揽货" if status == "delivered": return "司机已送达" return "物流状态更新" def _notify_status_change(self, session: Session, order, before_status: str, after_status: str, operator_id: int | None) -> None: """订单状态变更时生成提醒通知。 向业务员发送状态变更通知。 Args: session: 数据库会话 order: 订单 ORM 对象 before_status: 变更前状态 after_status: 变更后状态 operator_id: 操作人 ID """ _STATUS_LABELS = { "pending_logistics": "待绑定物流单号", "in_transit": "运输中", } before_label = _STATUS_LABELS.get(before_status, before_status) after_label = _STATUS_LABELS.get(after_status, after_status) receiver = order.salesman_id or 0 if receiver > 0: reminder_content = f"订单 {order.order_no} 状态由「{before_label}」变更为「{after_label}」。" self.reminder_repository.create_reminder( session, { "reminder_type": "order_status_change", "receiver_user_id": receiver, "reminder_title": "订单状态变更", "reminder_content": reminder_content, "biz_type": "sales_order", "biz_id": order.id, }, ) def _build_trace_payload(self, task_id: int, node_desc: str, node_type: str, payload: dict | None) -> dict: """构建物流轨迹节点数据载荷。 Args: task_id: 任务 ID node_desc: 节点描述 node_type: 节点类型 payload: 可选的请求体(照片、视频等) Returns: 轨迹节点数据字典 """ return { "task_id": task_id, "node_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"), "node_desc": node_desc, "node_type": node_type, "source_platform": "driver", "remark": (payload or {}).get("remark"), "photo_files": (payload or {}).get("photo_files", []), "video_files": (payload or {}).get("video_files", []), } def _persist_task_attachments( self, session: Session, task, target_status: str, payload: dict, current_user: dict | None, ) -> None: """持久化司机操作时上传的照片和视频附件。 将附件统一存入现有附件表,便于审计和详情展示复用。 Args: session: 数据库会话 task: 任务 ORM 对象 target_status: 目标任务状态 payload: 包含 photo_files、video_files 的请求体 current_user: 当前登录用户信息 """ file_payloads = [ *[ { "biz_type": "logistics_task", "biz_id": task.id, "file_name": item["file_name"], "file_url": item["file_url"], "file_type": item.get("file_type") or "image/jpeg", "file_size": item.get("file_size"), } for item in payload.get("photo_files", []) ], *[ { "biz_type": "logistics_task", "biz_id": task.id, "file_name": item["file_name"], "file_url": item["file_url"], "file_type": item.get("file_type") or "video/mp4", "file_size": item.get("file_size"), } for item in payload.get("video_files", []) ], ] if not file_payloads: return # 司机操作附件统一落到现有附件表,便于后续审计和详情展示复用。 for attachment_payload in file_payloads: repository_payload = file_service.build_attachment_payload( attachment_payload, created_by=current_user.get("user_id") if current_user else None, require_file_size=False, ) file_service.repository.create_attachment(session, repository_payload) audit_service.write_log( session, { "operate_type": "logistics_task_attachment_create", "biz_type": "logistics_task", "biz_id": task.id, "before_value": None, "after_value": { "task_status": target_status, "photo_count": len(payload.get("photo_files", [])), "video_count": len(payload.get("video_files", [])), }, "remark": f"司机任务 {task.task_no} 附件留痕", }, ) def _load_third_party_traces(self, order_id: int, session: Session | None = None) -> list[dict]: """加载第三方物流轨迹数据。 根据配置的 provider(internal/kuaidi100/通用 APPCODE)决定数据来源。 单号回退顺序:logistics_task -> 第一条 logistics_waybill -> sales_order Args: order_id: 订单 ID session: 数据库会话 Returns: 第三方轨迹节点列表,无配置或无运单号时返回空列表 """ provider = (self.settings.logistics_trace_provider or "internal").strip().lower() if provider == "internal": return [] tracking_number = "" express_company = "" if session is not None: task = self.logistics_repository.get_task_by_order_id(session, order_id) if task is not None: tracking_number = (task.tracking_number or "").strip() express_company = (task.express_company or "").strip() # 回退:如果任务表没有单号,尝试从运单子表获取 if not tracking_number: waybills = self.logistics_repository.list_waybills_by_task_id(session, task.id) if waybills: tracking_number = (waybills[0].tracking_number or "").strip() express_company = (waybills[0].express_company or "").strip() # 回退:如果没有物流任务或任务无单号,尝试从订单表获取 if not tracking_number: order = self.order_repository.get_order(session, order_id) if order is not None: tracking_number = (order.tracking_number or "").strip() # 快递公司:如果没有明确指定,根据单号自动识别 if tracking_number and not express_company: auto_code, _ = self._guess_express_by_tracking(tracking_number) if auto_code: express_company = auto_code if not tracking_number: return [] return self._fetch_remote_traces(tracking_number, express_company) def _fetch_remote_traces(self, tracking_number: str, express_company: str) -> list[dict]: """根据配置的 provider 调用第三方物流轨迹 API。 支持快递100和通用 APPCODE 两种模式。 Args: tracking_number: 运单号 express_company: 快递公司编码 Returns: 标准化后的轨迹节点列表 """ provider = (self.settings.logistics_trace_provider or "").strip().lower() # 兼容 "kdniao" 和 "kuaidi100" 两种配置别名 if provider in ("kuaidi100", "kdniao"): return self._fetch_kuaidi100_traces(tracking_number, express_company) # 通用 APPCODE 模式(兼容阿里云市场等第三方) if not self.settings.logistics_trace_endpoint.strip() or not self.settings.logistics_trace_path.strip(): raise AppException(code=ErrorCode.THIRD_PARTY_FAILED, message="第三方物流轨迹配置不完整", status_code=400) path = self.settings.logistics_trace_path.strip() if "{tracking_number}" in path: path = path.replace("{tracking_number}", tracking_number) endpoint = self.settings.logistics_trace_endpoint.strip().rstrip("/") if not path.startswith("/"): path = f"/{path}" url = f"{endpoint}{path}" req = request.Request( url=url, headers={ "Authorization": f"APPCODE {self.settings.logistics_trace_app_code.strip()}", "Content-Type": "application/json; charset=UTF-8", }, method="GET", ) try: with request.urlopen(req, timeout=20) as response: payload = json.loads(response.read().decode("utf-8") or "{}") except error.HTTPError as exc: detail = exc.read().decode("utf-8", errors="ignore") raise AppException( code=ErrorCode.THIRD_PARTY_FAILED, message=f"第三方物流轨迹调用失败:HTTP {exc.code} {detail}".strip(), status_code=400, ) from exc except error.URLError as exc: raise AppException( code=ErrorCode.THIRD_PARTY_FAILED, message=f"第三方物流轨迹网络请求失败:{exc.reason}", status_code=400, ) from exc except json.JSONDecodeError as exc: raise AppException(code=ErrorCode.THIRD_PARTY_FAILED, message="第三方物流轨迹返回格式异常", status_code=400) from exc return self._parse_generic_response(payload) def _fetch_kuaidi100_traces(self, tracking_number: str, express_company: str) -> list[dict]: """调用快递100 API 查询物流轨迹。 使用 MD5 签名认证,返回标准化后的轨迹节点列表。 Args: tracking_number: 运单号 express_company: 快递公司编码 Returns: 标准化后的轨迹节点列表 Raises: AppException: 配置不完整或 API 调用失败时抛出 """ import hashlib from urllib.parse import urlencode as _urlencode settings = self.settings if not settings.kuaidi100_key or not settings.kuaidi100_customer: raise AppException( code=ErrorCode.THIRD_PARTY_FAILED, message="快递100 配置不完整:缺少 KUAIDI100_KEY 或 KUAIDI100_CUSTOMER", status_code=400, ) # 快递100 传统查询接口(GET 方式) endpoint = (settings.logistics_trace_endpoint or "https://api.kuaidi100.com").strip().rstrip("/") path = (settings.logistics_trace_path or "/api").strip() if not path.startswith("/"): path = f"/{path}" url = f"{endpoint}{path}" # 签名:MD5(key + 运单号 + 快递公司编码 + customer) sign_str = f"{settings.kuaidi100_key}{tracking_number}{express_company}{settings.kuaidi100_customer}" sign = hashlib.md5(sign_str.encode("utf-8")).hexdigest() params = { "key": settings.kuaidi100_key, "com": express_company, "nu": tracking_number, "show": "0", "muti": "1", "order": "desc", "customer": settings.kuaidi100_customer, "sign": sign, } url = f"{url}?{_urlencode(params)}" req = request.Request( url=url, headers={"Content-Type": "application/x-www-form-urlencoded"}, method="GET", ) try: with request.urlopen(req, timeout=20) as response: raw = response.read().decode("utf-8") or "{}" payload = json.loads(raw) except error.HTTPError as exc: detail = exc.read().decode("utf-8", errors="ignore") raise AppException( code=ErrorCode.THIRD_PARTY_FAILED, message=f"快递100 查询失败:HTTP {exc.code} {detail}".strip(), status_code=400, ) from exc except error.URLError as exc: raise AppException( code=ErrorCode.THIRD_PARTY_FAILED, message=f"快递100 网络请求失败:{exc.reason}", status_code=400, ) from exc except json.JSONDecodeError as exc: raise AppException( code=ErrorCode.THIRD_PARTY_FAILED, message=f"快递100 返回格式异常:{raw[:200]}", status_code=400, ) from exc status_code = str(payload.get("status", "")) if status_code not in ("0", "1"): msg = payload.get("message") or payload.get("msg") or "未知错误" raise AppException( code=ErrorCode.THIRD_PARTY_FAILED, message=f"快递100 返回错误(status={status_code}):{msg}", status_code=400, ) # 响应格式:data 直接是轨迹数组,或 data 包含 trace 子数组 traces_raw = payload.get("data") or [] if isinstance(traces_raw, dict): traces_raw = traces_raw.get("trace") or traces_raw.get("list") or [] state = payload.get("state", "") state_map = {"0": "无轨迹", "1": "已揽收", "2": "在途中", "3": "已签收", "4": "问题件"} state_desc = state_map.get(str(state), "") normalized: list[dict] = [] for item in traces_raw: if not isinstance(item, dict): continue normalized.append({ "task_id": None, "node_time": str(item.get("ftime") or item.get("time") or ""), "node_desc": str(item.get("context") or ""), "node_type": state_desc or "kuaidi100", "source_platform": "kuaidi100", "remark": item.get("location", ""), }) return normalized def _parse_generic_response(self, payload: dict) -> list[dict]: """解析通用 APPCODE 模式的第三方物流响应。 从 payload 的 trace_list 字段提取轨迹数据并标准化。 Args: payload: API 原始响应字典 Returns: 标准化后的轨迹节点列表 """ rows = payload.get("trace_list") if isinstance(payload, dict) else [] if not isinstance(rows, list): return [] normalized: list[dict] = [] for item in rows: if not isinstance(item, dict): continue normalized.append({ "task_id": item.get("task_id"), "node_time": str(item.get("node_time") or ""), "node_desc": str(item.get("node_desc") or item.get("content") or ""), "node_type": str(item.get("node_type") or "third_party"), "source_platform": str(item.get("source_platform") or "third_party"), "remark": item.get("remark"), }) return normalized def _merge_trace_list(self, local_trace_list: list[dict], third_party_trace_list: list[dict]) -> list[dict]: """合并内部轨迹和第三方轨迹,按时间排序。 Args: local_trace_list: 内部轨迹列表 third_party_trace_list: 第三方轨迹列表 Returns: 合并后的轨迹列表,按 node_time 降序排列(最新的在前) """ merged = list(local_trace_list) + list(third_party_trace_list) merged.sort(key=lambda item: item.get("node_time") or "", reverse=True) return merged def check_delivery_by_tracking(self, session: Session | None = None) -> None: """定时检查运输中订单的物流状态,自动更新为已送达。 查询所有 in_transit 状态的订单,调用快递100 API 查询物流状态, 如果返回"已签收",自动将订单状态更新为 delivered。 Args: session: 数据库会话(可选,为None时使用默认会话) """ if session is None: return try: # 查询所有 in_transit 状态的订单 from backend.app.models.business import SalesOrder orders = session.query(SalesOrder).filter( SalesOrder.order_status == "in_transit" ).all() for order in orders: try: # 加载第三方物流轨迹 third_party_traces = self._load_third_party_traces(order.id, session) # 检查是否有"已签收"状态 for trace in third_party_traces: if trace.get("node_type") == "signed" or "已签收" in trace.get("node_desc", ""): # 自动更新为已送达 self._update_order_to_delivered(session, order) break except Exception as e: logger.warning("检查订单 %s 物流状态失败: %s", order.id, e) continue except SQLAlchemyError as e: logger.error("检查物流状态数据库异常: %s", e) def _update_order_to_delivered(self, session: Session, order) -> None: """将订单状态更新为已送达,并同步相关数据。 Args: session: 数据库会话 order: 订单 ORM 对象 """ before_status = order.order_status self.order_repository.update_order_status(session, order, "delivered") # 同步欠款数据 arrears_service.sync_order_arrears(session, order) # 记录审计日志 audit_service.write_log( session, { "operate_type": "order_auto_delivered", "biz_type": "sales_order", "biz_id": order.id, "before_value": {"order_status": before_status}, "after_value": {"order_status": "delivered"}, "remark": f"订单 {order.order_no} 物流签收自动更新为已送达", }, ) # 发送通知 self._notify_status_change(session, order, before_status, "delivered", None) cache_delete_pattern(f"order:detail:{order.id}") def _notify_status_change(self, session: Session, order, before_status: str, after_status: str, operator_id: int | None) -> None: """订单状态变更时生成提醒通知。 Args: session: 数据库会话 order: 订单 ORM 对象 before_status: 变更前状态 after_status: 变更后状态 operator_id: 操作人 ID """ from backend.app.services.order_service import OrderService order_service = OrderService() order_service._notify_status_change(session, order, before_status, after_status, operator_id) logistics_service = LogisticsService()