dingdanquanliucheng/backend/app/services/logistics_service.py
taiyi e7473fc757 feat: 运单号识别后自动推断快递公司并显示中文名称
- 新增运单号前缀映射表覆盖安能/圆通/顺丰/申通/中通/韵达等
- 前端运单卡片显示快递公司中文名而非编码

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

1388 lines
57 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.

"""物流管理服务层。
负责司机运输任务的创建、状态流转、物流轨迹查询等业务逻辑。
支持内部轨迹和第三方物流快递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.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()
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 {"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()
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 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()
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
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()
cache_delete_pattern("logistics:*")
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)
# 校验至少有一条运单号
if session is not None:
waybills = self.logistics_repository.list_waybills_by_task_id(session, task_id)
if not waybills:
raise AppException(code=ErrorCode.PARAM_ERROR, message="请至少录入一个物流运单号", status_code=400)
# 同步第一条运单号到 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)
else:
waybill_desc = ""
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:
"""司机送达确认。
任务状态从 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()
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)
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()
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="http://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:
result = json.loads(resp.read().decode("utf-8") or "{}")
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()
cache_delete_pattern("logistics:*")
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()
cache_delete_pattern("logistics:*")
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
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:
"""构建任务详情字典,在摘要基础上补充产品规格和业务员信息。
仅暴露司机工作所需信息:产品规格、业务员。
不暴露客户姓名、手机号、价格等敏感信息。
"""
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 _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]:
"""加载第三方物流轨迹数据。
根据配置的 providerinternal/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()