# ============================================================ # round-38 方向 N:WMS 事件流重放黄金测试 # 固化:事件可重放(seq/arrival 顺序);重复重放幂等不重复触发; # emit_shortage 幂等(idem_key);库存版本单调。 # ============================================================ from __future__ import annotations from pathlib import Path from server.aps_domain import wms_events from server.integrations.wms_stub import reset_wms_client from server.state.checkpoints import CheckpointStore from server.state.seed import seed_world class _MemStore: def __init__(self, data, checkpoint_path: Path): self.data = data self.checkpoints = CheckpointStore(str(checkpoint_path)) def next_id(self, kind: str) -> int: key = f"_c_{kind}" self.data[key] = self.data.get(key, 4000) + 1 return self.data[key] def save(self): pass def _seed_ledger(client): client.upsert_inventory([ {"materialCode": "WIRE-HV", "name": "高压线材", "unit": "米", "stock": 12000, "inTransit": 0, "safetyStock": 1000}, {"materialCode": "TERM-HV", "name": "高压端子", "unit": "个", "stock": 6000, "inTransit": 2000, "safetyStock": 500}, {"materialCode": "BUSBAR", "name": "铜排", "unit": "根", "stock": 1500, "inTransit": 400, "safetyStock": 200}, ]) def test_emit_shortage_idempotent_by_idem_key(tmp_path: Path): """emit_shortage:相同 idem_key 只生成一次事件(重复调用返回同一 eventId)。""" reset_wms_client(tmp_path / "wms_mirror.json") client = reset_wms_client(tmp_path / "wms_mirror.json") _seed_ledger(client) ev1 = client.emit_shortage("WIRE-HV", new_stock=500, shortage_qty=220, occurred_at="2026-08-02 09:00:00", idem_key="k-wire") ev2 = client.emit_shortage("WIRE-HV", new_stock=500, shortage_qty=220, occurred_at="2026-08-02 09:00:00", idem_key="k-wire") assert ev1["duplicate"] is False and ev2["duplicate"] is True assert ev2["eventId"] == ev1["eventId"] assert ev2["seq"] == ev1["seq"] # 缺省 idem_key(materialCode + occurredAt)同样幂等 ev3 = client.emit_shortage("TERM-HV", new_stock=3000, occurred_at="2026-08-02 10:00:00") ev4 = client.emit_shortage("TERM-HV", new_stock=3000, occurred_at="2026-08-02 10:00:00") assert ev4["duplicate"] is True and ev4["eventId"] == ev3["eventId"] # 事件流单调:seq 递增、库存版本递增 assert [e["seq"] for e in client.replay()] == [1, 2] assert client.inventory_version() == 2 def test_replay_stream_order_limit_and_version(tmp_path: Path): """replay:seq 逻辑顺序 / arrival 到达顺序 / limit 截断;版本单调。""" reset_wms_client(tmp_path / "wms_mirror.json") client = reset_wms_client(tmp_path / "wms_mirror.json") _seed_ledger(client) ev1 = client.emit_shortage("WIRE-HV", new_stock=500, shortage_qty=220, occurred_at="2026-08-02 09:00:00") ev2 = client.emit_shortage("TERM-HV", new_stock=3000, occurred_at="2026-08-02 10:00:00") ev3 = client.emit_shortage("BUSBAR", new_stock=200, shortage_qty=180, occurred_at="2026-08-02 11:00:00") seq_order = client.replay(order="seq") assert [e["eventId"] for e in seq_order] == [ev1["eventId"], ev2["eventId"], ev3["eventId"]] arrival = client.replay(order="arrival") assert [e["offset"] for e in arrival] == [0, 1, 2] assert [e["eventId"] for e in client.replay(limit=2)] == [ev1["eventId"], ev2["eventId"]] assert client.inventory_version() == 3 assert client.status()["eventCount"] == 3 assert client.status()["shortageCount"] == 3 # 事件带乱序容忍元数据(seq 逻辑顺序 / offset 到达顺序) assert all(e.get("seq") is not None and e.get("offset") is not None for e in seq_order) def test_replayed_events_consume_idempotently(tmp_path: Path): """重放的事件可重复消费:第一次消费落证据链,第二次全部幂等跳过(不重复触发)。""" reset_wms_client(tmp_path / "wms_mirror.json") store = _MemStore(seed_world(), tmp_path / "checkpoints.json") from server.aps_domain.flex import run_flex_schedule run_flex_schedule(store, sort_mode="BOTTLENECK", actor="test") client = reset_wms_client(tmp_path / "wms_mirror.json") _seed_ledger(client) for i, (mat, ns, sq) in enumerate([("WIRE-HV", 500, 220), ("TERM-HV", 3000, 0), ("BUSBAR", 200, 180)]): client.emit_shortage(mat, new_stock=ns, shortage_qty=sq, occurred_at=f"2026-08-02 0{i + 9}:00:00") stream = client.replay(order="seq") # 第一次消费:3 条全部消费,库存版本 1→2→3,证据链逐条落 inventory-version for ev in stream: r = wms_events.consume_event(store, ev, actor="test") assert r["consumed"] is True and r["duplicate"] is False assert store.data["inventoryVersion"] == 3 assert len(store.data["wmsConsumed"]) == 3 consume_refs = [e["evidenceRefs"] for e in store.data["auditEvents"] if e["action"] == "wms.event.consume"] assert sorted(r[1] for r in consume_refs) == \ ["inventory-version:1", "inventory-version:2", "inventory-version:3"] assert len(store.data["wmsPending"]) == 0 # 第二次重放消费:全部幂等(duplicate),不再产生 consume 证据、版本不涨 for ev in stream: r = wms_events.consume_event(store, ev, actor="test") assert r["duplicate"] is True and r["consumed"] is False assert store.data["inventoryVersion"] == 3 consume_after = [e for e in store.data["auditEvents"] if e["action"] == "wms.event.consume"] assert len(consume_after) == 3 # 不重复触发 dup_count = sum(1 for e in store.data["auditEvents"] if e["action"] == "wms.event.duplicate") assert dup_count == 3 # 重放投影接口 replay_view = wms_events.replay_events(limit=None, order="seq") assert replay_view["count"] == 3 assert replay_view["inventoryVersion"] == 3