dingdanquanliucheng/backend/app/services/order_service.py

2101 lines
95 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.

"""销售订单服务层。
负责销售订单的全生命周期管理,包括创建、编辑、提交审核、审批、取消、
发厂文案生成、状态流转等核心业务逻辑。
依赖多个 Repository 进行数据持久化,依赖 audit_service 记录操作审计日志,
依赖 arrears_service 同步欠款数据,依赖 ReminderRepository 发送状态变更通知。
"""
from datetime import datetime
import json
import re
from sqlalchemy.exc import SQLAlchemyError
from sqlalchemy.orm import Session
from backend.app.core.cache import cache_delete, 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.customer_repository import CustomerRepository
from backend.app.repositories.config_repository import ConfigRepository
from backend.app.repositories.file_repository import FileRepository
from backend.app.repositories.logistics_repository import LogisticsRepository
from backend.app.repositories.order_approve_log_repository import OrderApproveLogRepository
from backend.app.repositories.order_repository import OrderRepository
from backend.app.repositories.supplier_repository import SupplierRepository
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
class OrderService:
"""订单服务,封装销售订单的完整业务流程。
依赖:
OrderRepository - 订单数据访问
CustomerRepository - 客户数据访问(订单录入时自动创建客户)
ConfigRepository - 系统配置读取
OrderApproveLogRepository - 审批日志
LogisticsRepository - 物流任务和轨迹
FileRepository - 附件查询
ReminderRepository - 状态变更提醒通知
arrears_service - 欠款同步
audit_service - 操作审计日志
"""
def __init__(self) -> None:
self.order_repository = OrderRepository()
self.customer_repository = CustomerRepository()
self.config_repository = ConfigRepository()
self.order_approve_log_repository = OrderApproveLogRepository()
self.logistics_repository = LogisticsRepository()
self.supplier_repository = SupplierRepository()
self.file_repository = FileRepository()
self.reminder_repository = ReminderRepository()
def list_orders(self, session: Session | None = None, filters: dict | None = None, current_user: dict | None = None) -> dict:
"""查询订单列表。
根据过滤条件查询订单,并根据用户角色过滤敏感字段(如利润、佣金对业务员不可见)。
Args:
session: 数据库会话
filters: 查询过滤条件
current_user: 当前登录用户信息,用于角色级数据过滤
Returns:
包含 total 和 list 的字典
被调用路由: orders.py - GET /orders, salesman.py - GET /salesman/orders
"""
if session is not None:
# 读缓存
cache_key = make_cache_key("order:list", user_id=(current_user or {}).get("user_id"), **(filters or {}))
cached = cache_get(cache_key)
if cached is not None:
return cached
try:
result = self.order_repository.list_orders_by_filters(session, filters or {})
total = result["total"] or 0
orders = result["list"]
role_code = (current_user or {}).get("role_code")
# 批量查询业务员信息
from backend.app.models.system import User as SysUser
salesman_ids = list(set(item.salesman_id for item in orders if item.salesman_id))
salesman_map = {}
if salesman_ids:
salesmen = session.query(SysUser).filter(SysUser.id.in_(salesman_ids)).all()
salesman_map = {s.id: (s.real_name or s.username or "") for s in salesmen}
rows = []
for item in orders:
# 获取订单明细摘要
order_items = self.order_repository.list_order_items(session, item.id)
items_summary = ""
if order_items:
summary_parts = [f"{oi.product_name}×{int(oi.quantity or 0)}" for oi in order_items[:2]]
items_summary = "".join(summary_parts)
if len(order_items) > 2:
items_summary += "..."
# 获取物流信息
task = self.logistics_repository.get_task_by_order_id(session, item.id)
logistics_info = ""
if task:
logistics_info = task.express_company + " " + (task.tracking_number or "") if task.express_company else task.status or ""
row = self._filter_order_row(
{
"order_id": item.id,
"order_no": item.order_no,
"customer_name": item.customer_name,
"customer_mobile": item.customer_mobile,
"salesman_name": salesman_map.get(item.salesman_id, ""),
"order_status": item.order_status,
"order_source": item.order_source,
"order_type": item.order_type,
"contract_amount": float(item.contract_amount or 0),
"cost_price_total": float(item.cost_price_total or 0),
"profit_total": float(item.profit_total or 0),
"commission_amount": float(item.commission_amount or 0),
"need_invoice": item.need_invoice or 0,
"created_at": item.created_at.strftime("%Y-%m-%d %H:%M:%S") if item.created_at else "",
"items_summary": items_summary,
"logistics_info": logistics_info,
},
role_code,
)
rows.append(row)
# 写缓存
result_data = {"total": total, "list": rows}
cache_set(cache_key, result_data, ttl=30)
return result_data
except SQLAlchemyError as exc:
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return {"total": 0, "list": []}
def get_order(
self,
order_id: int,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""查询单个订单详情。
聚合订单主档、明细、审批日志、物流任务、物流轨迹、附件、计算明细等信息。
根据用户角色过滤敏感字段。
Args:
order_id: 订单 ID
session: 数据库会话
current_user: 当前登录用户信息
Returns:
订单完整详情字典,订单不存在时返回 None
被调用路由: orders.py - GET /orders/{order_id}
"""
if session is not None:
# 读缓存
cache_key = f"order:detail:{order_id}"
cached = cache_get(cache_key)
if cached is not None:
return self._filter_order_detail_by_role(cached, current_user)
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user)
supplier = self.order_repository.get_supplier(session, order.factory_id)
items = self.order_repository.list_order_items(session, order.id)
approve_logs = self.order_approve_log_repository.list_by_order_id(session, order.id)
task = self.logistics_repository.get_task_by_order_id(session, order.id)
traces = self.logistics_repository.list_traces_by_order_id(session, order.id)
order_attachments = self.file_repository.list_attachments(session, "sales_order", order.id)
task_attachments = self.file_repository.list_attachments(session, "logistics_task", task.id) if task else []
# 查询业务员姓名
salesman_name = ""
if order.salesman_id:
from backend.app.models.system import User as SysUser
salesman = session.get(SysUser, order.salesman_id)
if salesman:
salesman_name = salesman.real_name or salesman.username or ""
result = {
"order_id": order.id,
"order_no": order.order_no,
"customer_id": order.customer_id,
"customer_name": order.customer_name,
"customer_mobile": order.customer_mobile,
"customer_address": order.customer_address,
"salesman_id": order.salesman_id,
"salesman_name": salesman_name,
"order_source": order.order_source,
"order_status": order.order_status,
"delivery_type": order.delivery_type,
"factory_id": order.factory_id,
"factory_name": supplier.supplier_name if supplier else "",
"contract_amount": float(order.contract_amount or 0),
"sale_price_total": float(order.sale_price_total or 0),
"cost_price_total": float(order.cost_price_total or 0),
"rebate_total": float(order.rebate_total or 0),
"freight_total": float(order.freight_total or 0),
"tax_total": float(order.tax_total or 0),
"other_fee_total": float(order.other_fee_total or 0),
"profit_total": float(order.profit_total or 0),
"profit_rate": float(order.profit_rate or 0),
"commission_amount": float(order.commission_amount or 0),
"payment_method": order.payment_method,
"tax_amount": float(order.tax_amount or 0),
"need_invoice": order.need_invoice or 0,
"order_type": order.order_type,
"self_delivery": order.self_delivery or 0,
"tracking_number": order.tracking_number,
"items": [
{
"item_id": item.id,
"product_id": item.product_id,
"product_name": item.product_name,
"specification": item.specification,
"unit": item.unit,
"quantity": float(item.quantity or 0),
"sale_price": float(item.sale_price or 0),
"cost_price": float(item.cost_price or 0),
"rebate_amount": float(item.rebate_amount or 0),
"freight_amount": float(item.freight_amount or 0),
"tax_amount": float(item.tax_amount or 0),
"other_fee_amount": float(item.other_fee_amount or 0),
"remark": item.remark,
"demand_specification": getattr(item, "demand_specification", None),
"pricing_type": getattr(item, "pricing_type", None),
"length_m": float(item.length_m) if getattr(item, "length_m", None) else None,
"width_m": float(item.width_m) if getattr(item, "width_m", None) else None,
"area_sqm": float(item.area_sqm) if getattr(item, "area_sqm", None) else None,
"surcharge_detail": getattr(item, "surcharge_detail", None),
"processing_detail": getattr(item, "processing_detail", None),
"supplier_id": getattr(item, "supplier_id", None),
"supplier_model": getattr(item, "supplier_model", None),
"price_tier": getattr(item, "price_tier", None),
}
for item in items
],
"approve_logs": [
{
"log_id": log.id,
"approve_type": log.approve_type,
"approve_result": log.approve_result,
"before_status": log.before_status,
"after_status": log.after_status,
"approve_opinion": log.approve_opinion,
"operator_id": log.operator_id,
"approve_time": log.created_at.strftime("%Y-%m-%d %H:%M:%S") if log.created_at else "",
}
for log in approve_logs
],
"logistics_info": {
"task": self._build_task_summary(task),
"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
],
},
"attachments": [
*[self._build_attachment_row(item) for item in order_attachments],
*[self._build_attachment_row(item) for item in task_attachments],
],
"calculation_detail": self._build_calculation_detail(order, items),
"profit_alert_threshold": float(self._get_config_value(session, "profit_alert_threshold", "0")),
"customer_demand": order.customer_demand,
"remark": order.remark,
}
# 写缓存(缓存完整结果,角色过滤在读取后执行)
cache_set(cache_key, result, ttl=60)
return self._filter_order_detail_by_role(result, current_user)
except SQLAlchemyError as exc:
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def update_order(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""编辑订单。
仅草稿和已退回状态的订单允许编辑。编辑时若客户不存在会自动创建。
重新计算销售总额、成本总额、利润等汇总数据。
Args:
order_id: 订单 ID
payload: 包含订单主档和明细的请求体
session: 数据库会话
current_user: 当前登录用户信息
Returns:
更新结果字典,包含订单汇总数据
被调用路由: orders.py - PUT /orders/{order_id}
"""
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user, owner_only=True)
self._ensure_status(order.order_status, {"draft", "rejected", "pending_approve"}, "当前状态不允许编辑订单")
if not payload["customer_name"].strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="客户姓名不能为空", status_code=400)
if not payload["customer_mobile"].strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="客户手机号不能为空", status_code=400)
customer = self.customer_repository.find_by_name_and_mobile(
session,
payload["customer_name"],
payload["customer_mobile"],
)
if customer is None:
customer = self.customer_repository.create_customer(
session,
{
"customer_name": payload["customer_name"],
"mobile": payload["customer_mobile"],
"address": payload.get("customer_address"),
"settlement_type": None,
"settlement_days": 0,
"customer_type": None,
"salesman_id": current_user.get("user_id") if current_user else payload.get("salesman_id"),
"credit_limit": 0,
"remark": "订单编辑时自动创建",
},
)
items = payload["items"]
if not items:
raise AppException(code=ErrorCode.PARAM_ERROR, message="订单明细不能为空", status_code=400)
sale_total = sum(item["quantity"] * item["sale_price"] for item in items)
cost_total = sum(item["quantity"] * item["cost_price"] for item in items)
# 使用合同金额作为收入计算利润,如果没有合同金额则使用系统计算的销售金额
income_amount = payload.get("contract_amount", 0) or sale_total
profit_total = income_amount - cost_total - payload["rebate_total"] - payload["freight_total"] - payload["tax_total"] - payload["other_fee_total"]
profit_rate = round((profit_total / income_amount) * 100, 2) if income_amount else 0
before_status = order.order_status
self.order_repository.update_order_with_items(
session,
order,
{
"customer_id": customer.id,
"customer_name": payload["customer_name"],
"customer_mobile": payload["customer_mobile"],
"customer_address": payload.get("customer_address"),
"salesman_id": current_user.get("user_id") if current_user else payload.get("salesman_id"),
"order_source": payload.get("order_source"),
"delivery_type": payload.get("delivery_type"),
"factory_id": payload.get("factory_id"),
"contract_amount": payload.get("contract_amount", 0),
"sale_price_total": sale_total,
"cost_price_total": cost_total,
"rebate_total": payload["rebate_total"],
"freight_total": payload["freight_total"],
"tax_total": payload["tax_total"],
"other_fee_total": payload["other_fee_total"],
"profit_total": profit_total,
"profit_rate": profit_rate,
"commission_amount": payload["commission_amount"],
"payment_method": payload.get("payment_method"),
"tax_amount": payload.get("tax_amount", 0),
"need_invoice": payload.get("need_invoice", 0),
"order_type": payload.get("order_type"),
"self_delivery": payload.get("self_delivery", 0),
"tracking_number": payload.get("tracking_number"),
"remark": payload.get("remark"),
"customer_demand": payload.get("customer_demand"),
},
items,
)
audit_service.write_log(
session,
{
"operate_type": "order_update",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"order_status": before_status},
"after_value": {
"order_no": order.order_no,
"order_status": order.order_status,
"sale_price_total": round(sale_total, 2),
"cost_price_total": round(cost_total, 2),
"profit_total": round(profit_total, 2),
"profit_rate": profit_rate,
"commission_amount": payload["commission_amount"],
},
"remark": f"更新订单 {order.order_no}",
},
)
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_no": order.order_no,
"order_status": order.order_status,
"sale_price_total": round(sale_total, 2),
"cost_price_total": round(cost_total, 2),
"profit_total": round(profit_total, 2),
"profit_rate": profit_rate,
"commission_amount": payload["commission_amount"],
"payment_method": payload.get("payment_method"),
"tax_amount": payload.get("tax_amount", 0),
"need_invoice": payload.get("need_invoice", 0),
"order_type": payload.get("order_type"),
"self_delivery": payload.get("self_delivery", 0),
"tracking_number": payload.get("tracking_number"),
}
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 submit_order(
self,
order_id: int,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""提交订单审核。
仅草稿和已退回状态的订单允许提交。提交后状态变为 pending_approve
并向经理角色发送审批提醒通知。
Args:
order_id: 订单 ID
session: 数据库会话
current_user: 当前登录用户信息
Returns:
提交结果字典
被调用路由: orders.py - POST /orders/{order_id}/submit
"""
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user, owner_only=True)
self._ensure_status(order.order_status, {"draft", "rejected", "pending_approve"}, "当前状态不允许提交审核")
before_status = order.order_status
self.order_repository.update_order_status(session, order, "pending_approve")
audit_service.write_log(
session,
{
"operate_type": "order_submit",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"order_status": before_status},
"after_value": {"order_status": order.order_status},
"remark": f"提交订单审核 {order.order_no}",
},
)
self._notify_status_change(session, order, before_status, order.order_status, current_user.get("user_id") if current_user else None)
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_status": order.order_status,
"submitted_at": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
}
except AppException:
session.rollback()
raise
except SQLAlchemyError as exc:
session.rollback()
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def cancel_order(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""发起订单取消。
根据订单当前状态决定取消方式:草稿/待审核可直接取消,
已通过/履约中的订单需走取消审批流程。必须填写取消原因。
Args:
order_id: 订单 ID
payload: 包含 cancel_reason、cancel_opinion 的请求体
session: 数据库会话
current_user: 当前登录用户信息
Returns:
取消结果字典,包含状态变更信息
被调用路由: orders.py - POST /orders/{order_id}/cancel
"""
self._validate_cancel_payload(payload)
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user, owner_only=True)
if order.order_status in {"canceled", "settled"}:
raise AppException(code=ErrorCode.INVALID_STATUS, message="当前状态不允许取消订单", status_code=400)
allow_cancel_config = self._get_config_value(session, "allow_cancel_statuses", "draft,pending_approve,approved,pending_factory")
allow_cancel_statuses = {s.strip() for s in allow_cancel_config.split(",") if s.strip()}
if order.order_status not in allow_cancel_statuses:
raise AppException(code=ErrorCode.INVALID_STATUS, message=f"当前状态「{order.order_status}」不允许取消", status_code=400)
_DIRECT_CANCEL_STATUSES = {"draft", "pending_approve"}
_APPROVED_CANCEL_STATUSES = {"approved"}
_FULFILLMENT_CANCEL_STATUSES = {"pending_factory", "pending_driver", "accepted", "picked_up", "pending_logistics", "in_transit", "shipped"}
if order.order_status in _DIRECT_CANCEL_STATUSES:
next_status = "canceled"
elif order.order_status in _APPROVED_CANCEL_STATUSES:
next_status = "cancel_pending"
elif order.order_status in _FULFILLMENT_CANCEL_STATUSES:
next_status = "cancel_fulfillment_pending"
else:
next_status = "cancel_pending"
previous_status = order.order_status
self.order_repository.update_cancel_fields(
session,
order,
{
"order_status": next_status,
"cancel_requested_by": current_user.get("user_id") if current_user else None,
"cancel_requested_at": datetime.now(),
"cancel_reason": payload["cancel_reason"].strip(),
"cancel_opinion": payload.get("cancel_opinion"),
"cancel_previous_status": None if next_status == "canceled" else previous_status,
},
)
if next_status == "canceled":
self.order_repository.clear_cancel_request(session, order)
self.order_repository.update_order_status(session, order, "canceled")
audit_service.write_log(
session,
{
"operate_type": "order_cancel",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"order_status": previous_status},
"after_value": {
"order_status": order.order_status,
"cancel_reason": order.cancel_reason,
"cancel_previous_status": order.cancel_previous_status,
"cancel_requested_by": order.cancel_requested_by,
},
"remark": f"发起订单取消 {order.order_no}",
},
)
self._notify_status_change(session, order, previous_status, order.order_status, current_user.get("user_id") if current_user else None)
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_status": order.order_status,
"previous_status": previous_status,
"canceled_at": order.cancel_requested_at.strftime("%Y-%m-%d %H:%M:%S") if order.cancel_requested_at else None,
}
except AppException:
session.rollback()
raise
except SQLAlchemyError as exc:
session.rollback()
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def approve_order(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""审批订单(通过或退回)。
仅待审核状态的订单允许审批。通过后进入待下发工厂状态,
退回后状态恢复为 rejected。驳回时必须填写审批意见。
自动记录审批日志并通知业务员。
Args:
order_id: 订单 ID
payload: 包含 approve_resultpass/reject和 approve_opinion 的请求体
session: 数据库会话
current_user: 当前登录用户信息
Returns:
审批结果字典
被调用路由: orders.py - POST /orders/{order_id}/approve
"""
self._validate_approve_payload(payload)
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user)
self._ensure_status(order.order_status, {"pending_approve"}, "当前状态不允许审批")
previous_status = order.order_status
if payload["approve_result"] == "pass":
# 自发订单(已填快递单号)审批通过后直接进入运输中,无需工厂/司机环节
if order.tracking_number:
target_status = "in_transit"
elif payload.get("factory_id") and payload.get("driver_id"):
# 有工厂+司机:审批通过并直接创建物流任务,跳到待司机接单
target_status = "pending_driver"
elif payload.get("factory_id"):
# 仅有工厂:审批通过并下发工厂
target_status = "pending_factory"
else:
target_status = "approved"
else:
target_status = "rejected"
# 如果有工厂ID设置到订单上
if payload.get("factory_id") and target_status in ("pending_factory",):
order.factory_id = payload["factory_id"]
# 同时选择了工厂和司机:在同一个事务内创建物流任务,直接进入待司机接单
if payload.get("factory_id") and payload.get("driver_id") and target_status == "pending_factory":
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["factory_id"],
"pickup_address": payload.get("pickup_address", ""),
"delivery_address": payload.get("delivery_address", ""),
"pickup_content": payload.get("pickup_content", ""),
"quantity": payload.get("quantity", 1),
"status": "pending",
"created_by": current_user.get("user_id") if current_user else None,
},
)
target_status = "pending_driver"
self.order_repository.update_order_status(session, order, target_status)
# 有工厂+司机时,在同一事务内创建物流任务
if payload["approve_result"] == "pass" and target_status == "pending_driver" and payload.get("driver_id"):
# 获取工厂地址作为取货地址
pickup_address = ""
factory = self.supplier_repository.get_supplier(session, payload["factory_id"])
if factory and factory.address:
pickup_address = factory.address
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["factory_id"],
"pickup_address": pickup_address,
"delivery_address": "",
"pickup_content": "",
"quantity": 1,
"status": "pending",
"created_by": current_user.get("user_id") if current_user else None,
},
)
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)
self.order_approve_log_repository.create_log(
session,
{
"order_id": order.id,
"approve_type": "order",
"approve_result": payload["approve_result"],
"before_status": previous_status,
"after_status": target_status,
"approve_opinion": payload.get("approve_opinion"),
"operator_id": current_user.get("user_id") if current_user else None,
},
)
audit_service.write_log(
session,
{
"operate_type": "order_approve",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"order_status": previous_status},
"after_value": {
"order_status": order.order_status,
"approve_result": payload["approve_result"],
"approve_opinion": payload.get("approve_opinion"),
},
"remark": f"订单审批 {order.order_no}",
},
)
self._notify_status_change(session, order, previous_status, order.order_status, current_user.get("user_id") if current_user else None, {"approve_opinion": payload.get("approve_opinion")})
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_status": order.order_status,
"profit_total": float(order.profit_total or 0),
"profit_rate": float(order.profit_rate or 0),
"approve_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
"approve_result": payload["approve_result"],
"approve_opinion": payload.get("approve_opinion"),
}
except AppException:
session.rollback()
raise
except SQLAlchemyError as exc:
session.rollback()
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def cancel_approve_order(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""审批订单取消申请(通过或驳回)。
仅取消审批中或履约取消审批中的订单允许操作。
通过则订单变为已取消,驳回则恢复为取消前的状态。
驳回时必须填写审批意见。
Args:
order_id: 订单 ID
payload: 包含 approve_resultpass/refuse和 approve_opinion 的请求体
session: 数据库会话
current_user: 当前登录用户信息
Returns:
审批结果字典
被调用路由: orders.py - POST /orders/{order_id}/cancel-approve
"""
self._validate_cancel_approve_payload(payload)
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user)
self._ensure_status(order.order_status, {"cancel_pending", "cancel_fulfillment_pending"}, "当前状态不允许取消审批")
previous_status = order.order_status
target_status = "canceled" if payload["approve_result"] == "pass" else (order.cancel_previous_status or "approved")
# 设置工厂到订单
if payload.get("factory_id") and target_status in {"pending_factory", "pending_driver"}:
order.factory_id = payload["factory_id"]
self.order_repository.update_order_status(session, order, target_status)
if target_status != "canceled":
self.order_repository.clear_cancel_request(session, order)
self.order_approve_log_repository.create_log(
session,
{
"order_id": order.id,
"approve_type": "cancel_order",
"approve_result": payload["approve_result"],
"before_status": previous_status,
"after_status": target_status,
"approve_opinion": payload.get("approve_opinion"),
"operator_id": current_user.get("user_id") if current_user else None,
},
)
audit_service.write_log(
session,
{
"operate_type": "order_cancel_approve",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {
"order_status": previous_status,
"cancel_previous_status": order.cancel_previous_status,
},
"after_value": {
"order_status": order.order_status,
"approve_result": payload["approve_result"],
"approve_opinion": payload.get("approve_opinion"),
},
"remark": f"订单取消审批 {order.order_no}",
},
)
self._notify_status_change(session, order, previous_status, order.order_status, current_user.get("user_id") if current_user else None)
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_status": order.order_status,
"approve_result": payload["approve_result"],
"approve_opinion": payload.get("approve_opinion"),
}
except AppException:
session.rollback()
raise
except SQLAlchemyError as exc:
session.rollback()
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def get_supplier_text(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""生成发厂文案。
根据订单明细和供应商模板类型生成文本,仅已通过状态的订单允许操作。
Args:
order_id: 订单 ID
payload: 包含 supplier_id 的请求体(可选)
session: 数据库会话
current_user: 当前登录用户信息
Returns:
包含 text_content 的发厂文案字典
被调用路由: orders.py - POST /orders/{order_id}/supplier-text
"""
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user)
self._ensure_status(order.order_status, {"approved"}, "当前状态不允许生成发厂文案")
supplier_id = payload.get("supplier_id") or order.factory_id
if not supplier_id:
raise AppException(code=ErrorCode.NOT_FOUND, message="订单未绑定工厂", status_code=404)
supplier = self.order_repository.get_supplier(session, supplier_id)
if supplier is None:
raise AppException(code=ErrorCode.NOT_FOUND, message="工厂不存在", status_code=404)
items = self.order_repository.list_order_items(session, order.id)
template_type = supplier.template_type or "default"
text_content = self._build_supplier_text(order, items, template_type)
return {
"order_id": order.id,
"supplier_id": supplier.id,
"template_type": template_type,
"text_content": text_content,
}
except SQLAlchemyError as exc:
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def confirm_supplier_text(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""确认发厂文案。
确认后订单状态从 approved 变更为 pending_factory进入履约阶段。
自动记录发厂文案日志和审计日志。
Args:
order_id: 订单 ID
payload: 包含 text_content、supplier_id 等的请求体
session: 数据库会话
current_user: 当前登录用户信息
Returns:
确认结果字典
被调用路由: orders.py - POST /orders/{order_id}/confirm-supplier-text
"""
if not (payload.get("text_content") or "").strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="发厂文案不能为空", status_code=400)
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user)
self._ensure_status(order.order_status, {"approved"}, "当前状态不允许确认发厂")
supplier_id = payload.get("supplier_id") or order.factory_id
if not supplier_id:
raise AppException(code=ErrorCode.NOT_FOUND, message="订单未绑定工厂", status_code=404)
supplier = self.order_repository.get_supplier(session, supplier_id)
if supplier is None:
raise AppException(code=ErrorCode.NOT_FOUND, message="工厂不存在", status_code=404)
before_status = order.order_status
self.order_repository.update_supplier_text_confirm(session, order, current_user.get("user_id") if current_user else None)
self._record_supplier_text_log(session, order, {**payload, "supplier_id": supplier_id}, current_user)
audit_service.write_log(
session,
{
"operate_type": "order_supplier_text_confirm",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"order_status": before_status},
"after_value": {
"order_status": order.order_status,
"supplier_id": supplier_id,
"text_content": payload["text_content"],
"remark": payload.get("remark"),
},
"remark": f"确认发厂文案 {order.order_no}",
},
)
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_status": order.order_status,
"confirmed_at": order.supplier_text_confirmed_at.strftime("%Y-%m-%d %H:%M:%S")
if order.supplier_text_confirmed_at
else datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
}
except AppException:
session.rollback()
raise
except SQLAlchemyError as exc:
session.rollback()
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def settle_order(
self,
order_id: int,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""结算订单提成(将 pending_settle 状态转为 settled
Args:
order_id: 订单 ID
session: 数据库会话
current_user: 当前用户信息
Returns:
更新后的订单状态信息
被调用路由: POST /api/orders/{id}/settle
"""
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
raise AppException(code=ErrorCode.NOT_FOUND, message="订单不存在", status_code=404)
if order.order_status != "pending_settle":
raise AppException(
code=ErrorCode.INVALID_STATUS_TRANSITION,
message=f"当前状态 {order.order_status} 不允许结算操作",
status_code=400,
)
self.order_repository.update_order_status(session, order, "settled")
audit_service.write_log(
session,
{
"operate_type": "order_settle",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"order_status": "pending_settle"},
"after_value": {"order_status": "settled"},
"remark": f"结算订单 {order.order_no} 提成",
},
)
session.commit()
self._invalidate_order_cache(order.id)
return {"order_id": order.id, "order_status": "settled"}
except AppException:
session.rollback()
raise
except SQLAlchemyError as exc:
session.rollback()
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def update_tracking_number(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict:
"""修改自发订单快递单号。
审批通过后仍可修改,仅排除 draft/pending_approve/canceled 状态。
业务员只能修改自己订单的快递单号。
Args:
order_id: 订单 ID
payload: 包含 tracking_number、remark 的字典
session: 数据库会话
current_user: 当前登录用户信息
Returns:
更新后的订单信息字典
"""
if session is None:
raise AppException(code=ErrorCode.PARAM_ERROR, message="数据库会话未初始化", status_code=500)
order = self.order_repository.get_order(session, order_id)
if order is None:
raise AppException(code=ErrorCode.NOT_FOUND, message="订单不存在", status_code=404)
self._ensure_order_access(order, current_user)
if order.order_status in {"draft", "pending_approve", "canceled"}:
raise AppException(
code=ErrorCode.INVALID_STATUS,
message="当前状态不允许修改快递单号",
status_code=400,
)
new_tracking = (payload.get("tracking_number") or "").strip()
if not new_tracking:
raise AppException(code=ErrorCode.PARAM_ERROR, message="快递单号不能为空", status_code=400)
old_tracking = order.tracking_number or ""
order.tracking_number = new_tracking
# 如果订单处于待绑定物流单号状态,绑定快递单号后自动变为运输中
if order.order_status == "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": "order_tracking_update",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"tracking_number": old_tracking},
"after_value": {"tracking_number": new_tracking},
"remark": payload.get("remark") or f"修改快递单号:{old_tracking}{new_tracking}",
},
)
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_no": order.order_no,
"tracking_number": new_tracking,
}
def change_order_status(
self,
order_id: int,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict | None:
"""手动变更订单状态。
校验状态流转合法性(基于 _VALID_STATUS_TRANSITIONS同步欠款数据。
Args:
order_id: 订单 ID
payload: 包含 target_status 的请求体
session: 数据库会话
current_user: 当前登录用户信息
Returns:
状态变更结果字典
被调用路由: orders.py - POST /orders/{order_id}/status
"""
if session is not None:
try:
order = self.order_repository.get_order(session, order_id)
if order is None:
return None
self._ensure_order_access(order, current_user)
target_status = (payload.get("target_status") or "").strip()
if not target_status:
raise AppException(code=ErrorCode.PARAM_ERROR, message="目标状态不能为空", status_code=400)
allowed_statuses = {
"approved",
"pending_factory",
"pending_driver",
"accepted",
"picked_up",
"pending_logistics",
"in_transit",
"shipped",
"completed",
"pending_settle",
"settled",
"canceled",
}
if target_status not in allowed_statuses:
raise AppException(code=ErrorCode.PARAM_ERROR, message="目标状态不支持", status_code=400)
self._validate_status_transition(order.order_status, target_status)
# 如果是推进到待工厂状态,支持同时设置工厂
if target_status == "pending_factory" and "factory_id" in payload:
factory_id = payload.get("factory_id")
if factory_id:
order.factory_id = int(factory_id)
before_status = order.order_status
# 如果有工厂ID设置到订单上
if payload.get("factory_id") and target_status == "pending_factory":
order.factory_id = payload["factory_id"]
self.order_repository.update_order_status(session, order, target_status)
arrears_service.sync_order_arrears(session, order)
audit_service.write_log(
session,
{
"operate_type": "order_status_change",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"order_status": before_status},
"after_value": {"order_status": order.order_status},
"remark": payload.get("remark") or f"订单状态变更为 {target_status}",
},
)
self._notify_status_change(session, order, before_status, order.order_status, current_user.get("user_id") if current_user else None)
session.commit()
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_status": order.order_status,
}
except AppException:
session.rollback()
raise
except SQLAlchemyError as exc:
session.rollback()
raise AppException(code=ErrorCode.INTERNAL_ERROR, message="数据库异常", status_code=500) from exc
return None
def create_order(
self,
payload: dict,
session: Session | None = None,
current_user: dict | None = None,
) -> dict:
"""创建新订单。
若客户不存在会自动创建(基于姓名+手机号匹配)。
生成唯一订单号SO + 时间戳),计算销售/成本/利润汇总。
初始状态为 pending_approve直接进入审批流程。
Args:
payload: 包含客户信息、订单明细、费用等的请求体
session: 数据库会话
current_user: 当前登录用户信息
Returns:
创建结果字典,包含订单号、状态和汇总金额
被调用路由: orders.py - POST /orders
"""
if session is not None:
try:
if not payload["customer_name"].strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="客户姓名不能为空", status_code=400)
if not payload["customer_mobile"].strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="客户手机号不能为空", status_code=400)
customer = self.customer_repository.find_by_name_and_mobile(
session,
payload["customer_name"],
payload["customer_mobile"],
)
if customer is None:
customer = self.customer_repository.create_customer(
session,
{
"customer_name": payload["customer_name"],
"mobile": payload["customer_mobile"],
"address": payload.get("customer_address"),
"settlement_type": None,
"settlement_days": 0,
"customer_type": None,
"salesman_id": current_user.get("user_id") if current_user else payload.get("salesman_id"),
"credit_limit": 0,
"remark": "订单录入时自动创建",
},
)
items = payload["items"]
if not items:
raise AppException(code=ErrorCode.PARAM_ERROR, message="订单明细不能为空", status_code=400)
sale_total = sum(item["quantity"] * item["sale_price"] for item in items)
cost_total = sum(item["quantity"] * item["cost_price"] for item in items)
# 使用合同金额作为收入计算利润,如果没有合同金额则使用系统计算的销售金额
income_amount = payload.get("contract_amount", 0) or sale_total
profit_total = (
income_amount
- cost_total
- payload["rebate_total"]
- payload["freight_total"]
- payload["tax_total"]
- payload["other_fee_total"]
)
profit_rate = round((profit_total / income_amount) * 100, 2) if income_amount else 0
order = self.order_repository.create_order(
session,
{
"order_no": f"SO{datetime.now().strftime('%Y%m%d%H%M%S')}",
"customer_id": customer.id,
"customer_name": payload["customer_name"],
"customer_mobile": payload["customer_mobile"],
"customer_address": payload.get("customer_address"),
"salesman_id": current_user.get("user_id") if current_user else payload.get("salesman_id"),
"order_status": "pending_approve",
"order_source": payload.get("order_source"),
"delivery_type": payload.get("delivery_type"),
"factory_id": payload.get("factory_id"),
"contract_amount": payload.get("contract_amount", 0),
"sale_price_total": sale_total,
"cost_price_total": cost_total,
"rebate_total": payload["rebate_total"],
"freight_total": payload["freight_total"],
"tax_total": payload["tax_total"],
"other_fee_total": payload["other_fee_total"],
"profit_total": profit_total,
"profit_rate": profit_rate,
"commission_amount": payload["commission_amount"],
"payment_method": payload.get("payment_method"),
"tax_amount": payload.get("tax_amount", 0),
"need_invoice": payload.get("need_invoice", 0),
"order_type": payload.get("order_type"),
"self_delivery": payload.get("self_delivery", 0),
"tracking_number": payload.get("tracking_number"),
"remark": payload.get("remark"),
"customer_demand": payload.get("customer_demand"),
},
items,
)
audit_service.write_log(
session,
{
"operate_type": "order_create",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": None,
"after_value": {
"order_no": order.order_no,
"order_status": order.order_status,
"sale_price_total": round(sale_total, 2),
"cost_price_total": round(cost_total, 2),
"profit_total": round(profit_total, 2),
"profit_rate": profit_rate,
"commission_amount": payload["commission_amount"],
},
"remark": f"创建订单 {order.order_no}",
},
)
session.commit()
# 触发订单创建事件,通知管理层审批(失败不影响主流程)
try:
salesman_name = ""
if current_user:
salesman_name = current_user.get("real_name", "")
event_bus.emit("order_created", {
"order_id": order.id,
"order_no": order.order_no,
"salesman_id": order.salesman_id,
"salesman_name": salesman_name,
"customer_name": payload["customer_name"],
"amount": round(sale_total, 2),
"biz_type": "sales_order",
"biz_id": order.id,
}, session)
except Exception:
pass
self._invalidate_order_cache(order.id)
return {
"order_id": order.id,
"order_no": order.order_no,
"order_status": order.order_status,
"sale_price_total": round(sale_total, 2),
"cost_price_total": round(cost_total, 2),
"profit_total": round(profit_total, 2),
"profit_rate": profit_rate,
"commission_amount": payload["commission_amount"],
"payment_method": payload.get("payment_method"),
"tax_amount": payload.get("tax_amount", 0),
"need_invoice": payload.get("need_invoice", 0),
"order_type": payload.get("order_type"),
"self_delivery": payload.get("self_delivery", 0),
"tracking_number": payload.get("tracking_number"),
}
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_list_filters(self, filters: dict | None, current_user: dict | None = None) -> dict:
"""标准化订单查询过滤条件。
对业务员角色自动注入 salesman_id确保只能查看自己的订单。
Args:
filters: 原始过滤条件
current_user: 当前登录用户信息
Returns:
标准化后的过滤条件字典
被调用路由: orders.py - GET /orders
"""
normalized = dict(filters or {})
if not current_user:
return normalized
if current_user.get("role_code") == "salesman":
normalized["salesman_id"] = current_user.get("user_id")
return normalized
_STATUS_LABELS = {
"draft": "草稿",
"pending_approve": "待审核",
"approved": "已通过",
"rejected": "已退回",
"pending_factory": "待下发工厂",
"pending_driver": "待司机接单",
"accepted": "已接单",
"picked_up": "已揽货",
"pending_logistics": "待绑定物流单号",
"in_transit": "运输中",
"delivered": "已送达",
"canceled": "已取消",
"cancel_pending": "取消审批中",
"cancel_fulfillment_pending": "履约取消审批中",
"shipped": "已发货",
"completed": "已完成",
"pending_settle": "待结算",
"settled": "已结算",
}
def _notify_status_change(self, session: Session, order, before_status: str, after_status: str, operator_id: int | None, extra_info: dict | None = None) -> None:
"""订单状态变更时生成提醒通知。
向业务员发送状态变更通知,待审核/取消审批类状态额外通知所有经理。
Args:
session: 数据库会话
order: 订单 ORM 对象
before_status: 变更前状态
after_status: 变更后状态
operator_id: 操作人 ID
extra_info: 额外信息(如审批意见)
"""
before_label = self._STATUS_LABELS.get(before_status, before_status)
after_label = self._STATUS_LABELS.get(after_status, after_status)
_NOTIFY_MANAGER_STATUSES = {"pending_approve", "cancel_pending", "cancel_fulfillment_pending"}
_NOTIFY_SALESMAN_STATUSES = {"approved", "rejected", "canceled", "pending_logistics", "in_transit", "completed", "settled", "pending_factory", "pending_driver", "accepted", "picked_up", "shipped"}
receiver = order.salesman_id or 0
if receiver > 0 and after_status in _NOTIFY_SALESMAN_STATUSES | _NOTIFY_MANAGER_STATUSES:
reminder_content = f"订单 {order.order_no} 状态由「{before_label}」变更为「{after_label}」。"
if after_status == "rejected" and extra_info and extra_info.get("approve_opinion"):
reminder_content += f" 退回原因:{extra_info['approve_opinion']}"
self.reminder_repository.create_reminder(
session,
{
"reminder_type": "order_status_change",
"biz_type": "sales_order",
"biz_id": order.id,
"receiver_user_id": receiver,
"reminder_title": f"订单状态变更 - {order.order_no}",
"reminder_content": reminder_content,
"status": "pending",
"sent_at": datetime.now(),
},
skip_wechat=True,
)
if after_status in _NOTIFY_MANAGER_STATUSES:
managers = self._find_manager_user_ids(session)
for manager_id in managers:
self.reminder_repository.create_reminder(
session,
{
"reminder_type": "order_approval_needed",
"biz_type": "sales_order",
"biz_id": order.id,
"receiver_user_id": manager_id,
"reminder_title": f"订单待审批 - {order.order_no}",
"reminder_content": f"订单 {order.order_no}(客户:{order.customer_name})需要审批处理。",
"status": "pending",
"sent_at": datetime.now(),
},
)
# 推送 WebSocket 实时通知(失败不影响主流程)
try:
event_bus.emit("order_status_changed", {
"order_no": order.order_no,
"biz_type": "sales_order",
"biz_id": order.id,
"salesman_id": order.salesman_id or 0,
"operator_id": operator_id or 0,
"before_status_text": before_label,
"after_status_text": after_label,
})
except Exception:
pass
def _find_manager_user_ids(self, session: Session) -> list[int]:
"""查询所有经理角色的用户 ID 列表。
Args:
session: 数据库会话
Returns:
经理用户 ID 列表,查询失败时返回空列表
"""
try:
from backend.app.models.system import User as SysUser, SysRole
from sqlalchemy import select
stmt = (
select(SysUser.id)
.join(SysRole, SysRole.id == SysUser.role_id)
.where(SysRole.role_code == "manager", SysUser.status == 1, SysUser.deleted == 0)
)
return list(session.execute(stmt).scalars())
except Exception:
return []
def _ensure_order_access(self, order: object, current_user: dict | None, owner_only: bool = False) -> None:
"""校验当前用户是否有权访问指定订单。
管理员和经理可访问所有订单,业务员只能访问自己负责的订单。
Args:
order: 订单 ORM 对象
current_user: 当前登录用户信息
owner_only: 是否仅允许所有者访问(当前未使用)
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 == "salesman":
if 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 _filter_order_detail_by_role(self, order: dict, current_user: dict | None) -> dict:
"""根据用户角色过滤订单详情中的敏感字段。
司机只能看到物流相关信息,业务员看不到成本和利润信息。
Args:
order: 订单详情字典
current_user: 当前登录用户信息
Returns:
过滤后的订单详情字典
"""
if not current_user:
return order
role_code = current_user.get("role_code")
if role_code == "driver":
filtered = dict(order)
filtered["customer_name"] = None
filtered["customer_mobile"] = None
filtered["customer_address"] = None
filtered["salesman_id"] = None
filtered["salesman_name"] = ""
filtered["sale_price_total"] = None
filtered["cost_price_total"] = None
filtered["rebate_total"] = None
filtered["freight_total"] = None
filtered["tax_total"] = None
filtered["other_fee_total"] = None
filtered["profit_total"] = None
filtered["profit_rate"] = None
filtered["commission_amount"] = None
filtered["order_source"] = None
filtered["items"] = []
filtered["approve_logs"] = []
filtered["remark"] = None
return filtered
if role_code == "salesman":
filtered = dict(order)
filtered["cost_price_total"] = None
filtered["rebate_total"] = None
filtered["profit_total"] = None
filtered["profit_rate"] = None
filtered["calculation_detail"] = None
filtered["items"] = [
{
**item,
"cost_price": None,
"rebate_amount": None,
}
for item in order.get("items", [])
]
return filtered
return order
def _invalidate_order_cache(self, order_id: int | None = None) -> None:
"""清除订单相关缓存。"""
if order_id:
cache_delete(f"order:detail:{order_id}")
cache_delete_pattern("order:list:*")
cache_delete_pattern("dashboard:*")
def _filter_order_row(self, row: dict, role_code: str | None) -> dict:
"""根据用户角色过滤订单列表行中的敏感字段。
业务员看不到利润和佣金字段。
Args:
row: 订单行数据字典
role_code: 用户角色编码
Returns:
过滤后的订单行数据字典
"""
if role_code == "salesman":
filtered = dict(row)
filtered["profit_total"] = None
filtered["commission_amount"] = None
return filtered
return row
# 需求 6.1 定义的合法状态流转路径
_VALID_STATUS_TRANSITIONS = {
"draft": {"pending_approve", "canceled"},
"pending_approve": {"approved", "rejected", "canceled", "pending_factory", "in_transit"},
"rejected": {"draft", "pending_approve"},
"approved": {"pending_factory", "cancel_pending", "canceled"},
"pending_factory": {"pending_driver", "cancel_fulfillment_pending", "cancel_pending"},
"pending_driver": {"accepted", "cancel_fulfillment_pending"},
"accepted": {"picked_up", "cancel_fulfillment_pending"},
"picked_up": {"pending_logistics", "cancel_fulfillment_pending"},
"pending_logistics": {"in_transit", "cancel_fulfillment_pending"},
"in_transit": {"shipped", "cancel_fulfillment_pending"},
"shipped": {"completed", "cancel_fulfillment_pending"},
"completed": {"pending_settle", "settled"},
"pending_settle": {"settled"},
"settled": set(),
"canceled": set(),
"cancel_pending": {"canceled"},
"cancel_fulfillment_pending": {"canceled"},
}
def _validate_status_transition(self, from_status: str, to_status: str) -> None:
"""校验订单状态流转是否合法。
基于 _VALID_STATUS_TRANSITIONS 白名单校验。
Args:
from_status: 当前状态
to_status: 目标状态
Raises:
AppException: 非法流转时抛出 INVALID_STATUS
"""
allowed = self._VALID_STATUS_TRANSITIONS.get(from_status, set())
if to_status not in allowed:
raise AppException(
code=ErrorCode.INVALID_STATUS,
message=f"不允许从「{from_status}」变更为「{to_status}",
status_code=400,
)
def _ensure_status(self, current_status: str, allowed_statuses: set[str], message: str) -> None:
"""校验订单当前状态是否在允许的操作状态集合中。
Args:
current_status: 当前订单状态
allowed_statuses: 允许的操作状态集合
message: 不满足时的错误提示信息
Raises:
AppException: 状态不匹配时抛出 INVALID_STATUS
"""
if current_status not in allowed_statuses:
raise AppException(code=ErrorCode.INVALID_STATUS, message=message, status_code=400)
def _get_config_value(self, session: Session, config_key: str, default: str) -> str:
"""从系统配置表中读取配置值。
Args:
session: 数据库会话
config_key: 配置键名
default: 默认值
Returns:
配置值或默认值
"""
config = self.config_repository.get_by_key(session, config_key)
return config.config_value if config is not None and config.config_value else default
def _validate_cancel_payload(self, payload: dict) -> None:
"""校验取消订单的请求参数。
Args:
payload: 请求参数
Raises:
AppException: 取消原因为空时抛出 PARAM_ERROR
"""
if not (payload.get("cancel_reason") or "").strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="取消原因不能为空", status_code=400)
def _validate_approve_payload(self, payload: dict) -> None:
"""校验订单审批的请求参数。
审批结果仅支持 pass 或 reject驳回时必须填写意见。
Args:
payload: 请求参数
Raises:
AppException: 参数不合法时抛出 PARAM_ERROR
"""
approve_result = (payload.get("approve_result") or "").strip()
if approve_result not in {"pass", "reject"}:
raise AppException(code=ErrorCode.PARAM_ERROR, message="审批结果仅支持 pass 或 reject", status_code=400)
if approve_result == "reject" and not (payload.get("approve_opinion") or "").strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="审批驳回时必须填写意见", status_code=400)
def _validate_cancel_approve_payload(self, payload: dict) -> None:
"""校验取消审批的请求参数。
审批结果仅支持 pass 或 refuse驳回时必须填写意见。
Args:
payload: 请求参数
Raises:
AppException: 参数不合法时抛出 PARAM_ERROR
"""
approve_result = (payload.get("approve_result") or "").strip()
if approve_result not in {"pass", "refuse"}:
raise AppException(code=ErrorCode.PARAM_ERROR, message="取消审批结果仅支持 pass 或 refuse", status_code=400)
if approve_result == "refuse" and not (payload.get("approve_opinion") or "").strip():
raise AppException(code=ErrorCode.PARAM_ERROR, message="取消审批驳回时必须填写意见", status_code=400)
def _build_calculation_detail(self, order, items) -> dict:
"""构建订单利润计算明细。
包含公式说明、各费用汇总和每条明细的金额拆分,支持报价引擎快照数据。
Args:
order: 订单 ORM 对象
items: 订单明细 ORM 对象列表
Returns:
包含计算公式、汇总金额和明细的字典
"""
sale_total = float(order.sale_price_total or 0)
cost_total = float(order.cost_price_total or 0)
rebate_total = float(order.rebate_total or 0)
freight_total = float(order.freight_total or 0)
tax_total = float(order.tax_total or 0)
other_fee_total = float(order.other_fee_total or 0)
profit_total = float(order.profit_total or 0)
profit_rate = float(order.profit_rate or 0)
item_details = []
for item in items:
qty = float(item.quantity or 0)
sale_price = float(item.sale_price or 0)
cost_price = float(item.cost_price or 0)
item_sale_amount = round(qty * sale_price, 2)
item_cost_amount = round(qty * cost_price, 2)
detail_row = {
"product_name": item.product_name,
"quantity": qty,
"sale_price": sale_price,
"cost_price": cost_price,
"sale_amount": item_sale_amount,
"cost_amount": item_cost_amount,
"tax_amount": float(item.tax_amount or 0),
}
# 报价引擎快照
pricing_type = getattr(item, "pricing_type", None)
if pricing_type:
detail_row["pricing_type"] = pricing_type
detail_row["length_m"] = float(item.length_m) if getattr(item, "length_m", None) else None
detail_row["width_m"] = float(item.width_m) if getattr(item, "width_m", None) else None
detail_row["area_sqm"] = float(item.area_sqm) if getattr(item, "area_sqm", None) else None
surcharge_detail = getattr(item, "surcharge_detail", None)
if surcharge_detail:
try:
detail_row["surcharge_items"] = json.loads(surcharge_detail)
except (json.JSONDecodeError, TypeError):
detail_row["surcharge_items"] = []
detail_row["supplier_model"] = getattr(item, "supplier_model", None)
detail_row["price_tier"] = getattr(item, "price_tier", None)
item_details.append(detail_row)
# 税计算说明
tax_detail = None
if tax_total > 0:
# 反推税率tax_total / (sale_total - tax_total) 近似
base_for_tax = sale_total - tax_total if sale_total > tax_total else sale_total
approx_rate = round(tax_total / base_for_tax * 100, 2) if base_for_tax else 0
tax_detail = {
"tax_total": tax_total,
"approx_tax_rate": approx_rate,
"desc": f"税费 ¥{tax_total:.2f}(约占销售额 {approx_rate}%",
}
# 获取合同金额
contract_amount = float(order.contract_amount or 0)
return {
"formula": "利润 = 合同金额 - 成本价 - 回扣 - 运费 - 税费 - 其他费用",
"rate_formula": "利润率 = 利润 / 合同金额 x 100%",
"contract_amount": contract_amount,
"sale_total": sale_total,
"cost_total": cost_total,
"rebate_total": rebate_total,
"freight_total": freight_total,
"tax_total": tax_total,
"tax_detail": tax_detail,
"other_fee_total": other_fee_total,
"profit_total": profit_total,
"profit_rate": profit_rate,
"item_details": item_details,
}
def _build_task_summary(self, task) -> dict:
"""构建物流任务摘要字典。
Args:
task: 物流任务 ORM 对象,可为 None
Returns:
任务摘要字典task 为 None 时返回空字典
"""
if task is None:
return {}
return {
"task_id": task.id,
"task_no": task.task_no,
"driver_id": task.driver_id,
"factory_id": task.factory_id,
"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,
"express_company": task.express_company,
"remark": task.remark,
"created_at": task.created_at.strftime("%Y-%m-%d %H:%M:%S") if task.created_at else "",
}
def _build_attachment_row(self, attachment) -> dict:
"""构建附件行数据字典。
Args:
attachment: 附件 ORM 对象
Returns:
附件摘要字典
"""
return {
"attachment_id": attachment.id,
"biz_type": attachment.biz_type,
"biz_id": attachment.biz_id,
"file_name": attachment.file_name,
"file_url": attachment.file_url,
"file_type": attachment.file_type,
"file_size": attachment.file_size,
"file_category": attachment.file_category,
"object_key": attachment.object_key,
"created_at": attachment.created_at.strftime("%Y-%m-%d %H:%M:%S") if attachment.created_at else "",
}
def _record_supplier_text_log(self, session: Session, order, payload: dict, current_user: dict | None) -> None:
"""记录发厂文案确认日志。
Args:
session: 数据库会话
order: 订单 ORM 对象
payload: 包含 supplier_id、text_content 等的请求体
current_user: 当前登录用户信息
"""
self.order_repository.create_supplier_text_log(
session,
{
"order_id": order.id,
"supplier_id": payload["supplier_id"],
"template_type": payload.get("template_type") or "default",
"text_content": payload["text_content"],
"confirmed": 1,
"confirmed_by": current_user.get("user_id") if current_user else None,
"confirmed_at": datetime.now(),
"created_by": current_user.get("user_id") if current_user else None,
},
)
def _build_supplier_text(self, order, items, template_type: str) -> str:
"""根据模板类型生成发厂文案文本。
支持 default完整模板和 brief简略模板两种格式。
Args:
order: 订单 ORM 对象
items: 订单明细 ORM 对象列表
template_type: 模板类型default 或 brief
Returns:
格式化后的发厂文案字符串
"""
item_lines = []
for item in items:
qty = float(item.quantity or 0)
spec = item.specification or ""
unit = item.unit or ""
item_lines.append(f" {item.product_name} {spec} {qty}{unit}")
items_text = "\n".join(item_lines) if item_lines else " (无明细)"
if template_type == "brief":
total_qty = sum(float(item.quantity or 0) for item in items)
return f"订单 {order.order_no}\n客户:{order.customer_name}\n总数量:{total_qty}\n请安排发货。"
return (
f"订单编号:{order.order_no}\n"
f"客户:{order.customer_name}\n"
f"发货方式:{order.delivery_type or '待定'}\n"
f"产品明细:\n{items_text}\n"
f"备注:{order.remark or ''}\n"
f"请安排生产与发货。"
)
# ------------------------------------------------------------------
# 批量匹配物流
# ------------------------------------------------------------------
_TRACKING_RE = re.compile(r'(?<!\d)(\d{10,20})(?!\d)')
_RECIPIENT_PREFIX_RE = re.compile(r'收件人\s*[:]?\s*')
def parse_logistics_text(self, text):
results = []
for line in text.strip().splitlines():
line = line.strip()
if not line:
continue
tracking_match = self._TRACKING_RE.search(line)
if not tracking_match:
continue
tracking_number = tracking_match.group(1)
before = line[:tracking_match.start()].strip()
after = line[tracking_match.end():].strip()
express_company = before if before else None
recipient_name = self._RECIPIENT_PREFIX_RE.sub('', after).strip()
recipient_name = recipient_name.strip(':.。、 ') or None
results.append({
"tracking_number": tracking_number,
"express_company": express_company,
"recipient_name": recipient_name,
})
return results
def parse_logistics_table_images(self, image_urls, session):
"""识别物流表格截图,用 LLM 从 OCR 文本中提取快递单号和收件人。"""
from backend.app.services.ai_service import ai_service
from backend.app.core.config import get_settings
settings = get_settings()
results = []
seen_tracking = set()
for url in image_urls:
try:
ctx = ai_service._ocr_and_parse_image(url)
raw_text = ctx.get("raw_text", "")
if not raw_text.strip():
continue
# 用 LLM 从表格文本中提取快递单号和收件人
llm_items = self._llm_extract_logistics(raw_text, settings)
for item in llm_items:
tn = (item.get("tracking_number") or "").strip()
if not tn or len(tn) < 10 or tn in seen_tracking:
continue
seen_tracking.add(tn)
results.append({
"tracking_number": tn,
"express_company": item.get("express_company") or None,
"recipient_name": item.get("recipient_name") or None,
})
except Exception:
pass
return results
def _llm_extract_logistics(self, text, settings):
"""调用 LLM 从物流表格文本中提取快递单号和收件人。"""
import json as _json
from urllib import request as urllib_request
system_prompt = (
"你是物流信息提取助手。从以下物流表格文本中提取每一行的快递单号和收件人姓名。\n"
"表格通常包含:运单号、运单状态、件数、体积、计费重量、付款方式、运费、"
"保价费、包装服务费、信息费、代收货款、签收费、声明价值、产品类型、"
"增值服务、服务方式、业务属性、托寄物、寄件人、收件人、目的网点 等列。\n\n"
"严格按以下 JSON 数组格式输出,不要添加任何其他内容:\n"
'[{"tracking_number":"运单号","express_company":"快递公司名或null","recipient_name":"收件人姓名"}]'
)
payload = _json.dumps({
"model": settings.llm_parse_model,
"messages": [
{"role": "system", "content": system_prompt},
{"role": "user", "content": f"请从以下物流表格文本中提取快递单号和收件人:\n\n{text}"},
],
"temperature": 0.1,
"max_tokens": 2048,
}).encode("utf-8")
api_key = settings.llm_parse_api_key or settings.aliyun_ai_access_key_id
headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {api_key}",
}
req = urllib_request.Request(
url=settings.llm_parse_api_url, data=payload, headers=headers, method="POST"
)
with urllib_request.urlopen(req, timeout=30) as resp:
result = _json.loads(resp.read().decode("utf-8"))
content = result["choices"][0]["message"]["content"]
return self._extract_json_array(content)
def _extract_json_array(self, text):
"""从 LLM 响应中提取 JSON 数组。"""
import json as _json
text = text.strip()
if text.startswith("```"):
text = text.split("\n", 1)[1]
text = text.rsplit("```", 1)[0]
text = text.strip()
if text.startswith("json"):
text = text[4:].strip()
try:
data = _json.loads(text)
return data if isinstance(data, list) else []
except _json.JSONDecodeError:
# 尝试截断到最后一个 ]
last_bracket = text.rfind("]")
while last_bracket > 0:
try:
return _json.loads(text[:last_bracket + 1])
except _json.JSONDecodeError:
last_bracket = text.rfind("]", 0, last_bracket - 1)
return []
def parse_express_images(self, image_urls, session):
"""Parse express waybill photos via kuaidi100 OCR."""
from backend.app.services.logistics_service import logistics_service
results = []
for url in image_urls:
try:
result = logistics_service.recognize_waybill({"image_url": url}, session)
tracking_number = result.get("tracking_number") or result.get("no") or None
if not tracking_number:
continue
express_company = result.get("express_company") or result.get("com") or None
results.append({
"tracking_number": tracking_number,
"express_company": express_company,
"recipient_name": None,
})
except Exception:
pass
return results
def _match_orders_by_name(self, name, session):
"""Match orders by exact customer_name."""
from backend.app.models.business import SalesOrder
orders = session.query(SalesOrder).filter(
SalesOrder.customer_name == name,
SalesOrder.deleted == 0,
).all()
return [
{
"order_id": o.id,
"order_no": o.order_no,
"customer_name": o.customer_name,
"order_status": o.order_status,
"tracking_number": o.tracking_number,
}
for o in orders
]
def confirm_batch_match(self, matches, session):
"""批量确认物流匹配结果,为订单写入快递单号。
Args:
matches: 匹配结果列表,每项包含 order_id 和 tracking_number
session: 数据库会话
Returns:
包含 success_count、fail_count 和 failures 的结果字典
"""
success_count = 0
fail_count = 0
failures = []
for m in matches:
order_id = m.get("order_id")
tracking_number = (m.get("tracking_number") or "").strip()
if not tracking_number:
fail_count += 1
failures.append({"tracking_number": tracking_number, "order_id": order_id, "reason": "快递单号为空"})
continue
try:
order = session.query(SalesOrder).filter(
SalesOrder.id == order_id, SalesOrder.deleted == 0
).first()
if order is None:
fail_count += 1
failures.append({"tracking_number": tracking_number, "order_id": order_id, "reason": "订单不存在"})
continue
# 唯一性检查:其他未删除订单是否已占用该快递单号
conflict = session.query(SalesOrder).filter(
SalesOrder.tracking_number == tracking_number,
SalesOrder.id != order_id,
SalesOrder.deleted == 0,
).first()
if conflict:
fail_count += 1
failures.append({"tracking_number": tracking_number, "order_id": order_id, "reason": f"快递单号已绑定订单 {conflict.order_no}"})
continue
old_tracking = order.tracking_number or ""
order.tracking_number = tracking_number
audit_service.write_log(
session,
{
"operate_type": "order_tracking_batch_update",
"biz_type": "sales_order",
"biz_id": order.id,
"before_value": {"tracking_number": old_tracking},
"after_value": {"tracking_number": tracking_number},
"remark": f"批量匹配快递单号:{old_tracking} -> {tracking_number}",
},
)
self._invalidate_order_cache(order_id)
success_count += 1
except Exception as e:
session.rollback()
fail_count += 1
failures.append({"tracking_number": tracking_number, "order_id": order_id, "reason": str(e)})
return {"success_count": success_count, "fail_count": fail_count, "failures": failures}
order_service = OrderService()