dingdanquanliucheng/backend/app/services/logistics_service.py

227 lines
9.8 KiB
Python

from datetime import datetime
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.orm import Session
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.services.demo_store import demo_store
class LogisticsService:
def __init__(self) -> None:
self.logistics_repository = LogisticsRepository()
self.order_repository = OrderRepository()
self.trace_records: dict[int, list[dict]] = {}
def list_tasks(self, session: Session | None = None, filters: dict | None = None) -> dict:
if session is not None:
try:
tasks = self.logistics_repository.list_tasks_by_filters(session, filters or {})
return {
"total": len(tasks),
"page_no": 1,
"page_size": 20,
"list": [self._map_task_summary(task, session) for task in tasks],
}
except SQLAlchemyError:
pass
return demo_store.list_logistics_tasks()
def create_task(self, payload: dict, session: Session | None = None) -> dict:
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",
"created_by": None,
"remark": payload.get("remark"),
},
)
# 司机任务创建后,订单进入待司机接单阶段。
self.order_repository.update_order_status(session, order, "pending_driver")
session.commit()
return self._map_task_detail(task, session)
except SQLAlchemyError:
session.rollback()
return demo_store.create_logistics_task(payload)
def get_task(self, task_id: int, session: Session | 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)
return self._map_task_detail(task, session)
except SQLAlchemyError:
pass
task = demo_store.get_logistics_task(task_id)
if task is None:
raise AppException(code=ErrorCode.NOT_FOUND, message="司机任务不存在", status_code=404)
return task
def list_driver_tasks(self, session: Session | None = None) -> dict:
return self.list_tasks(session=session, filters={})
def get_driver_task(self, task_id: int, session: Session | None = None) -> dict:
return self.get_task(task_id, session=session)
def accept_task(self, task_id: int, payload: dict | None = None, session: Session | None = None) -> dict:
result = self._change_task_status(task_id, "pending", "accepted", "accepted", session)
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) -> dict:
result = self._change_task_status(task_id, "accepted", "picked_up", "picked_up", session)
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": (payload or {}).get("remark"),
"photo_files": (payload or {}).get("photo_files", []),
"video_files": (payload or {}).get("video_files", []),
},
)
return result
def deliver_task(self, task_id: int, payload: dict | None = None, session: Session | None = None) -> dict:
result = self._change_task_status(task_id, "picked_up", "delivered", "delivered", session)
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": "delivered",
"source_platform": "driver",
"remark": (payload or {}).get("remark"),
"photo_files": (payload or {}).get("photo_files", []),
"video_files": (payload or {}).get("video_files", []),
},
)
return result
def get_trace(self, order_id: int) -> dict:
return {
"order_id": order_id,
"trace_list": self.trace_records.get(
order_id,
[
{
"node_time": "2026-05-14 12:00:00",
"node_desc": "司机已揽货",
"node_type": "picked_up",
"source_platform": "manual",
}
],
),
}
def create_trace(self, order_id: int, payload: dict) -> dict:
self._append_trace(order_id, payload)
return {"order_id": order_id, "created": True}
def _change_task_status(
self,
task_id: int,
current_status: str,
target_status: str,
order_status: str,
session: Session | 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)
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)
session.commit()
return {"task_id": task.id, "status": task.status, "order_id": task.order_id}
except SQLAlchemyError:
session.rollback()
task = demo_store.get_logistics_task(task_id)
return {
"task_id": task_id,
"status": target_status,
"order_id": task["order_id"] if task else 0,
}
def _map_task_summary(self, task, session: Session) -> dict:
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,
"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)
summary["salesman_name"] = ""
summary["customer_name"] = order.customer_name if order else ""
return summary
def _append_trace(self, order_id: int, payload: dict) -> None:
trace_list = self.trace_records.setdefault(order_id, [])
trace_list.append(payload)
logistics_service = LogisticsService()