629 lines
30 KiB
Python
629 lines
30 KiB
Python
# ============================================================
|
||
# WMS 缺料事件消费管线(moduleId: domain-wms-events, 可重生 ✅)
|
||
# round-38 方向 N:缺料 → 重排端到端闭环 Mock
|
||
# ① 幂等去重:eventId 全局一次消费,乱序/重复不重复触发,落证据链;
|
||
# ② 库存版本随事件更新并写入证据链(evidenceRefs 带 inventory-version);
|
||
# ③ 缺料 → 影响半径评估(复用 flex 沙盒路径,只读)→ Explore 方案卡(P2 确认);
|
||
# ④ 确认后走 MES 下发(复用 mes.dispatch,P3 门禁)→ 回执镜像写回审计。
|
||
# 说明:方案卡动作复用既有 flex.reschedule(P2,L4 全量重排),
|
||
# 不改 workflow/harness 既有函数;WMS 侧证据由 wms.* 审计事件携带。
|
||
# ============================================================
|
||
from __future__ import annotations
|
||
|
||
import copy
|
||
import uuid
|
||
from datetime import datetime
|
||
from typing import Any
|
||
|
||
from server.agent_core.audit import write_audit
|
||
from server.engines import PoolEngine
|
||
from server.timeutil import fmt_dt
|
||
|
||
World = dict[str, Any]
|
||
|
||
_FLEX_ACTIVE_STATUSES = ("RELEASED", "CREATED")
|
||
|
||
|
||
def _now() -> str:
|
||
return fmt_dt(datetime.now())
|
||
|
||
|
||
def _ensure_tables(world: World) -> None:
|
||
world.setdefault("inventoryVersion", 0)
|
||
world.setdefault("wmsConsumed", {}) # eventId -> {seq, inventoryVersion, at}
|
||
world.setdefault("wmsPending", []) # 乱序缓冲:按 seq 连续性排空
|
||
world.setdefault("wmsLastSeq", 0)
|
||
world.setdefault("wmsInventoryLedger", []) # WMS 库存台账投影(world 侧)
|
||
world.setdefault("wmsTriggered", {}) # eventId -> {confirmId, inventoryVersion, materialCode, at}
|
||
world.setdefault("wmsReceipts", []) # MES 回执镜像(闭环末端)
|
||
|
||
|
||
def _normalize_event(event: dict[str, Any]) -> dict[str, Any]:
|
||
ev = dict(event or {})
|
||
if not str(ev.get("materialCode") or "").strip():
|
||
raise ValueError("WMS 事件缺少 materialCode")
|
||
ev["materialCode"] = str(ev["materialCode"]).strip()
|
||
ev["type"] = str(ev.get("type") or "SHORTAGE").upper()
|
||
if ev.get("newStock") is not None:
|
||
ev["newStock"] = float(ev["newStock"])
|
||
if ev.get("shortageQty") is not None:
|
||
ev["shortageQty"] = float(ev["shortageQty"])
|
||
if ev.get("seq") is not None:
|
||
ev["seq"] = int(ev["seq"])
|
||
return ev
|
||
|
||
|
||
def _is_shortage(ev: dict[str, Any]) -> bool:
|
||
return ev.get("type") == "SHORTAGE" or float(ev.get("shortageQty") or 0) > 0
|
||
|
||
|
||
# ============================================================
|
||
# ① 幂等去重 + ② 库存版本随事件更新(乱序按 seq 缓冲)
|
||
# ============================================================
|
||
def consume_event(store, event: dict[str, Any], actor: str = "wms") -> dict[str, Any]:
|
||
"""消费一条 WMS 事件:eventId 全局一次;乱序按 seq 缓冲排空;版本单调并落证据链。
|
||
|
||
返回 {duplicate, consumed, eventId, seq, materialCode, inventoryVersion, applied}。
|
||
applied 为本次排空批次的消费记录(含各自 inventoryVersion),乱序时可能 >1 条。
|
||
"""
|
||
world = store.data
|
||
_ensure_tables(world)
|
||
ev = _normalize_event(event)
|
||
event_id = str(ev.get("eventId") or "").strip()
|
||
if not event_id:
|
||
raise ValueError("WMS 事件缺少 eventId")
|
||
seq = int(ev.get("seq") or 0)
|
||
consumed = world["wmsConsumed"]
|
||
if event_id in consumed:
|
||
write_audit(world, store.next_id, actor=actor, category="INTEGRATION",
|
||
action="wms.event.duplicate", target={"type": "WMS_EVENT", "id": event_id},
|
||
power="P1", rationale={"seq": seq, "materialCode": ev["materialCode"],
|
||
"firstVersion": consumed[event_id].get("inventoryVersion")},
|
||
result="DUPLICATE", evidence_refs=[f"wms-event:{event_id}"])
|
||
store.save()
|
||
return {"duplicate": True, "consumed": False, "eventId": event_id, "seq": seq,
|
||
"materialCode": ev["materialCode"], "inventoryVersion": world["inventoryVersion"],
|
||
"applied": [], "reason": "already-consumed"}
|
||
# 乱序容忍:追加缓冲,仅排空 seq 连续的批次(不依赖到达顺序)
|
||
world["wmsPending"].append(ev)
|
||
applied = _drain_pending(store, actor)
|
||
store.save()
|
||
mine = next((a for a in applied if a["eventId"] == event_id), None)
|
||
return {"duplicate": False, "consumed": mine is not None, "eventId": event_id, "seq": seq,
|
||
"materialCode": ev["materialCode"], "inventoryVersion": world["inventoryVersion"],
|
||
"applied": applied}
|
||
|
||
|
||
def _drain_pending(store, actor: str) -> list[dict[str, Any]]:
|
||
"""按 seq 连续性排空缓冲:lastSeq+1 起逐条应用(乱序到达按逻辑顺序落版本)。"""
|
||
world = store.data
|
||
applied: list[dict[str, Any]] = []
|
||
while True:
|
||
nxt = int(world.get("wmsLastSeq") or 0) + 1
|
||
idx = next((i for i, e in enumerate(world["wmsPending"])
|
||
if int(e.get("seq") or 0) == nxt), None)
|
||
if idx is None:
|
||
break
|
||
ev = world["wmsPending"].pop(idx)
|
||
applied.append(_apply_event(store, ev, actor))
|
||
world["wmsLastSeq"] = nxt
|
||
return applied
|
||
|
||
|
||
def _apply_event(store, ev: dict[str, Any], actor: str) -> dict[str, Any]:
|
||
"""应用单条事件:库存版本 +1、台账/物料库存更新、冲突入 flexConflicts、证据链审计。"""
|
||
world = store.data
|
||
event_id = ev["eventId"]
|
||
mat_code = ev["materialCode"]
|
||
new_version = int(world.get("inventoryVersion") or 0) + 1
|
||
world["inventoryVersion"] = new_version
|
||
new_stock = ev.get("newStock")
|
||
|
||
# 库存台账(world 侧 WMS 投影)
|
||
ledger = world["wmsInventoryLedger"]
|
||
row = next((r for r in ledger if r["materialCode"] == mat_code), None)
|
||
if row is None:
|
||
row = {"materialCode": mat_code, "name": ev.get("name"), "unit": ev.get("unit"),
|
||
"stock": None, "inTransit": 0, "safetyStock": 0, "version": 0,
|
||
"source": "WMS", "updatedAt": _now()}
|
||
ledger.append(row)
|
||
if new_stock is not None:
|
||
row["stock"] = float(new_stock)
|
||
row["version"] = new_version
|
||
row["updatedAt"] = _now()
|
||
row["lastEventId"] = event_id
|
||
# 同步引擎读的物料库存(flexMaterials / materials 两轨)
|
||
if new_stock is not None:
|
||
for table in ("flexMaterials", "materials"):
|
||
for m in world.get(table, []):
|
||
if m.get("code") == mat_code:
|
||
m["stock"] = float(new_stock)
|
||
break
|
||
# 缺料冲突入冲突中心(WMS 来源,带 eventId/版本;有版本时挂最新 flex 版本)
|
||
if _is_shortage(ev):
|
||
latest_vid = None
|
||
versions = world.get("flexScheduleVersions") or []
|
||
if versions:
|
||
latest_vid = versions[-1].get("id")
|
||
world.setdefault("flexConflicts", []).append({
|
||
"id": store.next_id("flexConflict"),
|
||
"conflictType": "MATERIAL_SHORTAGE",
|
||
"severity": "MAJOR",
|
||
"resourceType": "MATERIAL",
|
||
"resourceName": row.get("name") or mat_code,
|
||
"orderNo": "",
|
||
"description": (f"WMS 缺料事件 {event_id}:{row.get('name') or mat_code} "
|
||
f"库存降至 {new_stock}" +
|
||
(f",缺 {ev['shortageQty']} {row.get('unit') or ''}" if ev.get("shortageQty") else "")),
|
||
"suggestedSolution": "缺料重排:确认 Explore 方案卡后 L4 重排并 MES 下发",
|
||
"versionId": latest_vid,
|
||
"isResolved": False,
|
||
"source": "WMS-EVENT",
|
||
"eventId": event_id,
|
||
"inventoryVersion": new_version,
|
||
})
|
||
# 证据链:eventId + 库存版本
|
||
write_audit(world, store.next_id, actor=actor, category="INTEGRATION",
|
||
action="wms.event.consume", target={"type": "WMS_EVENT", "id": event_id},
|
||
power="P1", rationale={"seq": ev.get("seq"), "offset": ev.get("offset"),
|
||
"materialCode": mat_code, "newStock": new_stock,
|
||
"shortageQty": ev.get("shortageQty"),
|
||
"oldStock": ev.get("oldStock"), "type": ev.get("type")},
|
||
evidence_refs=[f"wms-event:{event_id}", f"inventory-version:{new_version}"])
|
||
world["wmsConsumed"][event_id] = {"seq": int(ev.get("seq") or 0),
|
||
"inventoryVersion": new_version, "at": _now()}
|
||
return {"eventId": event_id, "seq": int(ev.get("seq") or 0),
|
||
"inventoryVersion": new_version, "materialCode": mat_code,
|
||
"newStock": new_stock, "shortageQty": ev.get("shortageQty")}
|
||
|
||
|
||
# ============================================================
|
||
# ③ 影响半径评估(复用 flex 沙盒路径,只读主干)
|
||
# ============================================================
|
||
def _sandbox_counter():
|
||
counters: dict[str, int] = {}
|
||
|
||
def next_id(kind: str) -> int:
|
||
counters[kind] = counters.get(kind, 0) + 1000000
|
||
return counters[kind]
|
||
|
||
return next_id
|
||
|
||
|
||
def _material_row(world: World, code: str) -> dict[str, Any] | None:
|
||
for table in ("flexMaterials", "materials"):
|
||
for m in world.get(table, []):
|
||
if m.get("code") == code:
|
||
return m
|
||
return None
|
||
|
||
|
||
def _bom_products_for_material(world: World, code: str) -> list[str]:
|
||
products: list[str] = []
|
||
for b in world.get("flexBom", []):
|
||
if b.get("materialCode") == code and b.get("productCode") not in products:
|
||
products.append(str(b["productCode"]))
|
||
for bi in world.get("bomItems", []):
|
||
mat = next((m for m in world.get("materials", []) if m.get("id") == bi.get("materialId")), None)
|
||
if mat is not None and mat.get("code") == code:
|
||
prod = next((p for p in world.get("products", []) if p.get("id") == bi.get("productId")), None)
|
||
if prod is not None and prod.get("code") not in products:
|
||
products.append(str(prod["code"]))
|
||
return products
|
||
|
||
|
||
def evaluate_shortage_impact(world: World, event: dict[str, Any]) -> dict[str, Any]:
|
||
"""缺料影响半径(P1 只读):深拷贝沙盒内 flex 试排对比(before 旧库存 vs after 新库存),主干零接触。
|
||
|
||
复用 rush/flex 的沙盒评估路径(深拷贝 + PoolEngine 试排 + 指标差分)。
|
||
"""
|
||
ev = _normalize_event(event)
|
||
mat_code = ev["materialCode"]
|
||
mat = _material_row(world, mat_code)
|
||
products = set(_bom_products_for_material(world, mat_code))
|
||
affected_orders = [
|
||
o for o in world.get("flexOrders", [])
|
||
if o.get("productCode") in products and o.get("status") in _FLEX_ACTIVE_STATUSES
|
||
]
|
||
need_qty = 0.0
|
||
for o in affected_orders:
|
||
for b in world.get("flexBom", []):
|
||
if b.get("productCode") == o.get("productCode") and b.get("materialCode") == mat_code:
|
||
need_qty += float(b.get("quantity") or 0) * float(o.get("quantity") or 0)
|
||
stock = ev.get("newStock")
|
||
if stock is None and mat is not None:
|
||
stock = float(mat.get("stock") or 0)
|
||
in_transit = float(mat.get("inTransit") or 0) if mat is not None else 0.0
|
||
shortage = max(0.0, round(need_qty - float(stock or 0) - in_transit, 3))
|
||
eval_id = uuid.uuid4().hex[:10]
|
||
|
||
if not affected_orders or not world.get("flexOrders"):
|
||
return {"evalId": eval_id, "materialCode": mat_code,
|
||
"materialName": (mat or {}).get("name") or ev.get("name") or mat_code,
|
||
"unit": (mat or {}).get("unit") or ev.get("unit"),
|
||
"oldStock": ev.get("oldStock"), "stock": stock,
|
||
"needQty": round(need_qty, 3), "shortageQty": shortage,
|
||
"affectedOrderCount": 0, "affectedOrders": [],
|
||
"baseline": None, "after": None, "delayDelta": 0.0, "conflictDelta": 0,
|
||
"strategy": "FLEX-BOTTLENECK", "readOnly": True,
|
||
"note": "无受影响柔性订单,无需重排"}
|
||
|
||
def _solve(stock_override):
|
||
sandbox = copy.deepcopy(world)
|
||
if stock_override is not None:
|
||
for m in sandbox.get("flexMaterials", []):
|
||
if m.get("code") == mat_code:
|
||
m["stock"] = float(stock_override)
|
||
break
|
||
sandbox["flexScheduleVersions"] = []
|
||
sandbox["flexVirtualLines"] = []
|
||
sandbox["flexWorkOrders"] = []
|
||
sandbox["flexConflicts"] = []
|
||
result = PoolEngine().solve(sandbox, _sandbox_counter(), sort_mode="BOTTLENECK")
|
||
promised = {
|
||
vl.get("orderNo"): vl.get("plannedEnd")
|
||
for vl in sandbox.get("flexVirtualLines", [])
|
||
if vl.get("orderNo") and vl.get("plannedEnd")
|
||
}
|
||
return result, promised
|
||
|
||
old_stock = ev.get("oldStock")
|
||
before_res, before_promised = _solve(old_stock if old_stock is not None else (stock or 0.0))
|
||
after_res, after_promised = _solve(stock)
|
||
|
||
# 影响半径:物料整体缺料(聚合需求 > 可用量)时,所有在排消费订单都受影响;
|
||
# 附加沙盒对比中完工时点后移的订单(delay 半径),便于方案卡展示严重度。
|
||
if shortage > 0:
|
||
affected = []
|
||
for o in affected_orders:
|
||
on = o.get("orderNo")
|
||
b_end, a_end = before_promised.get(on), after_promised.get(on)
|
||
row = {"orderNo": on, "productCode": o.get("productCode"),
|
||
"quantity": o.get("quantity"), "dueDate": o.get("dueDate"),
|
||
"beforeEnd": b_end, "afterEnd": a_end,
|
||
"endShifted": bool(b_end and a_end and a_end > b_end)}
|
||
affected.append(row)
|
||
affected.sort(key=lambda x: str(x.get("afterEnd") or ""), reverse=True)
|
||
else:
|
||
affected = []
|
||
|
||
return {
|
||
"evalId": eval_id, "materialCode": mat_code,
|
||
"materialName": (mat or {}).get("name") or ev.get("name") or mat_code,
|
||
"unit": (mat or {}).get("unit") or ev.get("unit"),
|
||
"oldStock": old_stock, "stock": stock,
|
||
"needQty": round(need_qty, 3), "shortageQty": shortage,
|
||
"affectedOrderCount": len(affected),
|
||
"affectedOrders": affected[:12],
|
||
"baseline": {"orderCount": before_res["orderCount"],
|
||
"conflictCount": before_res["conflictCount"],
|
||
"totalTardiness": round(before_res["totalTardiness"], 1)},
|
||
"after": {"orderCount": after_res["orderCount"],
|
||
"conflictCount": after_res["conflictCount"],
|
||
"totalTardiness": round(after_res["totalTardiness"], 1)},
|
||
"delayDelta": round(after_res["totalTardiness"] - before_res["totalTardiness"], 1),
|
||
"conflictDelta": after_res["conflictCount"] - before_res["conflictCount"],
|
||
"strategy": "FLEX-BOTTLENECK",
|
||
"readOnly": True,
|
||
}
|
||
|
||
|
||
# ============================================================
|
||
# ③ Explore 方案卡(DRAFT,P2 确认;动作复用 flex.reschedule)
|
||
# ============================================================
|
||
def stage_shortage_solution(store, event: dict[str, Any],
|
||
session_id: str = "wms", actor: str = "wms") -> dict[str, Any]:
|
||
"""缺料 → 影响半径(沙盒只读)→ 生成 Explore 方案卡(P2 确认,L4 重排)。
|
||
|
||
幂等触发:同一 eventId 只出一张卡;重复/乱序重放不会重复触发。
|
||
"""
|
||
from server.agent_core import harness
|
||
world = store.data
|
||
_ensure_tables(world)
|
||
ev = _normalize_event(event)
|
||
event_id = str(ev.get("eventId") or "").strip()
|
||
if not event_id:
|
||
raise ValueError("WMS 事件缺少 eventId")
|
||
if not _is_shortage(ev):
|
||
return {"staged": False, "reason": "not-a-shortage", "eventId": event_id}
|
||
triggered = world["wmsTriggered"].get(event_id)
|
||
if triggered:
|
||
return {"staged": False, "reason": "already-triggered", "eventId": event_id,
|
||
"confirmId": triggered.get("confirmId")}
|
||
consumed = world["wmsConsumed"].get(event_id)
|
||
if not consumed:
|
||
return {"staged": False, "reason": "event-not-consumed", "eventId": event_id}
|
||
inv_version = int(consumed["inventoryVersion"])
|
||
impact = evaluate_shortage_impact(world, ev)
|
||
if not impact.get("affectedOrders"):
|
||
return {"staged": False, "reason": "no-impact", "eventId": event_id, "impact": impact}
|
||
refs = [f"wms-event:{event_id}", f"inventory-version:{inv_version}"]
|
||
affected_nos = [o["orderNo"] for o in impact["affectedOrders"]]
|
||
params = {
|
||
"level": "L4", "sortMode": "BOTTLENECK", # 确认后 flex.reschedule → L4 全量重排
|
||
"eventId": event_id, "materialCode": ev["materialCode"],
|
||
"inventoryVersion": inv_version, "shortageQty": ev.get("shortageQty"),
|
||
"evalId": impact["evalId"], "affectedOrderNos": affected_nos,
|
||
"impact": impact,
|
||
}
|
||
base, aft = impact["baseline"], impact["after"]
|
||
block = harness.stage_confirmation(
|
||
session_id, "flex.reschedule", params,
|
||
title=f"缺料重排方案({impact['materialName']} 缺 {impact['shortageQty']} {impact.get('unit') or ''})",
|
||
summary_lines=[
|
||
f"WMS 事件 {event_id}:{impact['materialName']} 库存 {impact.get('oldStock')} → {impact.get('stock')}(影响半径 {impact['affectedOrderCount']} 单)",
|
||
"受影响订单:" + "、".join(affected_nos[:6]) + ("…" if len(affected_nos) > 6 else ""),
|
||
(f"沙盒对比(只读):冲突 {base['conflictCount']} → {aft['conflictCount']};"
|
||
f"延迟 {base['totalTardiness']}h → {aft['totalTardiness']}h"),
|
||
"确认后执行 L4 全量重排(DRAFT 版本),随后可走 MES 下发(P3 门禁)",
|
||
],
|
||
evidence_refs=refs,
|
||
)
|
||
confirm_id = str(block.props["confirmId"])
|
||
world["wmsTriggered"][event_id] = {"confirmId": confirm_id, "inventoryVersion": inv_version,
|
||
"materialCode": ev["materialCode"], "at": _now(),
|
||
"evalId": impact["evalId"]}
|
||
write_audit(world, store.next_id, actor=actor, category="GATE",
|
||
action="wms.event.stage", target={"type": "WMS_SHORTAGE", "id": event_id},
|
||
power="P1", rationale={"confirmId": confirm_id, "materialCode": ev["materialCode"],
|
||
"affectedOrderCount": impact["affectedOrderCount"],
|
||
"inventoryVersion": inv_version, "evalId": impact["evalId"]},
|
||
evidence_refs=refs)
|
||
store.save()
|
||
return {"staged": True, "duplicate": False, "eventId": event_id,
|
||
"confirmId": confirm_id, "block": block.model_dump(), "impact": impact,
|
||
"inventoryVersion": inv_version}
|
||
|
||
|
||
# ============================================================
|
||
# 上报入口(POST /api/wms/events):消费 → 版本 → 影响 → 方案卡
|
||
# ============================================================
|
||
def consume_shortage_event(store, payload: dict[str, Any],
|
||
session_id: str = "wms", actor: str = "wms") -> dict[str, Any]:
|
||
"""WMS 缺料/库存事件上报入口(幂等)。
|
||
|
||
事件物化:外部 eventId 直接使用(并录入镜像流供重放);缺省经 stub.emit_shortage 生成。
|
||
返回:消费结果 + (缺料事件)Explore 方案卡 staged/impact/card。
|
||
"""
|
||
event = _materialize_event(payload)
|
||
result = consume_event(store, event, actor=actor)
|
||
if result.get("duplicate"):
|
||
return result
|
||
if _is_shortage(event):
|
||
staged = stage_shortage_solution(store, event, session_id=session_id, actor=actor)
|
||
result["staged"] = bool(staged.get("staged"))
|
||
result["stageReason"] = staged.get("reason") if not staged.get("staged") else None
|
||
if staged.get("staged"):
|
||
result["card"] = staged["block"]
|
||
result["confirmId"] = staged["confirmId"]
|
||
result["impact"] = staged["impact"]
|
||
result["evalId"] = staged["impact"]["evalId"]
|
||
else:
|
||
result["staged"] = False
|
||
result["stageReason"] = "not-a-shortage"
|
||
return result
|
||
|
||
|
||
def _materialize_event(payload: dict[str, Any]) -> dict[str, Any]:
|
||
"""事件物化(尽量少依赖外部 eventId;同时让镜像流可重放)。"""
|
||
from server.integrations.wms_stub import get_wms_client
|
||
client = get_wms_client()
|
||
p = dict(payload or {})
|
||
if p.get("eventId"):
|
||
event = dict(p)
|
||
event.setdefault("type", "SHORTAGE")
|
||
client.record_external(event)
|
||
return event
|
||
mat = str(p.get("materialCode") or "").strip()
|
||
if not mat:
|
||
raise ValueError("WMS 事件缺少 materialCode")
|
||
new_stock = p.get("newStock")
|
||
if new_stock is None:
|
||
raise ValueError("WMS 缺料事件必须提供 newStock")
|
||
idem_key = p.get("idemKey") or f"api:{mat}:{p.get('occurredAt') or ''}"
|
||
return client.emit_shortage(
|
||
mat, new_stock=float(new_stock), shortage_qty=p.get("shortageQty"),
|
||
occurred_at=p.get("occurredAt"), idem_key=idem_key,
|
||
name=p.get("name"), unit=p.get("unit"))
|
||
|
||
|
||
# ============================================================
|
||
# ④ 确认后 MES 下发(复用 mes.dispatch P3 门禁)+ 回执镜像写回审计
|
||
# ============================================================
|
||
def stage_dispatch(store, event_id: str,
|
||
session_id: str = "wms", actor: str = "wms") -> dict[str, Any]:
|
||
"""缺料闭环 ④a:为最新重排版本暂存 MES 下发(复用 mes.stage_dispatch,P3 门禁)。"""
|
||
from server.aps_domain.mes import stage_dispatch as mes_stage_dispatch
|
||
world = store.data
|
||
_ensure_tables(world)
|
||
triggered = world["wmsTriggered"].get(event_id)
|
||
if not triggered:
|
||
return {"staged": False, "message": f"WMS 事件 {event_id} 未触发方案卡", "block": None}
|
||
r = mes_stage_dispatch(store, track="flex", session_id=session_id, actor=actor)
|
||
if r.get("staged") and r.get("block"):
|
||
write_audit(world, store.next_id, actor=actor, category="GATE",
|
||
action="wms.mes.dispatch.stage",
|
||
target={"type": "WMS_EVENT", "id": event_id}, power="P1",
|
||
rationale={"confirmId": r["block"].props["confirmId"],
|
||
"inventoryVersion": triggered["inventoryVersion"],
|
||
"materialCode": triggered["materialCode"]},
|
||
evidence_refs=[f"wms-event:{event_id}",
|
||
f"inventory-version:{triggered['inventoryVersion']}"])
|
||
store.save()
|
||
return r
|
||
|
||
|
||
def execute_dispatch_receipt(store, *, event_id: str, confirm_id: str, execution_grant: str,
|
||
version_id: int, before_snapshot: str, checkpoint_store,
|
||
evidence_refs: list[str] | None = None,
|
||
actor: str = "wms") -> dict[str, Any]:
|
||
"""Execute governed MES dispatch and mirror the WMS receipt/evidence chain."""
|
||
|
||
from server.aps_domain.mes import apply_dispatch as mes_apply_dispatch
|
||
from server.aps_domain.mes import validate_dispatchable_version
|
||
|
||
world = store.data
|
||
_ensure_tables(world)
|
||
triggered = world["wmsTriggered"].get(event_id)
|
||
dispatch_evidence = list(evidence_refs or [])
|
||
version_ref = f"schedule-version:{version_id}"
|
||
if version_ref not in dispatch_evidence:
|
||
dispatch_evidence.insert(0, version_ref)
|
||
validation = validate_dispatchable_version(world, "flex", version_id)
|
||
evidence_ref = validation.get("evidenceRef")
|
||
if evidence_ref and evidence_ref not in dispatch_evidence:
|
||
dispatch_evidence.append(str(evidence_ref))
|
||
result = mes_apply_dispatch(
|
||
store,
|
||
track="flex",
|
||
actor=actor,
|
||
confirm_id=confirm_id,
|
||
execution_grant=execution_grant,
|
||
version_id=version_id,
|
||
before_snapshot=before_snapshot,
|
||
evidence_refs=dispatch_evidence,
|
||
checkpoint_store=checkpoint_store,
|
||
)
|
||
inventory_version = (
|
||
int(triggered.get("inventoryVersion"))
|
||
if triggered and triggered.get("inventoryVersion") is not None
|
||
else int(world.get("inventoryVersion") or 0)
|
||
)
|
||
external_ids = list(result.get("created") or [])
|
||
receipt = {
|
||
"receiptId": f"WMS-RCPT-{event_id}",
|
||
"eventId": event_id,
|
||
"inventoryVersion": inventory_version,
|
||
"scheduleVersionId": version_id,
|
||
"mesExternalWoIds": external_ids,
|
||
"created": len(external_ids),
|
||
"duplicates": len(result.get("duplicates") or []),
|
||
"at": _now(),
|
||
"source": "WMS-MOCK",
|
||
}
|
||
world.setdefault("wmsReceipts", []).append(receipt)
|
||
write_audit(
|
||
world,
|
||
store.next_id,
|
||
actor=actor,
|
||
category="INTEGRATION",
|
||
action="wms.mes.receipt",
|
||
target={"type": "MES_RECEIPT", "id": receipt["receiptId"]},
|
||
power="P3",
|
||
rationale={
|
||
"eventId": event_id,
|
||
"inventoryVersion": inventory_version,
|
||
"scheduleVersionId": version_id,
|
||
"externalWoIds": external_ids,
|
||
"created": len(external_ids),
|
||
},
|
||
evidence_refs=[
|
||
f"wms-event:{event_id}",
|
||
f"inventory-version:{inventory_version}",
|
||
*dispatch_evidence,
|
||
],
|
||
)
|
||
store.save()
|
||
return {
|
||
**result,
|
||
"receipt": receipt,
|
||
"message": f"{result.get('message')} 回执已入审计({receipt['receiptId']})。",
|
||
}
|
||
|
||
|
||
# ============================================================
|
||
# 只读投影(P0):库存版本/台账 / 事件流重放 / 连接状态
|
||
# ============================================================
|
||
def inventory_projection(world: World) -> dict[str, Any]:
|
||
"""WMS 库存版本/台账投影(P0 只读)。"""
|
||
_ensure_tables(world)
|
||
ledger = world.get("wmsInventoryLedger", [])
|
||
rows = []
|
||
for r in ledger:
|
||
mat = _material_row(world, r["materialCode"])
|
||
rows.append({**r, "apsStock": (float(mat.get("stock") or 0) if mat else r.get("stock"))})
|
||
return {
|
||
"connected": True, "mode": "stub",
|
||
"inventoryVersion": world.get("inventoryVersion", 0),
|
||
"consumedEvents": len(world.get("wmsConsumed", {})),
|
||
"triggeredCards": len(world.get("wmsTriggered", {})),
|
||
"receiptCount": len(world.get("wmsReceipts", [])),
|
||
"pendingEvents": len(world.get("wmsPending", [])),
|
||
"ledger": rows,
|
||
}
|
||
|
||
|
||
def wms_connection_status() -> dict[str, Any]:
|
||
from server.integrations.wms_stub import get_wms_client
|
||
return get_wms_client().status()
|
||
|
||
|
||
def replay_events(limit: int | None = None, order: str = "seq") -> dict[str, Any]:
|
||
"""WMS 事件流重放(P0):seq=逻辑顺序 / arrival=到达顺序。"""
|
||
from server.integrations.wms_stub import get_wms_client
|
||
client = get_wms_client()
|
||
status = client.status()
|
||
events = client.replay(limit=limit, order=order)
|
||
return {"system": status.get("system"), "plant": status.get("plant"),
|
||
"inventoryVersion": client.inventory_version(),
|
||
"count": len(events), "events": events}
|
||
|
||
|
||
# ============================================================
|
||
# round-39 方向 P:Saga/补偿接入点
|
||
# ① run_shortage_saga:WMS 缺料→重排→MES 下发→回执 作为第一个 Saga 实例编排
|
||
# ② rollback_receipt:MES 回执补偿(幂等 VOIDED,留痕不删除)
|
||
# ============================================================
|
||
def rollback_receipt(store, *, event_id: str | None = None, receipt_id: str | None = None,
|
||
actor: str = "saga", reason: str = "saga-compensation") -> dict[str, Any]:
|
||
"""MES 回执补偿(幂等):按 receiptId/eventId 将 wmsReceipts 条目标记 VOIDED(留痕不删除)。
|
||
|
||
- 幂等:同一 receiptId/eventId 已回滚过 → 直接返回记录(duplicate);
|
||
- 无匹配回执(如已随世界快照回滚)→ 返回 voided=0(幂等成功语义)。
|
||
"""
|
||
world = store.data
|
||
_ensure_tables(world)
|
||
world.setdefault("wmsReceiptVoids", {})
|
||
key = receipt_id or event_id
|
||
if key and key in world["wmsReceiptVoids"]:
|
||
return {**world["wmsReceiptVoids"][key], "duplicate": True}
|
||
targets = [r for r in world.get("wmsReceipts", [])
|
||
if (receipt_id and r.get("receiptId") == receipt_id)
|
||
or (event_id and r.get("eventId") == event_id)]
|
||
for r in targets:
|
||
r["status"] = "VOIDED"
|
||
r["voidedAt"] = _now()
|
||
r["voidReason"] = reason
|
||
write_audit(world, store.next_id, actor=actor, category="INTEGRATION",
|
||
action="wms.mes.receipt.rollback",
|
||
target={"type": "MES_RECEIPT", "id": key or "ALL"},
|
||
power="P1", rationale={"eventId": event_id, "receiptId": receipt_id,
|
||
"voided": len(targets), "reason": reason},
|
||
result="SUCCESS",
|
||
evidence_refs=[f"wms-receipt-void:{key}"] if key else [])
|
||
result = {"voided": len(targets), "eventId": event_id, "receiptId": receipt_id,
|
||
"reason": reason, "message": f"回执补偿 ✅ VOIDED {len(targets)} 条(幂等)。"}
|
||
if key:
|
||
world["wmsReceiptVoids"][key] = result
|
||
store.save()
|
||
return result
|
||
|
||
|
||
def run_shortage_saga(store, payload: dict[str, Any], *,
|
||
session_id: str = "wms", actor: str = "saga",
|
||
auto_approve: bool = True, timeout_sec: float = 30.0,
|
||
max_retries: int = 2, coordinator=None) -> dict[str, Any]:
|
||
"""WMS 缺料闭环以 Saga 编排执行(round-39 方向 P,矩阵 78/119)。
|
||
|
||
事件物化 → consume → impact → 方案卡 → 确认重排 → MES 下发 → 回执;
|
||
任一失败 → 重试超限走补偿链(重排失败回滚前版本;MES 失败撤销下发);
|
||
补偿也失败 → MANUAL_TAKEOVER(人工接管点,可经 retry 恢复)。
|
||
返回 saga 记录(含中间态 steps 明细 / 幂等键 / 审计引用)。
|
||
"""
|
||
from server.aps_domain import saga
|
||
event = _materialize_event(payload)
|
||
coord = coordinator or saga.get_coordinator(store, actor=actor)
|
||
record = saga.build_wms_shortage_saga(coord, event, session_id=session_id, actor=actor,
|
||
timeout_sec=timeout_sec, max_retries=max_retries)
|
||
return coord.run(record["id"], auto_approve=auto_approve)
|