"""销售订单服务层。 负责销售订单的全生命周期管理,包括创建、编辑、提交审核、审批、取消、 发厂文案生成、状态流转等核心业务逻辑。 依赖多个 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.audit_repository import AuditRepository 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 from backend.app.services.pricing_engine import pricing_engine 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() self.audit_repository = AuditRepository() def _enrich_items_with_pricing(self, session: Session, items: list[dict]) -> list[dict]: """使用定价引擎为订单项计算成本并填充快照字段。 对每个商品项: 1. 查找对应的定价规则 (ProductPricingRule) 2. 调用定价引擎计算 3. 用引擎结果覆盖 cost_price 并填充快照字段 如果找不到规则或计算失败,保留原始 cost_price 不变。 Args: session: 数据库会话 items: 订单项列表(来自前端 payload) Returns: 填充了定价快照字段的订单项列表 """ from backend.app.models.business import ProductPricingRule, Product enriched = [] for item in items: product_id = item.get("product_id") if not product_id: enriched.append(item) continue # 查找定价规则 rule = session.query(ProductPricingRule).filter( ProductPricingRule.product_id == product_id, ProductPricingRule.deleted == 0, ProductPricingRule.status == 1, ).first() if not rule: enriched.append(item) continue # 获取产品属性(厚度、克重等) product = session.query(Product).filter(Product.id == product_id).first() product_attrs = {} if product: if product.thickness: product_attrs["thickness"] = product.thickness if product.weight_gsm: product_attrs["weight_gsm"] = product.weight_gsm if product.default_width_m: product_attrs["default_width_m"] = float(product.default_width_m) # 构建用户输入:从 item 中提取定价相关字段 user_inputs = {} lm = item.get("length_m") wm = item.get("width_m") if lm and wm: user_inputs["length"] = lm user_inputs["length_unit"] = "m" user_inputs["width"] = wm user_inputs["width_unit"] = "m" else: # 前端未传尺寸时,从需求规格解析 demand_spec = item.get("demand_specification") parsed = self._parse_demand_spec(demand_spec) if demand_spec else {} if parsed: user_inputs["length"] = parsed["length_m"] user_inputs["length_unit"] = "m" user_inputs["width"] = parsed["width_m"] user_inputs["width_unit"] = "m" if item.get("quantity"): user_inputs["quantity"] = item["quantity"] if item.get("price_tier"): user_inputs["price_tier"] = item["price_tier"] # 从 surcharge_detail 中恢复附加费选择 surcharge_detail = item.get("surcharge_detail") if surcharge_detail: try: surcharges = json.loads(surcharge_detail) if isinstance(surcharge_detail, str) else surcharge_detail if isinstance(surcharges, list): for s in surcharges: key = s.get("key") if key: user_inputs[f"surcharge_{key}"] = True except (json.JSONDecodeError, TypeError): pass try: result = pricing_engine.calculate_for_order_item(rule, user_inputs, product_attrs) item_copy = dict(item) # 使用定价引擎计算的 cost_price(包含单位换算、公式计算、附加费) snapshot = result.get("item_snapshot", {}) if snapshot.get("cost_price") is not None: item_copy["cost_price"] = snapshot["cost_price"] if snapshot.get("pricing_type"): item_copy["pricing_type"] = snapshot["pricing_type"] if result.get("area_sqm"): item_copy["area_sqm"] = result["area_sqm"] if snapshot.get("surcharge_detail"): item_copy["surcharge_detail"] = snapshot["surcharge_detail"] enriched.append(item_copy) except Exception: # 计价失败时保留原始数据 enriched.append(item) return enriched @staticmethod def _parse_demand_spec(demand_spec: str) -> dict: """从需求规格字符串解析尺寸(宽×长),返回 {"width_m": x, "length_m": y}。 支持格式: "0.02m×10m", "25cm×30cm", "250mm×300mm" 等。 """ import re if not demand_spec: return {} match = re.match( r'([\d.]+)\s*(cm|毫米|mm|m|米)?\s*[×xX*]\s*([\d.]+)\s*(cm|毫米|mm|m|米)?', demand_spec.strip(), ) if not match: return {} v1, u1, v2, u2 = match.group(1), match.group(2), match.group(3), match.group(4) unit_map = {"m": 1.0, "米": 1.0, "cm": 0.01, "厘米": 0.01, "mm": 0.001, "毫米": 0.001} m1 = float(v1) * unit_map.get(u1, 0.01) m2 = float(v2) * unit_map.get(u2, 0.01) return {"width_m": m1, "length_m": m2} def _enrich_items_from_db(self, session: Session, items: list) -> list[dict]: """读取数据库订单项,对 cost_price=0 的项用定价引擎重算。 解决旧订单在定价引擎集成前创建、cost_price 存为 0 的问题。 仅在 cost_price 为 0 且存在对应定价规则时触发计算。 Args: session: 数据库会话 items: SalesOrderItem ORM 对象列表 Returns: 序列化后的订单项字典列表 """ from backend.app.models.business import ProductPricingRule, Product # 一次性加载所有有定价规则的产品属性和规则 product_ids = [item.product_id for item in items if item.product_id] rules_map = {} products_map = {} if product_ids: for rule in session.query(ProductPricingRule).filter( ProductPricingRule.product_id.in_(product_ids), ProductPricingRule.deleted == 0, ProductPricingRule.status == 1, ).all(): rules_map[rule.product_id] = rule for product in session.query(Product).filter(Product.id.in_(product_ids)).all(): products_map[product.id] = product result = [] for item in items: item_dict = { "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), } # cost_price 为 0 且有定价规则时,尝试用引擎重算 if item_dict["cost_price"] == 0 and item.product_id and item.product_id in rules_map: rule = rules_map[item.product_id] product = products_map.get(item.product_id) product_attrs = {} if product: if product.thickness: product_attrs["thickness"] = product.thickness if product.weight_gsm: product_attrs["weight_gsm"] = product.weight_gsm if product.default_width_m: product_attrs["default_width_m"] = float(product.default_width_m) user_inputs = {} lm = getattr(item, "length_m", None) wm = getattr(item, "width_m", None) if lm and wm: user_inputs["length"] = float(lm) user_inputs["length_unit"] = "m" user_inputs["width"] = float(wm) user_inputs["width_unit"] = "m" else: demand_spec = getattr(item, "demand_specification", None) parsed = self._parse_demand_spec(demand_spec) if demand_spec else {} if parsed: user_inputs["length"] = parsed["length_m"] user_inputs["length_unit"] = "m" user_inputs["width"] = parsed["width_m"] user_inputs["width_unit"] = "m" if item.quantity: user_inputs["quantity"] = float(item.quantity) if getattr(item, "price_tier", None): user_inputs["price_tier"] = item.price_tier surcharge_detail = getattr(item, "surcharge_detail", None) if surcharge_detail: try: surcharges = json.loads(surcharge_detail) if isinstance(surcharge_detail, str) else surcharge_detail if isinstance(surcharges, list): for s in surcharges: key = s.get("key") if key: user_inputs[f"surcharge_{key}"] = True except (json.JSONDecodeError, TypeError): pass try: calc_result = pricing_engine.calculate_for_order_item(rule, user_inputs, product_attrs) snapshot = calc_result.get("item_snapshot", {}) if snapshot.get("cost_price") is not None and snapshot["cost_price"] > 0: item_dict["cost_price"] = snapshot["cost_price"] item_dict["pricing_type"] = snapshot.get("pricing_type") or item_dict["pricing_type"] if calc_result.get("area_sqm"): item_dict["area_sqm"] = calc_result["area_sqm"] # 回写数据库,避免下次重复计算 try: item.cost_price = snapshot["cost_price"] if snapshot.get("pricing_type"): item.pricing_type = snapshot["pricing_type"] if calc_result.get("area_sqm"): item.area_sqm = calc_result["area_sqm"] session.flush() except Exception: pass except Exception: pass # 查询定价规则公式信息供前端展示(不受 cost_price 条件限制) if item.product_id and item.product_id in rules_map: rule = rules_map[item.product_id] item_dict["formula_expr"] = getattr(rule, "formula_expr", None) item_dict["formula_note"] = getattr(rule, "formula_note", None) item_dict["base_unit_price"] = float(rule.base_unit_price or 0) item_dict["pricing_unit"] = getattr(rule, "pricing_unit", None) # 构建人类可读的计算过程 if not item_dict.get("formula_detail"): try: unit_price = float(rule.base_unit_price or 0) pricing_type = getattr(rule, "pricing_type", "area") or "area" pricing_unit = getattr(rule, "pricing_unit", "㎡") or "㎡" qty = float(item.quantity or 0) length = float(getattr(item, "length_m", 0) or 0) width = float(getattr(item, "width_m", 0) or 0) if pricing_type == "area" and length and width and unit_price: area = length * width item_dict["formula_detail"] = f"{length}m × {width}m × ¥{unit_price}/{pricing_unit} = ¥{area * unit_price:.2f}" elif pricing_type in ("kg", "weight_g") and qty and unit_price: item_dict["formula_detail"] = f"{qty}{pricing_unit} × ¥{unit_price}/{pricing_unit} = ¥{qty * unit_price:.2f}" elif unit_price: item_dict["formula_detail"] = f"{qty}{pricing_unit} × ¥{unit_price}/{pricing_unit} = ¥{qty * unit_price:.2f}" except Exception: pass result.append(item_dict) return result 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 [] # 查询运单子表 waybills = [] if task: waybills = self.logistics_repository.list_waybills_by_task_id(session, task.id) # 查询业务员姓名 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 "" # 构建物流任务摘要,包含回退逻辑 task_summary = self._build_task_summary_with_fallback(task, waybills, order) 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": self._enrich_items_from_db(session, 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": task_summary, "waybills": [ { "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 ], "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, session), "profit_alert_threshold": float(self._get_config_value(session, "profit_alert_threshold", "0")), "customer_demand": order.customer_demand, "remark": order.remark, "revision_info": self._get_latest_revision(session, order.id), } # 写缓存(缓存完整结果,角色过滤在读取后执行) 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) # 调用定价引擎为每个订单项计算成本并生成快照 items = self._enrich_items_with_pricing(session, items) 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 revise_order( self, order_id: int, payload: dict, session: Session | None = None, current_user: dict | None = None, ) -> dict | None: """管理员修订订单(修改商品后自动重算)。 与 update_order 不同,此方法: 1. 记录修改前的快照 2. 调用定价引擎重新计算所有商品成本 3. 生成修订版本记录 4. 返回修改前后差异 仅 manager/admin 角色可调用。 Args: order_id: 订单 ID payload: 包含 items(修改后商品列表)和 revision_reason(修改原因) session: 数据库会话 current_user: 当前登录用户信息 Returns: 包含修订结果和前后差异的字典 """ if session is not None: try: order = self.order_repository.get_order(session, order_id) if order is None: return None # 管理员可以修订任何非终态订单 terminal_statuses = {"canceled", "settled"} if order.order_status in terminal_statuses: raise AppException( code=ErrorCode.PARAM_ERROR, message="已取消或已结算的订单不允许修订", status_code=400, ) items = payload.get("items", []) if not items: raise AppException(code=ErrorCode.PARAM_ERROR, message="订单明细不能为空", status_code=400) # 记录修改前的快照 old_items = self.order_repository.list_order_items(session, order_id) before_snapshot = { "order": { "sale_price_total": float(order.sale_price_total or 0), "cost_price_total": float(order.cost_price_total or 0), "profit_total": float(order.profit_total or 0), "profit_rate": float(order.profit_rate or 0), }, "items": [ { "product_name": item.product_name, "specification": item.specification, "quantity": float(item.quantity or 0), "cost_price": float(item.cost_price or 0), "sale_price": float(item.sale_price or 0), } for item in old_items ], } # 调用定价引擎重新计算 items = self._enrich_items_with_pricing(session, items) 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") or order.contract_amount or sale_total rebate_total = payload.get("rebate_total", float(order.rebate_total or 0)) freight_total = payload.get("freight_total", float(order.freight_total or 0)) tax_total = payload.get("tax_total", float(order.tax_total or 0)) other_fee_total = payload.get("other_fee_total", float(order.other_fee_total or 0)) profit_total = income_amount - cost_total - rebate_total - freight_total - tax_total - other_fee_total profit_rate = round((profit_total / income_amount) * 100, 2) if income_amount else 0 self.order_repository.update_order_with_items( session, order, { "sale_price_total": sale_total, "cost_price_total": cost_total, "rebate_total": rebate_total, "freight_total": freight_total, "tax_total": tax_total, "other_fee_total": other_fee_total, "profit_total": profit_total, "profit_rate": profit_rate, }, items, ) # 记录修改后的快照 after_snapshot = { "order": { "sale_price_total": round(sale_total, 2), "cost_price_total": round(cost_total, 2), "profit_total": round(profit_total, 2), "profit_rate": profit_rate, }, "items": [ { "product_name": item.get("product_name"), "specification": item.get("specification"), "quantity": item.get("quantity"), "cost_price": item.get("cost_price"), "sale_price": item.get("sale_price"), } for item in items ], } # 写审计日志 audit_service.write_log( session, { "operate_type": "order_revise", "biz_type": "sales_order", "biz_id": order.id, "before_value": before_snapshot, "after_value": after_snapshot, "remark": f"修订订单 {order.order_no}:{payload.get('revision_reason', '管理员修改')}", }, ) # 写审批日志(使修订记录出现在审批时间线中) self.order_approve_log_repository.create_log( session, { "order_id": order.id, "approve_type": "revise", "approve_result": "revised", "before_status": order.order_status, "after_status": order.order_status, "approve_opinion": payload.get("revision_reason", "管理员修订"), "operator_id": current_user.get("user_id") if current_user else None, }, ) 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, "revision_reason": payload.get("revision_reason", ""), "before": before_snapshot, "after": after_snapshot, } 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_driver,accepted,picked_up,pending_logistics,in_transit,shipped") 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_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_result(pass/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("driver_id"): # 有司机:审批通过并直接创建物流任务,跳到待司机接单 target_status = "pending_driver" else: # 无司机:直接进入待绑定物流单号状态 target_status = "pending_logistics" else: target_status = "rejected" # 如果有工厂ID,设置到订单上 if payload.get("factory_id") and target_status in ("pending_driver", "pending_logistics"): order.factory_id = payload["factory_id"] 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 = "" if payload.get("factory_id"): 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.get("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_result(pass/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_driver", "pending_logistics"}: 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_logistics,进入履约阶段。 自动记录发厂文案日志和审计日志。 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 in ("picked_up", "pending_logistics"): before_status = order.order_status self.order_repository.update_order_status(session, order, "in_transit") self._notify_status_change(session, order, before_status, "in_transit", current_user.get("user_id") if current_user else None) audit_service.write_log( session, { "operate_type": "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_driver", "accepted", "picked_up", "pending_logistics", "in_transit", "shipped", "delivered", "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) before_status = order.order_status # 如果有工厂ID,设置到订单上 if payload.get("factory_id") and target_status in ("pending_driver", "pending_logistics"): 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) # 调用定价引擎为每个订单项计算成本并生成快照 items = self._enrich_items_with_pricing(session, items) 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}", }, ) # 触发订单创建事件,通知管理层审批(失败不影响主流程) 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 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 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_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: """订单状态变更时生成提醒通知。 根据状态类型选择合适的事件: - 需要审批的状态:触发 order_created 事件(通知管理员) - 其他状态变更:触发 order_status_changed 事件(通知业务员) 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"} # 推送实时通知(失败不影响主流程) try: if after_status in _NOTIFY_MANAGER_STATUSES: # 需要审批的状态,通知管理员 salesman_name = "" if order.salesman_id: from backend.app.models.system import User salesman = session.query(User).filter(User.id == order.salesman_id).first() salesman_name = salesman.real_name if salesman else "" event_bus.emit("order_created", { "order_id": order.id, "order_no": order.order_no, "salesman_id": order.salesman_id or 0, "salesman_name": salesman_name, "customer_name": order.customer_name or "", "amount": float(order.sale_price_total or 0), "biz_type": "sales_order", "biz_id": order.id, }) else: # 其他状态变更,通知业务员 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, Role as SysRole from sqlalchemy import select stmt = ( select(SysUser.id) .join(SysRole, SysRole.id == SysUser.role_id) .where(SysRole.role_code == "manager", SysUser.status == 1) ) 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_logistics", "in_transit"}, "rejected": {"draft", "pending_approve"}, "approved": {"pending_driver", "pending_logistics", "cancel_pending", "canceled"}, "pending_driver": {"accepted", "cancel_fulfillment_pending"}, "accepted": {"picked_up", "cancel_fulfillment_pending"}, "picked_up": {"pending_logistics", "in_transit", "cancel_fulfillment_pending"}, "pending_logistics": {"in_transit", "completed", "cancel_fulfillment_pending"}, "in_transit": {"delivered", "shipped", "cancel_fulfillment_pending"}, "shipped": {"delivered", "completed", "cancel_fulfillment_pending"}, "delivered": {"completed", "pending_settle", "settled"}, "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, session=None) -> dict: """构建订单利润计算明细。 包含公式说明、各费用汇总和每条明细的金额拆分,支持报价引擎快照数据。 Args: order: 订单 ORM 对象 items: 订单明细 ORM 对象列表 session: 数据库会话(可选,用于查询定价规则公式) 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) # 从数据库读取快照字段 pricing_type = getattr(item, "pricing_type", None) area_sqm = float(item.area_sqm) if getattr(item, "area_sqm", None) else 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 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": pricing_type, "area_sqm": area_sqm, "length_m": length_m, "width_m": width_m, } # 附加费等附加快照 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) # 查询定价规则公式 if session and item.product_id: try: from backend.app.models.business import ProductPricingRule rule = session.query(ProductPricingRule).filter( ProductPricingRule.product_id == item.product_id, ProductPricingRule.deleted == 0, ).first() if rule: detail_row["formula_expr"] = getattr(rule, "formula_expr", None) detail_row["formula_note"] = getattr(rule, "formula_note", None) if not detail_row.get("pricing_type"): detail_row["pricing_type"] = getattr(rule, "pricing_type", None) except Exception: pass 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_task_summary_with_fallback(self, task, waybills: list, order) -> dict: """构建物流任务摘要字典,支持单号回退。 回退顺序:logistics_task -> 第一条 logistics_waybill -> sales_order Args: task: 物流任务 ORM 对象,可为 None waybills: 运单子表列表 order: 订单 ORM 对象 Returns: 任务摘要字典,task 为 None 时返回空字典 """ if task is None: return {} # 单号回退逻辑 tracking_number = task.tracking_number express_company = task.express_company if not tracking_number and waybills: tracking_number = waybills[0].tracking_number express_company = waybills[0].express_company if not tracking_number and order: tracking_number = order.tracking_number 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": tracking_number, "express_company": express_company, "remark": task.remark, "created_at": task.created_at.strftime("%Y-%m-%d %H:%M:%S") if task.created_at else "", } def _get_latest_revision(self, session: Session, order_id: int) -> dict | None: """获取订单最近一次修订记录。 从审计日志中查找 operate_type=order_revise 的最新记录, 返回修订前后的差异快照,供审批页展示。 Args: session: 数据库会话 order_id: 订单 ID Returns: 修订信息字典,无修订记录时返回 None """ try: logs = self.audit_repository.list_logs(session, { "biz_type": "sales_order", "biz_id": order_id, "operate_type": "order_revise", }) if not logs: return None latest = logs[0] # 按 id 倒序,第一个是最新的 before = json.loads(latest.before_value) if latest.before_value else None after = json.loads(latest.after_value) if latest.after_value else None return { "revision_time": latest.operate_time.strftime("%Y-%m-%d %H:%M:%S") if latest.operate_time else "", "revision_remark": latest.remark or "", "before": before, "after": after, } except Exception: return None 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'(? 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}", }, ) # 同步更新物流任务表(如果存在) task = self.logistics_repository.get_task_by_order_id(session, order.id) if task: task.tracking_number = tracking_number express_company = (m.get("express_company") or "").strip() if express_company: task.express_company = express_company # 如果订单处于待物流状态,推进到运输中 if order.order_status == "pending_logistics": before_status = order.order_status self.order_repository.update_order_status(session, order, "in_transit") from backend.app.services.logistics_service import logistics_service logistics_service._notify_status_change(session, order, before_status, "in_transit", None) self._invalidate_order_cache(order_id) cache_delete_pattern("logistics:*") cache_delete_pattern(f"order:detail:{order_id}") cache_delete_pattern("order:list:*") cache_delete_pattern("dashboard:*") 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()