ForgeFlow-ERP/backend/app/services/material_lifecycle_dag.py

457 lines
19 KiB
Python

from __future__ import annotations
from decimal import Decimal
from typing import Any
from sqlalchemy import false, select
from sqlalchemy.orm import Session, aliased
from app.models.master_data import Item
from app.models.operations import Delivery, DeliveryItem, OperationReport, ProductionBatchLedger, ProductionBatchLedgerTxn, StockLot
from app.models.org import Employee
from app.models.sales import Customer, SalesOrder
from app.schemas.dashboard import DashboardMaterialLotRead, LifecycleDetailItem, LifecycleEdge, LifecycleNode
from app.services.document_archives import production_document_batch_no_from_text
from app.services.system_config import get_smart_operation_report_enabled
def _num(value: Any) -> float:
if value is None:
return 0.0
if isinstance(value, Decimal):
return float(value)
try:
return float(value)
except (TypeError, ValueError):
return 0.0
def _fmt_qty(value: Any, digits: int = 2) -> str:
numeric = _num(value)
if abs(numeric - round(numeric)) < 0.000001:
return f"{int(round(numeric))}"
return f"{numeric:.{digits}f}".rstrip("0").rstrip(".")
def _fmt_weight(value: Any) -> str:
return f"{_fmt_qty(value, 3)} kg"
def _fmt_datetime(value: Any) -> str:
if not value:
return "-"
return str(value).replace("T", " ")[:19]
def _detail(label: str, value: Any) -> LifecycleDetailItem:
text = "-" if value is None or value == "" else str(value)
return LifecycleDetailItem(label=label, value=text)
def _node(
*,
node_id: str,
node_type: str,
title: str,
stage: str,
level: int,
subtitle: str | None = None,
status: str | None = None,
tone: str = "normal",
summary: str | None = None,
metrics: list[tuple[str, Any]] | None = None,
details: list[tuple[str, Any]] | None = None,
raw: dict[str, Any] | None = None,
) -> LifecycleNode:
return LifecycleNode(
id=node_id,
type=node_type,
title=title,
subtitle=subtitle,
status=status,
stage=stage,
level=level,
tone=tone,
summary=summary,
metrics=[_detail(label, value) for label, value in (metrics or [])],
details=[_detail(label, value) for label, value in (details or [])],
raw=raw or {},
)
def _edge(source: str, target: str, label: str | None = None) -> LifecycleEdge:
return LifecycleEdge(id=f"{source}->{target}", source=source, target=target, label=label)
def _append_node(nodes: list[LifecycleNode], node_ids: set[str], node: LifecycleNode) -> None:
if node.id not in node_ids:
node_ids.add(node.id)
nodes.append(node)
def _append_edge(edges: list[LifecycleEdge], edge_ids: set[str], node_ids: set[str], edge: LifecycleEdge) -> None:
if edge.source in node_ids and edge.target in node_ids and edge.id not in edge_ids:
edge_ids.add(edge.id)
edges.append(edge)
def _settlement_group_key(txn: ProductionBatchLedgerTxn) -> tuple[str, str]:
batch_no = production_document_batch_no_from_text(txn.remark)
if batch_no:
return f"batch:{batch_no}", batch_no
return f"txn:{txn.id}", f"{txn.txn_type}-{txn.id}"
def _txn_node_type(txn_type: str) -> str:
return {
"成品入库": "finished_inbound",
"生产余料入库": "surplus_inbound",
"生产废料入库": "scrap_inbound",
"重新开工": "ledger_reopen",
}.get(txn_type, "ledger_writeoff")
def build_material_lot_lifecycle_graph(
db: Session,
lot: DashboardMaterialLotRead,
) -> tuple[list[LifecycleNode], list[LifecycleEdge]]:
nodes: list[LifecycleNode] = []
edges: list[LifecycleEdge] = []
node_ids: set[str] = set()
edge_ids: set[str] = set()
raw_node_id = f"raw-lot-{lot.lot_id}"
_append_node(
nodes,
node_ids,
_node(
node_id=raw_node_id,
node_type="raw_material_lot",
title=lot.lot_no,
subtitle=f"{lot.material_code} · {lot.material_name}",
status=lot.status,
stage="原料入库",
level=0,
summary="原材料批次入库后,按生产台账追踪生产出库、库外余额、入库结算和销售去向。",
metrics=[
("入库重量", _fmt_weight(lot.inbound_weight_kg)),
("剩余重量", _fmt_weight(lot.remaining_weight_kg)),
],
details=[
("原材料编码", lot.material_code),
("原材料名称", lot.material_name),
("仓库", lot.warehouse_name),
],
raw=lot.model_dump(),
),
)
ledger_rows = db.execute(
select(
ProductionBatchLedger,
Item.item_code.label("product_code"),
Item.item_name.label("product_name"),
)
.join(Item, Item.id == ProductionBatchLedger.product_item_id)
.where(ProductionBatchLedger.material_lot_id == lot.lot_id)
.order_by(ProductionBatchLedger.id)
).all()
ledgers = [row.ProductionBatchLedger for row in ledger_rows]
ledger_ids = [int(ledger.id) for ledger in ledgers]
txns_by_ledger_id: dict[int, list[ProductionBatchLedgerTxn]] = {ledger_id: [] for ledger_id in ledger_ids}
if ledger_ids:
txns = db.scalars(
select(ProductionBatchLedgerTxn)
.where(ProductionBatchLedgerTxn.production_ledger_id.in_(ledger_ids))
.order_by(ProductionBatchLedgerTxn.biz_time, ProductionBatchLedgerTxn.id)
).all()
for txn in txns:
txns_by_ledger_id.setdefault(int(txn.production_ledger_id), []).append(txn)
finished_node_id_by_lot_id: dict[int, str] = {}
settlement_txn_types = {"成品入库", "生产余料入库", "生产废料入库", "结单核销", "该批材料结单", "重新开工"}
for row in ledger_rows:
ledger = row.ProductionBatchLedger
ledger_id = int(ledger.id)
ledger_node_id = f"ledger-{ledger_id}"
_append_node(
nodes,
node_ids,
_node(
node_id=ledger_node_id,
node_type="production_ledger",
title=f"{ledger.material_lot_no} · {row.product_name}",
subtitle=row.product_code,
status=ledger.status,
stage="生产台账",
level=1,
summary="一行生产台账代表同一材料库存批次号与产品组合,是材料批次闭环的主链路。",
metrics=[
("累计领料", _fmt_weight(ledger.total_issued_weight_kg)),
("库外余额", _fmt_weight(ledger.outside_weight_kg)),
("成品入库", _fmt_qty(ledger.finished_inbound_qty)),
("余料入库", _fmt_weight(ledger.surplus_inbound_weight_kg)),
("废料入库", _fmt_weight(ledger.scrap_inbound_weight_kg)),
],
details=[
("产品编码", row.product_code),
("产品名称", row.product_name),
("首次领料", _fmt_datetime(ledger.first_issue_time)),
("最近领料", _fmt_datetime(ledger.last_issue_time)),
("结单次数", ledger.lock_count),
("备注", ledger.remark),
],
raw={
"production_ledger_id": ledger.id,
"material_lot_id": ledger.material_lot_id,
"product_item_id": ledger.product_item_id,
},
),
)
_append_edge(edges, edge_ids, node_ids, _edge(raw_node_id, ledger_node_id, "进入生产台账"))
issue_txns = [txn for txn in txns_by_ledger_id.get(ledger_id, []) if txn.txn_type == "生产出库"]
issue_node_ids: list[str] = []
for txn in issue_txns:
issue_node_id = f"issue-txn-{txn.id}"
issue_node_ids.append(issue_node_id)
_append_node(
nodes,
node_ids,
_node(
node_id=issue_node_id,
node_type="production_issue_txn",
title=f"生产出库 {txn.id}",
subtitle=ledger.material_lot_no,
status="POSTED",
tone="success",
stage="生产出库",
level=2,
summary="生产出库事务累加到生产台账,并增加该批材料的库外余额。",
metrics=[
("出库重量", _fmt_weight(txn.weight_delta_kg)),
("出库后库外余额", _fmt_weight(txn.outside_weight_after_kg)),
],
details=[
("业务时间", _fmt_datetime(txn.biz_time)),
("来源单据", txn.source_doc_type),
("来源ID", txn.source_doc_id),
("备注", txn.remark),
],
raw={"txn_id": txn.id, "production_ledger_id": txn.production_ledger_id},
),
)
_append_edge(edges, edge_ids, node_ids, _edge(ledger_node_id, issue_node_id, "生产出库"))
outside_node_id = f"outside-balance-{ledger_id}"
_append_node(
nodes,
node_ids,
_node(
node_id=outside_node_id,
node_type="ledger_outside_balance",
title="库外余额",
subtitle=ledger.material_lot_no,
status=ledger.status,
stage="库外闭环",
level=3,
summary="库外余额用于衡量该材料批次尚未通过成品、余料、废料入库或结单核销闭环的重量。",
metrics=[
("库外余额", _fmt_weight(ledger.outside_weight_kg)),
("累计领料", _fmt_weight(ledger.total_issued_weight_kg)),
("结单核销", _fmt_weight(ledger.writeoff_weight_kg)),
],
raw={"production_ledger_id": ledger.id},
),
)
if issue_node_ids:
for issue_node_id in issue_node_ids:
_append_edge(edges, edge_ids, node_ids, _edge(issue_node_id, outside_node_id, "形成库外余额"))
else:
_append_edge(edges, edge_ids, node_ids, _edge(ledger_node_id, outside_node_id, "库外余额"))
settlement_groups: dict[tuple[str, str], list[ProductionBatchLedgerTxn]] = {}
for txn in txns_by_ledger_id.get(ledger_id, []):
if txn.txn_type in settlement_txn_types:
settlement_groups.setdefault(_settlement_group_key(txn), []).append(txn)
for group_index, ((group_key, group_title), group_txns) in enumerate(settlement_groups.items(), start=1):
settlement_node_id = f"settlement-{ledger_id}-{group_key}"
total_weight_delta = sum(_num(txn.weight_delta_kg) for txn in group_txns)
_append_node(
nodes,
node_ids,
_node(
node_id=settlement_node_id,
node_type="production_ledger_settlement",
title=group_title,
subtitle=f"{len(group_txns)} 笔台账事务",
status="POSTED",
tone="success",
stage="入库结算",
level=4,
summary="按生产单据批次号聚合成品、余料、废料入库及结单核销事务;没有批次号的事务独立成组。",
metrics=[
("事务数", len(group_txns)),
("重量变化", _fmt_weight(total_weight_delta)),
],
details=[
("分组", group_title),
("序号", group_index),
],
raw={"production_ledger_id": ledger.id, "group_key": group_key},
),
)
_append_edge(edges, edge_ids, node_ids, _edge(outside_node_id, settlement_node_id, "入库结算"))
txns_by_type: dict[str, list[ProductionBatchLedgerTxn]] = {}
for txn in group_txns:
txns_by_type.setdefault(str(txn.txn_type), []).append(txn)
for txn_type, type_txns in txns_by_type.items():
child_node_id = f"{_txn_node_type(txn_type)}-{ledger_id}-{group_key}"
child_weight_delta = sum(_num(txn.weight_delta_kg) for txn in type_txns)
child_qty_delta = sum(_num(txn.qty_delta) for txn in type_txns)
_append_node(
nodes,
node_ids,
_node(
node_id=child_node_id,
node_type=_txn_node_type(txn_type),
title=txn_type,
subtitle=group_title,
status="POSTED",
tone="success",
stage=txn_type,
level=5,
summary="该节点来自生产台账事务,并参与材料批次库外闭环。",
metrics=[
("数量变化", _fmt_qty(child_qty_delta)),
("重量变化", _fmt_weight(child_weight_delta)),
("事务数", len(type_txns)),
],
details=[
("最近库外余额", _fmt_weight(type_txns[-1].outside_weight_after_kg)),
("业务时间", _fmt_datetime(type_txns[-1].biz_time)),
],
raw={
"production_ledger_id": ledger.id,
"txn_ids": [txn.id for txn in type_txns],
"group_key": group_key,
},
),
)
_append_edge(edges, edge_ids, node_ids, _edge(settlement_node_id, child_node_id, txn_type))
if _txn_node_type(txn_type) == "finished_inbound":
for txn in type_txns:
if txn.source_line_id:
finished_node_id_by_lot_id[int(txn.source_line_id)] = child_node_id
if ledger_ids and get_smart_operation_report_enabled(db):
report_rows = db.execute(
select(
OperationReport,
Employee.employee_name.label("employee_name"),
)
.join(Employee, Employee.id == OperationReport.employee_id)
.where(OperationReport.production_ledger_id.in_(ledger_ids))
.order_by(OperationReport.start_time, OperationReport.id)
).all()
for row in report_rows:
report = row.OperationReport
node_id = f"report-{report.id}"
_append_node(
nodes,
node_ids,
_node(
node_id=node_id,
node_type="operation_report",
title=report.report_no,
subtitle=row.employee_name,
status="ABNORMAL" if report.is_abnormal else "NORMAL",
tone="danger" if report.is_abnormal else "success",
stage="现场报工证据",
level=2,
summary="只作为现场报工证据展示,不参与材料批次闭环判定。",
metrics=[
("报工数", _fmt_qty(report.report_qty)),
("合格数", _fmt_qty(report.good_qty)),
("报废数", _fmt_qty(report.scrap_qty)),
],
details=[
("员工", row.employee_name),
("开始", _fmt_datetime(report.start_time)),
("结束", _fmt_datetime(report.end_time)),
("来源", report.report_source),
],
raw={"report_id": report.id, "production_ledger_id": report.production_ledger_id},
),
)
_append_edge(edges, edge_ids, node_ids, _edge(f"ledger-{report.production_ledger_id}", node_id, "报工证据"))
finished_lot = aliased(StockLot)
order_customer = aliased(Customer)
direct_customer = aliased(Customer)
delivery_rows = db.execute(
select(
DeliveryItem.id.label("delivery_item_id"),
DeliveryItem.lot_id.label("delivery_lot_id"),
finished_lot.id.label("finished_lot_id"),
finished_lot.source_doc_id.label("production_ledger_id"),
finished_lot.lot_no.label("lot_no"),
Delivery.delivery_no.label("delivery_no"),
SalesOrder.order_no.label("order_no"),
direct_customer.customer_name.label("direct_customer_name"),
order_customer.customer_name.label("order_customer_name"),
Delivery.delivery_date.label("delivery_date"),
DeliveryItem.delivery_qty.label("delivery_qty"),
DeliveryItem.delivery_weight_kg.label("delivery_weight_kg"),
DeliveryItem.line_amount.label("line_amount"),
DeliveryItem.status.label("status"),
)
.join(finished_lot, finished_lot.id == DeliveryItem.lot_id)
.join(Item, Item.id == finished_lot.item_id)
.join(Delivery, Delivery.id == DeliveryItem.delivery_id)
.outerjoin(SalesOrder, SalesOrder.id == Delivery.sales_order_id)
.outerjoin(order_customer, order_customer.id == SalesOrder.customer_id)
.outerjoin(direct_customer, direct_customer.id == Delivery.customer_id)
.where(
finished_lot.source_doc_id.in_(ledger_ids) if ledger_ids else false(),
finished_lot.source_doc_type.in_(["PRODUCTION_LEDGER_IN", "生产台账"]),
Item.item_type == "FINISHED_GOOD",
)
.order_by(Delivery.delivery_date, DeliveryItem.id)
).mappings().all()
for row in delivery_rows:
delivery_node_id = f"delivery-{row['delivery_item_id']}"
_append_node(
nodes,
node_ids,
_node(
node_id=delivery_node_id,
node_type="delivery",
title=row["delivery_no"],
subtitle=row["order_no"],
status=row["status"],
stage="销售去向",
level=6,
summary="销售出库从生产台账生成的成品库存批次追溯到材料批次。",
metrics=[
("发货数量", _fmt_qty(row["delivery_qty"])),
("发货重量", _fmt_weight(row["delivery_weight_kg"])),
],
details=[
("客户", row["direct_customer_name"] or row["order_customer_name"]),
("库存批次号", row["lot_no"]),
("发货时间", _fmt_datetime(row["delivery_date"])),
],
raw=dict(row),
),
)
source_node_id = finished_node_id_by_lot_id.get(int(row["delivery_lot_id"] or 0))
if not source_node_id:
source_node_id = f"ledger-{row['production_ledger_id']}"
_append_edge(edges, edge_ids, node_ids, _edge(source_node_id, delivery_node_id, "销售出库"))
return nodes, edges