import re from datetime import date, datetime, timedelta from io import BytesIO from urllib.parse import quote from fastapi import APIRouter, Depends, File, HTTPException, Query, UploadFile from fastapi.responses import Response from openpyxl import Workbook, load_workbook from sqlalchemy.exc import IntegrityError from sqlalchemy import and_, delete as sa_delete, func, inspect, or_, select, text, update from sqlalchemy.orm import Session from app.db.session import get_db from app.models.master_data import Item, Material, Warehouse from app.models.operations import PurchaseOrder, PurchaseOrderItem, PurchaseOrderSalesOrderLink, PurchaseReceipt, PurchaseReceiptItem, StockLot, Supplier from app.models.planning import MaterialDemand from app.models.sales import SalesOrder, SalesOrderItem from app.schemas.operations import ( PurchaseOrderCreate, PurchaseOrderItemRead, PurchaseOrderRead, PurchaseOrderUpdate, PurchaseReceiptCreate, PurchaseReceiptItemRead, PurchaseReceiptRead, SupplierCreate, SupplierRead, SupplierUpdate, ) from app.services.auth import AuthContext, require_authenticated_user from app.services.contacts import normalize_contacts, parse_contacts, serialize_contacts from app.services.document_archives import DOCUMENT_TYPE_PURCHASE_ORDER, DOCUMENT_TYPE_PURCHASE_RECEIPT, generate_document_archive from app.services.logistics import normalize_logistics_fields from app.services.operations import ( build_inventory_lot_no, build_supplier_code, create_inventory_txn, get_location_or_default, get_purchase_order_items_query, get_purchase_order_query, get_purchase_receipt_items_query, get_purchase_receipt_query, sync_purchase_order_item_status, sync_purchase_order_status, to_decimal, upsert_stock_balance, ) from app.services.purchase_order_links import ensure_purchase_order_sales_link_table from app.services.stocktake import ensure_warehouses_unlocked from app.services.system_permissions import ensure_employee_has_permission router = APIRouter(dependencies=[Depends(require_authenticated_user)]) SUPPLIER_EXCEL_HEADERS = ["供应商名称", "联系人1", "联系电话1", "联系人2", "联系电话2", "联系人3", "联系电话3", "地址"] SUPPLIER_EXCEL_REQUIRED_HEADERS = ["供应商名称", "联系人1", "联系电话1", "地址"] LEGACY_SUPPLIER_EXCEL_HEADERS = ["供应商名称", "联系人", "联系电话", "地址"] PURCHASE_TARGET_WAREHOUSE_TYPES = {"RAW", "AUX"} RECEIVABLE_PURCHASE_ORDER_STATUSES = {"LOCKED", "PENDING_QC", "PARTIAL", "PARTIAL_RECEIVED", "REJECTED"} def _auth_context_user_id(context: AuthContext | object) -> int | None: user = getattr(context, "user", None) user_id = getattr(user, "id", None) return int(user_id) if user_id is not None else None def build_purchase_order_no(db: Session, supplier: Supplier) -> str: _ = supplier year_part = date.today().strftime("%Y") prefix = f"采购{year_part}" pattern = re.compile(rf"{re.escape(prefix)}-(\d{{5}})") max_no = 0 rows = db.scalars(select(PurchaseOrder.po_no).where(PurchaseOrder.po_no.like(f"{prefix}-%"))).all() for po_no in rows: match = pattern.fullmatch(str(po_no or "").strip()) if match: max_no = max(max_no, int(match.group(1))) return f"{prefix}-{max_no + 1:05d}" def build_purchase_receipt_no(db: Session, purchase_order: PurchaseOrder) -> str: _ = purchase_order year_part = date.today().strftime("%Y") prefix = f"入库{year_part}" pattern = re.compile(rf"{re.escape(prefix)}-(\d{{5}})") max_no = 0 rows = db.scalars(select(PurchaseReceipt.receipt_no).where(PurchaseReceipt.receipt_no.like(f"{prefix}-%"))).all() for receipt_no in rows: match = pattern.fullmatch(str(receipt_no or "").strip()) if match: max_no = max(max_no, int(match.group(1))) return f"{prefix}-{max_no + 1:05d}" def _clean_supplier_text(value: object | None) -> str: return str(value or "").strip() def _build_supplier_workbook(rows: list[Supplier]) -> Workbook: workbook = Workbook() worksheet = workbook.active worksheet.title = "供应商名录" worksheet.freeze_panes = "A2" worksheet.append(SUPPLIER_EXCEL_HEADERS) for row in rows: contacts = parse_contacts(row.contact_name, row.contact_phone) worksheet.append([ row.supplier_name, contacts[0]["contact_name"] if len(contacts) > 0 else "", contacts[0]["contact_phone"] if len(contacts) > 0 else "", contacts[1]["contact_name"] if len(contacts) > 1 else "", contacts[1]["contact_phone"] if len(contacts) > 1 else "", contacts[2]["contact_name"] if len(contacts) > 2 else "", contacts[2]["contact_phone"] if len(contacts) > 2 else "", row.address or "", ]) for index, width in enumerate([30, 18, 20, 18, 20, 18, 20, 46], start=1): worksheet.column_dimensions[worksheet.cell(row=1, column=index).column_letter].width = width return workbook def _supplier_excel_response(workbook: Workbook, filename: str) -> Response: output = BytesIO() workbook.save(output) encoded_filename = quote(filename) return Response( content=output.getvalue(), media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", headers={"Content-Disposition": f"attachment; filename=\"suppliers.xlsx\"; filename*=UTF-8''{encoded_filename}"}, ) def _read_supplier_excel_rows(content: bytes) -> list[dict[str, object]]: try: workbook = load_workbook(BytesIO(content), read_only=True, data_only=True) except Exception as exc: raise ValueError("Excel 文件读取失败,请确认上传的是供应商名录 xlsx 文件") from exc worksheet = workbook.active rows = list(worksheet.iter_rows(values_only=True)) if not rows: return [] headers = [_clean_supplier_text(value) for value in rows[0]] accepted_headers = SUPPLIER_EXCEL_REQUIRED_HEADERS if "联系人1" in headers else LEGACY_SUPPLIER_EXCEL_HEADERS missing = [header for header in accepted_headers if header not in headers] if missing: raise ValueError(f"Excel 缺少供应商名录必要列:{'、'.join(missing)}") return [ {headers[index]: value for index, value in enumerate(row) if index < len(headers)} for row in rows[1:] ] def _ensure_purchase_order_target_warehouse_type(db: Session) -> None: bind = db.get_bind() columns = {column["name"] for column in inspect(bind).get_columns("po_purchase_order")} changed = False if "target_warehouse_type" not in columns: if bind.dialect.name == "mysql": db.execute( text( "ALTER TABLE po_purchase_order " "ADD COLUMN target_warehouse_type VARCHAR(32) NOT NULL DEFAULT 'RAW' COMMENT '目标仓库类型:RAW原材料库/AUX辅料库'" ) ) else: db.execute(text("ALTER TABLE po_purchase_order ADD COLUMN target_warehouse_type VARCHAR(32) NOT NULL DEFAULT 'RAW'")) changed = True if "warning_lead_days" not in columns: if bind.dialect.name == "mysql": db.execute( text( "ALTER TABLE po_purchase_order " "ADD COLUMN warning_lead_days INT NOT NULL DEFAULT 0 COMMENT '预计到货提前预警天数'" ) ) else: db.execute(text("ALTER TABLE po_purchase_order ADD COLUMN warning_lead_days INTEGER NOT NULL DEFAULT 0")) changed = True if changed: db.commit() receipt_columns = {column["name"] for column in inspect(bind).get_columns("po_receipt")} receipt_column_defs = { "logistics_waybill_no": ("VARCHAR(100)", "运单号"), "logistics_freight_amount": ("DECIMAL(18, 2)", "运费"), "logistics_photo_url": ("VARCHAR(255)", "辅助照片"), } receipt_changed = False for column_name, (column_type, comment) in receipt_column_defs.items(): if column_name in receipt_columns: continue if bind.dialect.name == "mysql": db.execute(text(f"ALTER TABLE po_receipt ADD COLUMN {column_name} {column_type} NULL COMMENT '{comment}'")) else: db.execute(text(f"ALTER TABLE po_receipt ADD COLUMN {column_name} {column_type} NULL")) receipt_changed = True if receipt_changed: db.commit() def _normalize_purchase_target_warehouse_type(value: str | None) -> str: target_type = str(value or "RAW").strip().upper() if target_type not in PURCHASE_TARGET_WAREHOUSE_TYPES: raise HTTPException(status_code=400, detail="采购目标仓库类型只能选择原材料库或辅料库") return target_type def _purchase_target_label(value: str | None) -> str: return "辅料库" if _normalize_purchase_target_warehouse_type(value) == "AUX" else "原材料库" def _normalize_sales_order_ids(values: list[int] | None) -> list[int]: normalized_ids: set[int] = set() for value in values or []: sales_order_id = int(value) if sales_order_id <= 0: raise HTTPException(status_code=400, detail="销售订单ID必须大于0") normalized_ids.add(sales_order_id) return sorted(normalized_ids) def _ensure_sales_orders_exist(db: Session, sales_order_ids: list[int]) -> None: for sales_order_id in sales_order_ids: sales_order = db.get(SalesOrder, sales_order_id) if not sales_order: raise HTTPException(status_code=404, detail=f"销售订单不存在: {sales_order_id}") def _replace_purchase_order_sales_links(db: Session, purchase_order_id: int, sales_order_ids: list[int]) -> None: normalized_ids = _normalize_sales_order_ids(sales_order_ids) _ensure_sales_orders_exist(db, normalized_ids) db.execute( sa_delete(PurchaseOrderSalesOrderLink) .where(PurchaseOrderSalesOrderLink.purchase_order_id == purchase_order_id) .execution_options(synchronize_session=False) ) for sales_order_id in normalized_ids: db.add( PurchaseOrderSalesOrderLink( purchase_order_id=purchase_order_id, sales_order_id=sales_order_id, created_at=datetime.now(), updated_at=datetime.now(), ) ) def _coerce_date(value: object | None) -> date | None: if isinstance(value, date): return value if value: try: return date.fromisoformat(str(value)[:10]) except ValueError: return None return None def _enrich_purchase_order_warning(row: dict) -> None: expected_date = _coerce_date(row.get("expected_date")) warning_lead_days = max(int(row.get("warning_lead_days") or 0), 0) warning_start_date = expected_date - timedelta(days=warning_lead_days) if expected_date and warning_lead_days > 0 else None target_warehouse_type = _normalize_purchase_target_warehouse_type(row.get("target_warehouse_type")) total_order_value = to_decimal(row.get("total_order_qty") or 0) if target_warehouse_type == "AUX" else to_decimal(row.get("total_order_weight_kg") or 0) total_received_value = to_decimal(row.get("total_received_qty") or 0) if target_warehouse_type == "AUX" else to_decimal(row.get("total_received_weight_kg") or 0) row["warning_lead_days"] = warning_lead_days row["warning_start_date"] = warning_start_date row["is_warning_active"] = bool( warning_start_date and date.today() >= warning_start_date and total_order_value > 0 and total_received_value < total_order_value ) def _sales_order_ids_from_items(db: Session, items: list) -> list[int]: ids: set[int] = set() for item in items: if not item.source_demand_id: continue row = db.execute( select(SalesOrder.id) .join(SalesOrderItem, SalesOrderItem.sales_order_id == SalesOrder.id) .join(MaterialDemand, MaterialDemand.sales_order_item_id == SalesOrderItem.id) .where(MaterialDemand.id == item.source_demand_id) ).scalar_one_or_none() if row: ids.add(int(row)) return sorted(ids) def _enrich_purchase_order_sources(db: Session, rows: list[dict]) -> list[dict]: purchase_order_ids = [int(row["purchase_order_id"]) for row in rows] if not purchase_order_ids: return rows ensure_purchase_order_sales_link_table(db) link_rows = db.execute( select( PurchaseOrderSalesOrderLink.purchase_order_id.label("purchase_order_id"), SalesOrder.id.label("sales_order_id"), SalesOrder.order_no.label("order_no"), ) .join(SalesOrder, SalesOrder.id == PurchaseOrderSalesOrderLink.sales_order_id) .where(PurchaseOrderSalesOrderLink.purchase_order_id.in_(purchase_order_ids)) .order_by(PurchaseOrderSalesOrderLink.purchase_order_id, SalesOrder.order_no) ).mappings().all() links_by_po: dict[int, list[dict]] = {} for link in link_rows: links_by_po.setdefault(int(link["purchase_order_id"]), []).append(dict(link)) inferred_rows = db.execute( select( PurchaseOrderItem.purchase_order_id.label("purchase_order_id"), SalesOrder.id.label("sales_order_id"), SalesOrder.order_no.label("order_no"), ) .join(MaterialDemand, MaterialDemand.id == PurchaseOrderItem.source_demand_id) .join(SalesOrderItem, SalesOrderItem.id == MaterialDemand.sales_order_item_id) .join(SalesOrder, SalesOrder.id == SalesOrderItem.sales_order_id) .where(PurchaseOrderItem.purchase_order_id.in_(purchase_order_ids)) .order_by(PurchaseOrderItem.purchase_order_id, SalesOrder.order_no) ).mappings().all() inferred_by_po: dict[int, list[dict]] = {} for inferred in inferred_rows: links = inferred_by_po.setdefault(int(inferred["purchase_order_id"]), []) if any(int(link["sales_order_id"]) == int(inferred["sales_order_id"]) for link in links): continue links.append(dict(inferred)) for row in rows: _enrich_purchase_order_warning(row) links = links_by_po.get(int(row["purchase_order_id"]), []) if not links: links = inferred_by_po.get(int(row["purchase_order_id"]), []) if links: row["sales_order_ids"] = [int(link["sales_order_id"]) for link in links] row["source_order_nos"] = "、".join(str(link["order_no"]) for link in links) row["source_type"] = "SALES_ORDER" else: row["sales_order_ids"] = [] return rows def _import_supplier_excel(db: Session, content: bytes) -> dict[str, object]: rows = _read_supplier_excel_rows(content) if not rows: raise ValueError("Excel 中没有可导入的供应商数据") latest_by_name: dict[str, dict[str, str]] = {} skipped = 0 errors: list[str] = [] for index, row in enumerate(rows, start=2): if not any(_clean_supplier_text(value) for value in row.values()): continue supplier_name = _clean_supplier_text(row.get("供应商名称")) if not supplier_name: skipped += 1 errors.append(f"第{index}行缺少供应商名称") continue latest_by_name[supplier_name] = { "supplier_name": supplier_name, "contacts": normalize_contacts( [ {"contact_name": row.get("联系人1"), "contact_phone": row.get("联系电话1")}, {"contact_name": row.get("联系人2"), "contact_phone": row.get("联系电话2")}, {"contact_name": row.get("联系人3"), "contact_phone": row.get("联系电话3")}, ], fallback_name=_clean_supplier_text(row.get("联系人")), fallback_phone=_clean_supplier_text(row.get("联系电话")), ), "address": _clean_supplier_text(row.get("地址")), } if not latest_by_name: raise ValueError("Excel 中没有可导入的有效供应商数据") existing_rows = db.scalars( select(Supplier) .where(Supplier.supplier_name.in_(latest_by_name.keys())) .order_by(Supplier.id.asc()) ).all() existing_by_name: dict[str, Supplier] = {} for supplier in existing_rows: existing_by_name.setdefault(_clean_supplier_text(supplier.supplier_name), supplier) imported = 0 updated = 0 for supplier_name, payload in latest_by_name.items(): contact_name, contact_phone = serialize_contacts(payload["contacts"]) supplier = existing_by_name.get(supplier_name) if supplier: supplier.contact_name = contact_name supplier.contact_phone = contact_phone supplier.address = payload["address"] or None supplier.short_name = supplier.short_name or supplier_name supplier.lead_time_days = 0 supplier.default_tax_rate = to_decimal(0, "0.0001") supplier.status = "ACTIVE" db.add(supplier) updated += 1 continue db.add( Supplier( supplier_code=build_supplier_code(db), supplier_name=supplier_name, short_name=supplier_name, contact_name=contact_name, contact_phone=contact_phone, address=payload["address"] or None, lead_time_days=0, default_tax_rate=to_decimal(0, "0.0001"), status="ACTIVE", remark=None, ) ) db.flush() imported += 1 db.commit() return { "imported": imported, "updated": updated, "skipped": skipped, "errors": errors, "message": f"导入完成:新增 {imported} 条,更新 {updated} 条,跳过 {skipped} 条", } def _supplier_contact_fields(payload: SupplierCreate | SupplierUpdate) -> tuple[str | None, str | None]: try: contacts = normalize_contacts( list(payload.contacts or []), fallback_name=payload.contact_name, fallback_phone=payload.contact_phone, ) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc return serialize_contacts(contacts) def _supplier_read(supplier: Supplier) -> SupplierRead: return SupplierRead.model_validate( { "id": supplier.id, "supplier_code": supplier.supplier_code, "supplier_name": supplier.supplier_name, "short_name": supplier.short_name, "contact_name": supplier.contact_name, "contact_phone": supplier.contact_phone, "contacts": parse_contacts(supplier.contact_name, supplier.contact_phone), "address": supplier.address, "lead_time_days": supplier.lead_time_days, "default_tax_rate": supplier.default_tax_rate, "status": supplier.status, "remark": supplier.remark, } ) @router.get("/suppliers", response_model=list[SupplierRead]) def list_suppliers( limit: int = Query(default=100, ge=1, le=500), db: Session = Depends(get_db), ) -> list[SupplierRead]: rows = db.scalars(select(Supplier).order_by(Supplier.id.desc()).limit(limit)).all() return [_supplier_read(row) for row in rows] @router.post("/suppliers", response_model=SupplierRead) def create_supplier(payload: SupplierCreate, db: Session = Depends(get_db)) -> SupplierRead: contact_name, contact_phone = _supplier_contact_fields(payload) supplier = Supplier( supplier_code=build_supplier_code(db), supplier_name=payload.supplier_name, short_name=payload.short_name, contact_name=contact_name, contact_phone=contact_phone, address=payload.address, lead_time_days=0, default_tax_rate=to_decimal(payload.default_tax_rate, "0.0001"), status=payload.status, remark=payload.remark, ) db.add(supplier) db.commit() db.refresh(supplier) return _supplier_read(supplier) @router.get("/suppliers/export") def export_suppliers(db: Session = Depends(get_db)) -> Response: rows = db.scalars(select(Supplier).order_by(Supplier.supplier_name.asc(), Supplier.id.asc())).all() workbook = _build_supplier_workbook(rows) return _supplier_excel_response(workbook, f"供应商名录_{datetime.now().strftime('%Y%m%d%H%M%S')}.xlsx") @router.post("/suppliers/import") async def import_suppliers( file: UploadFile = File(...), db: Session = Depends(get_db), ) -> dict[str, object]: content = await file.read() try: result = _import_supplier_excel(db, content) except ValueError as exc: db.rollback() raise HTTPException(status_code=400, detail=str(exc)) from exc return result @router.put("/suppliers/{supplier_id}", response_model=SupplierRead) def update_supplier(supplier_id: int, payload: SupplierUpdate, db: Session = Depends(get_db)) -> SupplierRead: supplier = db.get(Supplier, supplier_id) if not supplier: raise HTTPException(status_code=404, detail="供应商不存在") contact_name, contact_phone = _supplier_contact_fields(payload) supplier.supplier_name = payload.supplier_name supplier.short_name = payload.short_name supplier.contact_name = contact_name supplier.contact_phone = contact_phone supplier.address = payload.address supplier.lead_time_days = 0 supplier.default_tax_rate = to_decimal(payload.default_tax_rate, "0.0001") supplier.status = payload.status supplier.remark = payload.remark db.add(supplier) db.commit() db.refresh(supplier) return _supplier_read(supplier) @router.delete("/suppliers/{supplier_id}") def delete_supplier(supplier_id: int, db: Session = Depends(get_db)) -> dict[str, str]: supplier = db.get(Supplier, supplier_id) if not supplier: raise HTTPException(status_code=404, detail="供应商不存在") references = [ ("采购订单", db.scalar(select(func.count(PurchaseOrder.id)).where(PurchaseOrder.supplier_id == supplier_id)) or 0), ] used = [f"{label}({count})" for label, count in references if count > 0] if used: raise HTTPException(status_code=400, detail=f"该供应商已被业务数据引用,不能删除:{'、'.join(used)}") supplier_name = supplier.supplier_name db.execute(update(Material).where(Material.default_supplier_id == supplier_id).values(default_supplier_id=None)) db.delete(supplier) db.commit() return {"message": f"供应商 {supplier_name} 已删除"} @router.get("/orders", response_model=list[PurchaseOrderRead]) def list_purchase_orders( limit: int = Query(default=100, ge=1, le=500), db: Session = Depends(get_db), ) -> list[PurchaseOrderRead]: _ensure_purchase_order_target_warehouse_type(db) rows = [dict(row) for row in db.execute(get_purchase_order_query(limit=limit)).mappings().all()] rows = _enrich_purchase_order_sources(db, rows) return [PurchaseOrderRead.model_validate(dict(row)) for row in rows] @router.get("/order-items", response_model=list[PurchaseOrderItemRead]) def list_purchase_order_items( purchase_order_id: int | None = Query(default=None), limit: int = Query(default=200, ge=1, le=500), db: Session = Depends(get_db), ) -> list[PurchaseOrderItemRead]: _ensure_purchase_order_target_warehouse_type(db) rows = db.execute(get_purchase_order_items_query(limit=limit, purchase_order_id=purchase_order_id)).mappings().all() return [PurchaseOrderItemRead.model_validate(dict(row)) for row in rows] def _read_purchase_order(db: Session, purchase_order_id: int) -> PurchaseOrderRead: _ensure_purchase_order_target_warehouse_type(db) row = db.execute(get_purchase_order_query(limit=1).where(PurchaseOrder.id == purchase_order_id)).mappings().first() if not row: raise HTTPException(status_code=500, detail="采购订单读取失败") rows = _enrich_purchase_order_sources(db, [dict(row)]) return PurchaseOrderRead.model_validate(rows[0]) def _read_purchase_order_after_archive_generation( db: Session, purchase_order_id: int, *, created_by: int | None, ) -> PurchaseOrderRead: try: archive_kwargs = {"created_by": created_by} if created_by is not None else {} archive_result = generate_document_archive( db, DOCUMENT_TYPE_PURCHASE_ORDER, int(purchase_order_id), **archive_kwargs, ) archive_fields = { "archive_status": archive_result.archive_status, "archive_version": archive_result.archive_version, "archive_error_message": archive_result.archive_error_message, } except Exception as exc: # The order is already committed; this only clears any failed archive transaction. try: db.rollback() except Exception: pass archive_fields = { "archive_status": "归档失败", "archive_version": None, "archive_error_message": f"归档失败:{exc}", } return _read_purchase_order(db, purchase_order_id).model_copy(update=archive_fields) def _get_purchase_order_display_status(db: Session, purchase_order_id: int) -> str: row = db.execute(get_purchase_order_query(limit=1).where(PurchaseOrder.id == purchase_order_id)).mappings().first() if not row: return "" return str(row.get("display_status") or row.get("status") or "").upper() def _assert_purchase_order_editable(db: Session, order: PurchaseOrder) -> list[PurchaseOrderItem]: if order.status != "OPEN": raise HTTPException(status_code=400, detail="只有未锁单的采购订单可以修改或删除") receipt_count = db.scalar(select(func.count(PurchaseReceipt.id)).where(PurchaseReceipt.purchase_order_id == order.id)) or 0 if receipt_count > 0: raise HTTPException(status_code=400, detail="采购订单已有到货入库记录,不能修改或删除") items = db.scalars(select(PurchaseOrderItem).where(PurchaseOrderItem.purchase_order_id == order.id)).all() if any(to_decimal(item.received_weight_kg) > 0 or to_decimal(item.received_qty) > 0 for item in items): raise HTTPException(status_code=400, detail="采购订单已有收货重量,不能修改或删除") return items def _reset_unreferenced_demands(db: Session, demand_ids: set[int]) -> None: for demand_id in demand_ids: still_referenced = db.scalar( select(func.count(PurchaseOrderItem.id)).where(PurchaseOrderItem.source_demand_id == demand_id) ) or 0 if still_referenced: continue demand = db.get(MaterialDemand, demand_id) if demand and demand.status == "ORDERED": demand.status = "NEW" db.add(demand) def _build_purchase_order_items( db: Session, *, order: PurchaseOrder, payload_items: list, target_warehouse_type: str, ) -> tuple[list[PurchaseOrderItem], set[int], object]: if not payload_items: raise HTTPException(status_code=400, detail="采购订单至少需要一条明细") rows: list[PurchaseOrderItem] = [] linked_demand_ids: set[int] = set() total_amount = to_decimal(0, "0.01") for index, item_payload in enumerate(payload_items, start=1): material_item = db.get(Item, item_payload.material_item_id) item_label = "辅料" if target_warehouse_type == "AUX" else "原材料" if not material_item: raise HTTPException(status_code=404, detail=f"{item_label}不存在: {item_payload.material_item_id}") material_ext = db.scalar(select(Material).where(Material.item_id == item_payload.material_item_id)) if not material_ext: raise HTTPException(status_code=404, detail=f"{item_label}缺少采购参数: {item_payload.material_item_id}") if target_warehouse_type == "AUX": order_qty = to_decimal(item_payload.order_qty) order_weight = to_decimal(0) else: order_qty = to_decimal(0) order_weight = to_decimal(item_payload.order_weight_kg) if target_warehouse_type == "AUX" and item_payload.source_demand_id: raise HTTPException(status_code=400, detail="辅料库采购只能做备料采购,不能关联销售订单采购需求") if item_payload.source_demand_id: demand = db.get(MaterialDemand, item_payload.source_demand_id) if not demand: raise HTTPException(status_code=404, detail=f"MRP 需求不存在: {item_payload.source_demand_id}") if demand.material_item_id != item_payload.material_item_id: raise HTTPException(status_code=400, detail=f"MRP 需求 {demand.demand_no} 与采购原材料不一致") if order_weight <= 0: order_weight = to_decimal(demand.suggested_purchase_weight_kg) demand.status = "ORDERED" db.add(demand) linked_demand_ids.add(int(demand.id)) if target_warehouse_type == "AUX": if order_qty <= 0: raise HTTPException(status_code=400, detail=f"辅料 {material_item.item_name} 的采购数量必须大于0") line_amount = to_decimal(order_qty, "0.000001") * to_decimal(item_payload.unit_price, "0.0001") else: if order_weight <= 0: raise HTTPException(status_code=400, detail=f"原材料 {material_item.item_name} 的采购重量必须大于0") line_amount = to_decimal(order_weight, "0.000001") * to_decimal(item_payload.unit_price, "0.0001") line_amount = to_decimal(line_amount, "0.01") rows.append( PurchaseOrderItem( purchase_order_id=order.id, line_no=index, material_item_id=item_payload.material_item_id, source_demand_id=item_payload.source_demand_id, order_qty=order_qty, order_weight_kg=order_weight, received_qty=to_decimal(0), received_weight_kg=to_decimal(0), unit_price=to_decimal(item_payload.unit_price, "0.0001"), line_amount=line_amount, status="OPEN", remark=item_payload.remark, ) ) total_amount += line_amount return rows, linked_demand_ids, total_amount @router.post("/orders", response_model=PurchaseOrderRead) def create_purchase_order( payload: PurchaseOrderCreate, db: Session = Depends(get_db), context: AuthContext = Depends(require_authenticated_user), ) -> PurchaseOrderRead: _ensure_purchase_order_target_warehouse_type(db) ensure_purchase_order_sales_link_table(db) explicit_sales_order_ids = _normalize_sales_order_ids(payload.sales_order_ids) target_warehouse_type = _normalize_purchase_target_warehouse_type(payload.target_warehouse_type) if target_warehouse_type == "AUX" and explicit_sales_order_ids: raise HTTPException(status_code=400, detail="辅料库采购只能做备料采购,不能显式关联销售订单") supplier = db.get(Supplier, payload.supplier_id) if not supplier: raise HTTPException(status_code=404, detail="供应商不存在") _ensure_sales_orders_exist(db, explicit_sales_order_ids) ensure_employee_has_permission( db, payload.purchaser_employee_id, "MENU_PURCHASE_ORDER", "采购员", required=False, ) order = PurchaseOrder( po_no=build_purchase_order_no(db, supplier), supplier_id=payload.supplier_id, order_date=payload.order_date or date.today(), expected_date=payload.expected_date, warning_lead_days=int(payload.warning_lead_days or 0), purchaser_employee_id=payload.purchaser_employee_id, target_warehouse_type=target_warehouse_type, tax_rate=to_decimal(payload.tax_rate, "0.0001"), total_amount=to_decimal(0, "0.01"), status="OPEN", remark=payload.remark, ) db.add(order) db.flush() item_rows, _linked_demand_ids, total_amount = _build_purchase_order_items( db, order=order, payload_items=payload.items, target_warehouse_type=target_warehouse_type, ) for item_row in item_rows: db.add(item_row) order.total_amount = total_amount linked_sales_order_ids = sorted({*explicit_sales_order_ids, *_sales_order_ids_from_items(db, payload.items)}) _replace_purchase_order_sales_links(db, order.id, linked_sales_order_ids) db.add(order) db.commit() return _read_purchase_order_after_archive_generation(db, int(order.id), created_by=_auth_context_user_id(context)) @router.put("/orders/{purchase_order_id}", response_model=PurchaseOrderRead) def update_purchase_order( purchase_order_id: int, payload: PurchaseOrderUpdate, db: Session = Depends(get_db), context: AuthContext = Depends(require_authenticated_user), ) -> PurchaseOrderRead: _ensure_purchase_order_target_warehouse_type(db) ensure_purchase_order_sales_link_table(db) explicit_sales_order_ids = _normalize_sales_order_ids(payload.sales_order_ids) target_warehouse_type = _normalize_purchase_target_warehouse_type(payload.target_warehouse_type) if target_warehouse_type == "AUX" and explicit_sales_order_ids: raise HTTPException(status_code=400, detail="辅料库采购只能做备料采购,不能显式关联销售订单") order = db.get(PurchaseOrder, purchase_order_id) if not order: raise HTTPException(status_code=404, detail="采购订单不存在") supplier = db.get(Supplier, payload.supplier_id) if not supplier: raise HTTPException(status_code=404, detail="供应商不存在") old_items = _assert_purchase_order_editable(db, order) _ensure_sales_orders_exist(db, explicit_sales_order_ids) ensure_employee_has_permission( db, payload.purchaser_employee_id, "MENU_PURCHASE_ORDER", "采购员", required=False, ) old_demand_ids = {int(item.source_demand_id) for item in old_items if item.source_demand_id} order.supplier_id = payload.supplier_id order.order_date = payload.order_date or order.order_date or date.today() order.expected_date = payload.expected_date order.warning_lead_days = int(payload.warning_lead_days or 0) order.purchaser_employee_id = payload.purchaser_employee_id order.target_warehouse_type = target_warehouse_type order.tax_rate = to_decimal(payload.tax_rate, "0.0001") order.remark = payload.remark db.execute( sa_delete(PurchaseOrderItem) .where(PurchaseOrderItem.purchase_order_id == order.id) .execution_options(synchronize_session=False) ) db.flush() for item in old_items: if item in db: db.expunge(item) item_rows, _linked_demand_ids, total_amount = _build_purchase_order_items( db, order=order, payload_items=payload.items, target_warehouse_type=target_warehouse_type, ) for item_row in item_rows: db.add(item_row) order.total_amount = total_amount order.status = "OPEN" linked_sales_order_ids = sorted({*explicit_sales_order_ids, *_sales_order_ids_from_items(db, payload.items)}) _replace_purchase_order_sales_links(db, order.id, linked_sales_order_ids) db.add(order) db.flush() _reset_unreferenced_demands(db, old_demand_ids) db.commit() return _read_purchase_order_after_archive_generation(db, int(order.id), created_by=_auth_context_user_id(context)) @router.delete("/orders/{purchase_order_id}") def delete_purchase_order( purchase_order_id: int, db: Session = Depends(get_db), ) -> dict[str, str]: _ensure_purchase_order_target_warehouse_type(db) ensure_purchase_order_sales_link_table(db) order = db.get(PurchaseOrder, purchase_order_id) if not order: raise HTTPException(status_code=404, detail="采购订单不存在") old_items = _assert_purchase_order_editable(db, order) old_demand_ids = {int(item.source_demand_id) for item in old_items if item.source_demand_id} po_no = order.po_no try: db.execute( sa_delete(PurchaseOrderItem) .where(PurchaseOrderItem.purchase_order_id == order.id) .execution_options(synchronize_session=False) ) db.flush() for item in old_items: if item in db: db.expunge(item) db.execute( sa_delete(PurchaseOrderSalesOrderLink) .where(PurchaseOrderSalesOrderLink.purchase_order_id == order.id) .execution_options(synchronize_session=False) ) db.delete(order) db.flush() _reset_unreferenced_demands(db, old_demand_ids) db.commit() except IntegrityError as exc: db.rollback() raise HTTPException(status_code=400, detail="采购订单已被后续业务引用,不能删除") from exc return {"message": f"采购订单 {po_no} 已删除"} @router.post("/orders/{purchase_order_id}/lock", response_model=PurchaseOrderRead) def lock_purchase_order( purchase_order_id: int, db: Session = Depends(get_db), ) -> PurchaseOrderRead: _ensure_purchase_order_target_warehouse_type(db) order = db.get(PurchaseOrder, purchase_order_id) if not order: raise HTTPException(status_code=404, detail="采购订单不存在") if order.status == "LOCKED": row = db.execute(get_purchase_order_query(limit=1).where(PurchaseOrder.id == order.id)).mappings().first() rows = _enrich_purchase_order_sources(db, [dict(row)]) return PurchaseOrderRead.model_validate(rows[0]) if order.status != "OPEN": raise HTTPException(status_code=400, detail="只有未开始的采购订单可以锁单") item_count = db.scalar(select(func.count(PurchaseOrderItem.id)).where(PurchaseOrderItem.purchase_order_id == order.id)) or 0 if item_count <= 0: raise HTTPException(status_code=400, detail="采购订单没有明细,不能锁单") order.status = "LOCKED" db.add(order) db.commit() row = db.execute(get_purchase_order_query(limit=1).where(PurchaseOrder.id == order.id)).mappings().first() if not row: raise HTTPException(status_code=500, detail="采购订单锁单后读取失败") rows = _enrich_purchase_order_sources(db, [dict(row)]) return PurchaseOrderRead.model_validate(rows[0]) @router.get("/receipts", response_model=list[PurchaseReceiptRead]) def list_purchase_receipts( limit: int = Query(default=100, ge=1, le=500), db: Session = Depends(get_db), ) -> list[PurchaseReceiptRead]: _ensure_purchase_order_target_warehouse_type(db) rows = db.execute(get_purchase_receipt_query(limit=limit)).mappings().all() return [PurchaseReceiptRead.model_validate(dict(row)) for row in rows] @router.get("/receipt-items", response_model=list[PurchaseReceiptItemRead]) def list_purchase_receipt_items( receipt_id: int | None = Query(default=None), limit: int = Query(default=200, ge=1, le=500), db: Session = Depends(get_db), ) -> list[PurchaseReceiptItemRead]: _ensure_purchase_order_target_warehouse_type(db) rows = db.execute(get_purchase_receipt_items_query(limit=limit, receipt_id=receipt_id)).mappings().all() return [PurchaseReceiptItemRead.model_validate(dict(row)) for row in rows] def _read_purchase_receipt(db: Session, purchase_receipt_id: int) -> PurchaseReceiptRead: _ensure_purchase_order_target_warehouse_type(db) row = db.execute(get_purchase_receipt_query(limit=1).where(PurchaseReceipt.id == purchase_receipt_id)).mappings().first() if not row: raise HTTPException(status_code=500, detail="到货入库单读取失败") return PurchaseReceiptRead.model_validate(dict(row)) def _read_purchase_receipt_after_archive_generation( db: Session, purchase_receipt_id: int, *, created_by: int | None, ) -> PurchaseReceiptRead: try: archive_kwargs = {"created_by": created_by} if created_by is not None else {} archive_result = generate_document_archive( db, DOCUMENT_TYPE_PURCHASE_RECEIPT, int(purchase_receipt_id), **archive_kwargs, ) archive_fields = { "archive_status": archive_result.archive_status, "archive_version": archive_result.archive_version, "archive_error_message": archive_result.archive_error_message, } except Exception as exc: try: db.rollback() except Exception: pass archive_fields = { "archive_status": "归档失败", "archive_version": None, "archive_error_message": f"归档失败:{exc}", } return _read_purchase_receipt(db, purchase_receipt_id).model_copy(update=archive_fields) @router.post("/receipts", response_model=PurchaseReceiptRead) def create_purchase_receipt( payload: PurchaseReceiptCreate, context: AuthContext = Depends(require_authenticated_user), db: Session = Depends(get_db), ) -> PurchaseReceiptRead: _ensure_purchase_order_target_warehouse_type(db) purchase_order = db.get(PurchaseOrder, payload.purchase_order_id) if not purchase_order: raise HTTPException(status_code=404, detail="采购订单不存在") display_status = _get_purchase_order_display_status(db, purchase_order.id) if purchase_order.status not in RECEIVABLE_PURCHASE_ORDER_STATUSES and display_status not in RECEIVABLE_PURCHASE_ORDER_STATUSES: raise HTTPException(status_code=400, detail="只有已锁单且未完成收货的采购订单可以到货入库") warehouse = db.get(Warehouse, payload.warehouse_id) if not warehouse: raise HTTPException(status_code=404, detail="仓库不存在") target_warehouse_type = _normalize_purchase_target_warehouse_type(purchase_order.target_warehouse_type) warehouse_type = str(warehouse.warehouse_type or "").upper() if warehouse_type != target_warehouse_type: raise HTTPException( status_code=400, detail=f"该采购订单目标到货仓库为{_purchase_target_label(target_warehouse_type)},不能入库到{warehouse.warehouse_name}", ) logistics_waybill_no, logistics_freight_amount, logistics_photo_url = normalize_logistics_fields( payload.waybill_no, payload.freight_amount, required=target_warehouse_type == "RAW", order_photo_url=payload.order_photo_url, freight_required=None if target_warehouse_type == "RAW" else False, ) ensure_warehouses_unlocked(db, [payload.warehouse_id], "到货入库") if not payload.items: raise HTTPException(status_code=400, detail="到货入库至少需要一条明细") ensure_employee_has_permission( db, payload.receiver_employee_id, "MENU_PURCHASE_RECEIPT", "接收人", required=False, ) receipt = PurchaseReceipt( receipt_no=build_purchase_receipt_no(db, purchase_order), purchase_order_id=payload.purchase_order_id, warehouse_id=payload.warehouse_id, receipt_date=payload.receipt_date or datetime.now(), receiver_employee_id=payload.receiver_employee_id, supplier_delivery_no=payload.supplier_delivery_no, logistics_waybill_no=logistics_waybill_no, logistics_freight_amount=logistics_freight_amount, logistics_photo_url=logistics_photo_url, status="PENDING_QC", remark=payload.remark, ) try: db.add(receipt) db.flush() selected_po_item_ids: set[int] = set() for index, item_payload in enumerate(payload.items, start=1): po_item = db.get(PurchaseOrderItem, item_payload.purchase_order_item_id) if not po_item or po_item.purchase_order_id != payload.purchase_order_id: raise HTTPException(status_code=400, detail="采购订单明细不存在或不属于当前采购订单") if item_payload.purchase_order_item_id in selected_po_item_ids: raise HTTPException(status_code=400, detail=f"采购行 {po_item.line_no} 在同一张到货单中重复提交") selected_po_item_ids.add(item_payload.purchase_order_item_id) if po_item.material_item_id != item_payload.material_item_id: raise HTTPException(status_code=400, detail=f"采购行 {po_item.line_no} 的原材料与当前选择不一致") material_item = db.get(Item, item_payload.material_item_id) item_label = "辅料" if target_warehouse_type == "AUX" else "原材料" if not material_item: raise HTTPException(status_code=404, detail=f"{item_label}不存在: {item_payload.material_item_id}") if target_warehouse_type == "AUX": received_qty = to_decimal(item_payload.received_qty) received_weight = to_decimal(0) if received_qty <= 0: raise HTTPException(status_code=400, detail=f"采购行 {po_item.line_no} 的收货数量必须大于0") else: received_qty = to_decimal(0) received_weight = to_decimal(item_payload.received_weight_kg) if received_weight <= 0: raise HTTPException(status_code=400, detail=f"采购行 {po_item.line_no} 的收货重量必须大于0") entered_unit_cost = to_decimal(item_payload.unit_cost, "0.0001") if entered_unit_cost < 0: raise HTTPException(status_code=400, detail=f"采购行 {po_item.line_no} 的入库单价不能小于0") if entered_unit_cost <= 0: entered_unit_cost = to_decimal(po_item.unit_price, "0.0001") location = get_location_or_default(db, payload.warehouse_id, item_payload.location_id) if location and location.warehouse_id != payload.warehouse_id: raise HTTPException(status_code=400, detail=f"采购行 {po_item.line_no} 选择的库位不属于当前仓库") final_unit_cost = entered_unit_cost lot_no = build_inventory_lot_no(db, item_payload.material_item_id) receipt_item = PurchaseReceiptItem( receipt_id=receipt.id, line_no=index, purchase_order_item_id=item_payload.purchase_order_item_id, material_item_id=item_payload.material_item_id, location_id=location.id if location else None, lot_no=lot_no, merge_to_lot_id=None, received_qty=received_qty, received_weight_kg=received_weight, accepted_qty=to_decimal(0), accepted_weight_kg=to_decimal(0), unit_cost=final_unit_cost, status="PENDING_QC", remark=item_payload.remark, ) db.add(receipt_item) db.flush() lot = StockLot( lot_no=lot_no, parent_lot_id=None, lot_role="INBOUND_RAW", material_sub_batch_no=None, item_id=item_payload.material_item_id, warehouse_id=payload.warehouse_id, location_id=location.id if location else None, source_doc_type="PURCHASE_RECEIPT", source_doc_id=receipt.id, source_line_id=receipt_item.id, source_material_lot_id=None, source_material_sub_batch_no=None, source_material_summary=None, inbound_qty=received_qty, inbound_weight_kg=received_weight, remaining_qty=received_qty, remaining_weight_kg=received_weight, locked_qty=received_qty, locked_weight_kg=received_weight, unit_cost=final_unit_cost, production_date=None, expire_date=None, quality_status="PENDING_QC", logistics_waybill_no=logistics_waybill_no, logistics_freight_amount=logistics_freight_amount, logistics_photo_url=logistics_photo_url, status="LOCKED", remark="到货待检", ) db.add(lot) db.flush() upsert_stock_balance( db, item_id=item_payload.material_item_id, warehouse_id=payload.warehouse_id, location_id=location.id if location else None, qty_delta=received_qty, weight_delta=received_weight, available_qty_delta=to_decimal(0), available_weight_delta=to_decimal(0), unit_cost=final_unit_cost, ) create_inventory_txn( db, txn_type="PURCHASE_IN", item_id=item_payload.material_item_id, warehouse_id=payload.warehouse_id, location_id=location.id if location else None, lot_id=lot.id, qty_change=received_qty, weight_change=received_weight, unit_cost=final_unit_cost, source_doc_type="PURCHASE_RECEIPT", source_doc_id=receipt.id, source_line_id=receipt_item.id, biz_time=receipt.receipt_date, operator_user_id=context.user.id, logistics_waybill_no=logistics_waybill_no, logistics_freight_amount=logistics_freight_amount, logistics_photo_url=logistics_photo_url, remark="采购到货入库", amount_basis="QTY" if target_warehouse_type == "AUX" else "WEIGHT", ) po_item.received_qty = to_decimal(po_item.received_qty) + received_qty po_item.received_weight_kg = to_decimal(po_item.received_weight_kg) + received_weight sync_purchase_order_item_status(po_item) db.add(po_item) if po_item.source_demand_id: demand = db.get(MaterialDemand, po_item.source_demand_id) if demand: demand.status = "ORDERED" db.add(demand) sync_purchase_order_status(db, payload.purchase_order_id) db.commit() except IntegrityError as exc: db.rollback() error_text = str(exc.orig or exc) if "lot_no" in error_text or "Duplicate entry" in error_text: raise HTTPException(status_code=400, detail="批次号重复,请更换批次号后重试") from exc raise HTTPException(status_code=400, detail="到货入库保存失败,请检查采购行、批次号、接收人与收货重量后重试") from exc return _read_purchase_receipt_after_archive_generation(db, receipt.id, created_by=_auth_context_user_id(context))