"""物流管理服务层。 负责司机运输任务的创建、状态流转、物流轨迹查询等业务逻辑。 支持内部轨迹和第三方物流(快递100、通用 APPCODE)轨迹查询。 依赖 LogisticsRepository、OrderRepository 进行数据持久化, 依赖 audit_service 记录操作审计日志,依赖 arrears_service 同步欠款数据, 依赖 file_service 管理司机任务附件。 """ from datetime import datetime import json 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.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.file_service import file_service 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() 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: try: tasks = self.logistics_repository.list_tasks_by_filters( session, self.normalize_task_filters(filters, current_user), ) return { "total": len(tasks), "page_no": 1, "page_size": 20, "list": [self._map_task_summary(task, session) for task in tasks], } 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 {"pending_factory", "approved"}: 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") self._notify_driver_task_assigned(session, task, order) 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() 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 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 状态, 则回退为 pending_factory。自动记录审计日志。 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, "pending_factory") 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() 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 task.express_company = (payload.get("express_company") or "").strip() or 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() 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) result = self._change_task_status( task_id, "accepted", "picked_up", "picked_up", session, current_user, operate_payload, ) 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": "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: """司机送达确认。 任务状态从 picked_up 变更为 delivered,订单状态同步变更。 追加物流轨迹节点"司机已送达",同步欠款数据。 Args: task_id: 任务 ID payload: 可选的请求体,包含 remark、photo_files 等 session: 数据库会话 current_user: 当前登录用户信息 Returns: 任务状态变更结果字典 被调用路由: logistics.py - POST /logistics/driver/tasks/{task_id}/deliver """ operate_payload = payload or {} result = self._change_task_status( task_id, "picked_up", "delivered", "delivered", session, current_user, operate_payload, ) self._append_trace(result["order_id"], self._build_trace_payload(task_id, "司机已送达", "delivered", operate_payload)) return result 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 ] third_party_trace_list = self._load_third_party_traces(order_id, session) 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() 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) 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 target_status == "delivered": arrears_service.sync_order_arrears(session, order) session.commit() 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 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 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 "", "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: """构建任务详情字典,在摘要基础上补充客户名称。 Args: task: 任务 ORM 对象 session: 数据库会话 Returns: 任务详情字典 """ summary = self._map_task_summary(task, session) order = self.order_repository.get_order(session, task.order_id) summary["salesman_name"] = "" summary["customer_name"] = order.customer_name if order else "" 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 _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)决定数据来源。 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: 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() if provider == "kuaidi100": 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 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, ) endpoint = (settings.logistics_trace_endpoint or "https://api.kuaidi100.com").strip().rstrip("/") path = (settings.logistics_trace_path or "/query").strip() if not path.startswith("/"): path = f"/{path}" url = f"{endpoint}{path}" param = { "com": express_company, "num": tracking_number, "phone": "", "from": "", "to": "", } # 快递100 签名:MD5(key + 排序参数拼接 + customer + key) sorted_params = "".join(f"{k}{v}" for k, v in sorted(param.items()) if v) sign_str = f"{settings.kuaidi100_key}{sorted_params}{settings.kuaidi100_customer}{settings.kuaidi100_key}" sign = hashlib.md5(sign_str.encode("utf-8")).hexdigest() body = json.dumps({ "customer": settings.kuaidi100_customer, "key": settings.kuaidi100_key, "sign": sign, "param": param, }, ensure_ascii=False).encode("utf-8") req = request.Request( url=url, data=body, headers={"Content-Type": "application/json; charset=UTF-8"}, method="POST", ) 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"快递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="快递100 返回格式异常", 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 "") return merged logistics_service = LogisticsService()