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