aps-agent/server/aps_domain/wms_events.py

629 lines
30 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# ============================================================
# 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)