# ============================================================ # Mock WMS 缺料/库存事件源(moduleId: integ-wms-stub, 可重生 ✅) # round-38 方向 N:缺料 → 重排端到端闭环 Mock # 与 mes_stub 同构:读/写本地镜像(库存台账 + 缺料事件流); # emit_shortage 幂等生成缺料事件(eventId + 乱序容忍 offset/seq); # replay 按事件流重放;库存版本号单调递增(monotonic)。 # ============================================================ from __future__ import annotations import json from datetime import datetime from pathlib import Path from typing import Any _MIRROR_PATH = Path(__file__).resolve().parents[1] / "data" / "wms_mirror.json" def default_mirror() -> dict[str, Any]: return { "system": "WMS-MOCK", "plant": "CNWH", "updatedAt": datetime.now().strftime("%Y-%m-%d %H:%M"), "inventory": [], # 库存台账 [{materialCode,name,unit,stock,inTransit,safetyStock,version}] "shortageEvents": [], # 缺料/库存事件流 [{eventId,seq,offset,type,materialCode,...}] "idempotency": {}, # dedupKey -> eventId "inventoryVersion": 0, # 单调库存版本号 } def _now_ts() -> str: return datetime.now().strftime("%Y-%m-%d %H:%M:%S") class MockWmsClient: """进程内 Mock WMS:读写本地镜像,模拟缺料事件上报与重放。""" def __init__(self, path: Path | None = None): self.path = path or _MIRROR_PATH self.path.parent.mkdir(parents=True, exist_ok=True) if not self.path.exists(): self._save(default_mirror()) def _load(self) -> dict: try: return json.loads(self.path.read_text(encoding="utf-8")) except (OSError, json.JSONDecodeError): data = default_mirror() self._save(data) return data def _save(self, data: dict) -> None: data["updatedAt"] = datetime.now().strftime("%Y-%m-%d %H:%M") self.path.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8") # ---------------- 只读:状态 / 台账 / 版本 ---------------- def status(self) -> dict: m = self._load() return { "connected": True, "mode": "stub", "system": m.get("system"), "plant": m.get("plant"), "updatedAt": m.get("updatedAt"), "inventoryVersion": int(m.get("inventoryVersion") or 0), "inventoryCount": len(m.get("inventory") or []), "eventCount": len(m.get("shortageEvents") or []), "shortageCount": sum(1 for e in m.get("shortageEvents") or [] if e.get("type") == "SHORTAGE"), } def inventory_version(self) -> int: """库存版本号(monotonic,随每次事件 +1)。""" return int(self._load().get("inventoryVersion") or 0) def ledger(self) -> list[dict]: """库存台账快照(P0 只读)。""" return list(self._load().get("inventory") or []) # ---------------- 写:库存台账(初始同步/测试种子) ---------------- def upsert_inventory(self, rows: list[dict[str, Any]]) -> list[dict]: """写入/更新库存台账(幂等按 materialCode 合并)。返回当前台账。""" m = self._load() existing = {r.get("materialCode"): r for r in m.get("inventory", [])} for row in rows: code = row.get("materialCode") if not code: continue if code in existing: existing[code].update({k: v for k, v in row.items() if v is not None}) else: m.setdefault("inventory", []).append(dict(row)) self._save(m) return self.ledger() # ---------------- 写:缺料事件(幂等 eventId + 乱序容忍 offset/seq) ---------------- def emit_shortage(self, material_code: str, *, new_stock: float, shortage_qty: float | None = None, occurred_at: str | None = None, idem_key: str | None = None, name: str | None = None, unit: str | None = None) -> dict: """生成一条缺料事件并写镜像事件流。 幂等:相同 idem_key(缺省 = materialCode + occurredAt)只生成一次, 重复调用返回同一事件(duplicate=True)。 乱序容忍:事件携带单调 seq(逻辑顺序)与 offset(到达顺序), 消费端按 seq 应用、按 eventId 去重,不依赖到达顺序。 """ m = self._load() key = idem_key or f"shortage:{material_code}:{occurred_at or ''}" seen = (m.get("idempotency") or {}).get(key) if seen: ev = next((e for e in m.get("shortageEvents", []) if e.get("eventId") == seen), None) if ev is not None: return {**ev, "duplicate": True} row = next((r for r in m.get("inventory", []) if r.get("materialCode") == material_code), None) seq = len(m.get("shortageEvents") or []) + 1 event_id = f"WMS-EVT-{seq:04d}" new_version = int(m.get("inventoryVersion") or 0) + 1 event = { "eventId": event_id, "seq": seq, "offset": len(m.get("shortageEvents") or []), "type": "SHORTAGE", "materialCode": material_code, "name": name or (row or {}).get("name"), "unit": unit or (row or {}).get("unit"), "oldStock": (float(row["stock"]) if row is not None and row.get("stock") is not None else None), "newStock": float(new_stock), "shortageQty": float(shortage_qty) if shortage_qty is not None else None, "occurredAt": occurred_at or _now_ts(), "dedupKey": key, "source": "WMS-MOCK", "wmsVersion": new_version, } m.setdefault("shortageEvents", []).append(event) m.setdefault("idempotency", {})[key] = event_id if row is None: m.setdefault("inventory", []).append({ "materialCode": material_code, "name": name or material_code, "unit": unit or "", "stock": float(new_stock), "inTransit": 0, "safetyStock": 0, "version": new_version, }) else: row["stock"] = float(new_stock) row["version"] = new_version m["inventoryVersion"] = new_version self._save(m) return {**event, "duplicate": False} def record_external(self, event: dict[str, Any]) -> dict: """把外部上报事件追加进镜像流(供 replay 覆盖),按 eventId 幂等。""" m = self._load() eid = str(event.get("eventId") or "").strip() if not eid: raise ValueError("外部 WMS 事件缺少 eventId") for e in m.get("shortageEvents", []): if e.get("eventId") == eid: return e seq = len(m.get("shortageEvents") or []) + 1 new_version = int(m.get("inventoryVersion") or 0) + 1 rec = { "eventId": eid, "seq": seq, "offset": len(m.get("shortageEvents") or []), "type": str(event.get("type") or "SHORTAGE").upper(), "materialCode": str(event.get("materialCode") or ""), "name": event.get("name"), "unit": event.get("unit"), "oldStock": event.get("oldStock"), "newStock": event.get("newStock"), "shortageQty": event.get("shortageQty"), "occurredAt": event.get("occurredAt") or _now_ts(), "dedupKey": f"external:{eid}", "source": str(event.get("source") or "WMS-EXTERNAL"), "wmsVersion": new_version, } m.setdefault("shortageEvents", []).append(rec) m["inventoryVersion"] = new_version self._save(m) return rec # ---------------- 读:事件流重放 ---------------- def replay(self, limit: int | None = None, order: str = "seq") -> list[dict]: """按事件流重放(P0)。order: seq=逻辑顺序(乱序容忍基准)/ arrival=到达顺序。""" events = list(self._load().get("shortageEvents") or []) if (order or "seq").lower() == "arrival": events.sort(key=lambda e: int(e.get("offset") or 0)) else: events.sort(key=lambda e: int(e.get("seq") or 0)) if limit is not None: events = events[: max(0, int(limit))] return events def reset(self) -> dict: data = default_mirror() self._save(data) return self.status() _client: MockWmsClient | None = None def get_wms_client() -> MockWmsClient: global _client if _client is None: _client = MockWmsClient() return _client def reset_wms_client(path: Path | None = None) -> MockWmsClient: global _client _client = MockWmsClient(path) _client.reset() return _client