aps-agent/tests/golden/test_saga_wms_closure.py

274 lines
13 KiB
Python
Raw Permalink Normal View History

# ============================================================
# round-39 方向 P:WMS 缺料闭环 Saga 黄金测试(矩阵 78/119 验收)
# 固化:全链成功 / 中途失败自动补偿(重排回滚 / MES 撤销)/ 人工接管 / 门禁暂停恢复
# ============================================================
from __future__ import annotations
from pathlib import Path
import pytest
from server.aps_domain import saga, wms_events
from server.integrations.mes_stub import get_mes_client, reset_mes_client
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) -> None:
pass
@pytest.fixture(autouse=True)
def _wf_checkpoints_in_tmp(monkeypatch, tmp_path: Path):
"""把 workflow.execute_confirmed 的内部检查点也隔离到 tmp(不污染 server/data)。"""
import server.aps_domain.workflow as workflow_module
cps = CheckpointStore(str(tmp_path / "wf_checkpoints.json"))
monkeypatch.setattr(workflow_module, "get_checkpoints", lambda: cps)
yield cps
def _baseline(store):
from server.aps_domain.flex import run_flex_schedule
run_flex_schedule(store, sort_mode="BOTTLENECK", actor="test")
return store.data["flexScheduleVersions"][-1]
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 _shortage_event(client, material_code: str, new_stock: float, shortage_qty: float,
occurred_at: str):
return client.emit_shortage(material_code, new_stock=new_stock,
shortage_qty=shortage_qty, occurred_at=occurred_at)
def _fresh_closure(tmp_path: Path, tag: str):
"""标准 WMS 缺料闭环现场:基线排产 + 库存台账 + 缺料事件。"""
reset_mes_client(tmp_path / f"mes_{tag}.json")
reset_wms_client(tmp_path / f"wms_{tag}.json")
store = _MemStore(seed_world(), tmp_path / f"cp_{tag}.json")
_baseline(store)
client = reset_wms_client(tmp_path / f"wms_{tag}.json")
_seed_ledger(client)
ev = _shortage_event(client, "WIRE-HV", 500, 220, "2026-08-02 09:00:00")
return store, ev
def _audit_actions(store) -> list[str]:
return [e["action"] for e in store.data["auditEvents"]]
# ------------------------------------------------------------
# ① 全链成功:consume → impact → 方案卡 → 重排 → MES 下发 → 回执
# ------------------------------------------------------------
def test_wms_closure_saga_full_chain_success(tmp_path: Path):
store, ev = _fresh_closure(tmp_path, "ok")
versions_before = len(store.data["flexScheduleVersions"])
coord = saga.get_coordinator(store, actor="tester")
record = saga.build_wms_shortage_saga(coord, ev, session_id="s-wms", actor="tester")
out = coord.run(record["id"], auto_approve=True)
assert out["status"] == "SUCCEEDED", out["steps"]
steps = {s["name"]: s for s in out["steps"]}
for name in ("consume_event", "evaluate_impact", "stage_solution",
"reschedule", "stage_dispatch", "dispatch"):
assert steps[name]["status"] == "SUCCEEDED", name
assert steps[name]["idemKey"] # 每步都定义了幂等键
assert steps[name]["maxRetries"] >= 0
assert steps[name]["timeoutSec"] > 0
# 幂等键登记
assert f"wms-event:{ev['eventId']}" in store.data["sagaIdem"]
assert f"dispatch:{ev['eventId']}" in store.data["sagaIdem"]
# 闭环落地
assert len(store.data["wmsReceipts"]) == 1
assert store.data["flexScheduleVersions"][-1]["status"] == "DISPATCHED"
assert len(store.data["flexScheduleVersions"]) == versions_before + 1
assert get_mes_client().status()["woCount"] >= 1
# 审计链完整(saga 事件 + 既有闭环动作)
actions = _audit_actions(store)
for expected in ("saga.created", "saga.succeeded", "wms.event.consume",
"flex.reschedule.approve", "mes.dispatch", "wms.mes.receipt"):
assert expected in actions, expected
# 每步审计含 idemKey
step_audit = next(e for e in store.data["auditEvents"]
if e["action"] == "saga.step.succeeded")
assert step_audit["rationale"]["idemKey"]
# ------------------------------------------------------------
# ② 重排失败(副作用已发生)→ 自动补偿:回滚前版本
# ------------------------------------------------------------
def test_wms_closure_saga_reschedule_failure_compensates(tmp_path: Path):
store, ev = _fresh_closure(tmp_path, "resched")
versions_before = len(store.data["flexScheduleVersions"])
coord = saga.get_coordinator(store, actor="tester")
real_flex = saga.default_registry()["flex.reschedule"]["fn"]
def failing_reschedule(ctx):
real_flex(ctx) # 副作用已发生(新 DRAFT 版本)
raise saga.SagaStepFailed("simulated reschedule failure after effect")
coord.registry["flex.reschedule"]["fn"] = failing_reschedule
record = saga.build_wms_shortage_saga(coord, ev, session_id="s-wms", actor="tester")
out = coord.run(record["id"], auto_approve=True)
assert out["status"] == "COMPENSATED"
res = next(s for s in out["steps"] if s["name"] == "reschedule")
assert res["status"] == "FAILED"
assert res["compensation"]["status"] == "SUCCEEDED"
# 重排失败 → 回滚前版本:版本数回到基线
assert len(store.data["flexScheduleVersions"]) == versions_before
# 无 MES 下发 / 无回执
assert not store.data.get("wmsReceipts")
assert get_mes_client().status()["woCount"] == 0
assert "saga.compensated" in _audit_actions(store)
# ------------------------------------------------------------
# ③ MES 下发失败(外部已提交)→ 自动补偿:撤销下发 + 回执补偿 + 世界回滚
# ------------------------------------------------------------
def test_wms_closure_saga_dispatch_failure_cancels_mes(tmp_path: Path):
store, ev = _fresh_closure(tmp_path, "dispatch")
versions_before = len(store.data["flexScheduleVersions"])
coord = saga.get_coordinator(store, actor="tester")
real_dispatch = saga.default_registry()["mes.dispatch"]["fn"]
def failing_dispatch(ctx):
result = real_dispatch(ctx) # 外部副作用已提交(MES 工单 + 回执)
ctx["step"]["result"] = result # 保留部分结果供补偿撤销
raise saga.SagaStepFailed("simulated dispatch failure after external commit")
coord.registry["mes.dispatch"]["fn"] = failing_dispatch
record = saga.build_wms_shortage_saga(coord, ev, session_id="s-wms", actor="tester")
out = coord.run(record["id"], auto_approve=True)
assert out["status"] == "COMPENSATED"
disp = next(s for s in out["steps"] if s["name"] == "dispatch")
assert disp["status"] == "FAILED"
assert disp["compensation"]["status"] == "SUCCEEDED"
# 外部 MES 工单已全部撤销(幂等 CANCELLED)
wos = get_mes_client()._load().get("workOrders", [])
assert wos and all(w["status"] == "CANCELLED" for w in wos)
# 补偿链落地:撤销下发 + 回执补偿审计
actions = _audit_actions(store)
assert "mes.cancel_dispatch" in actions
assert "wms.mes.receipt.rollback" in actions
# 世界回滚:无回执、版本回到基线(全链补偿到初始态,可整体重试)
assert not store.data.get("wmsReceipts")
assert len(store.data["flexScheduleVersions"]) == versions_before
# ------------------------------------------------------------
# ④ 补偿也失败 → MANUAL_TAKEOVER(人工接管点明确)
# ------------------------------------------------------------
def test_wms_closure_saga_compensation_failure_manual_takeover(tmp_path: Path):
store, ev = _fresh_closure(tmp_path, "takeover")
coord = saga.get_coordinator(store, actor="tester")
def failing_reschedule(_ctx):
raise saga.SagaStepFailed("simulated reschedule failure")
def broken_rollback(_ctx):
raise saga.SagaCompensationError("checkpoint restore unavailable")
coord.registry["flex.reschedule"]["fn"] = failing_reschedule
coord.registry["flex.reschedule.rollback"]["fn"] = broken_rollback
record = saga.build_wms_shortage_saga(coord, ev, session_id="s-wms", actor="tester")
out = coord.run(record["id"], auto_approve=True)
assert out["status"] == "MANUAL_TAKEOVER"
assert out["manualTakeover"] is True
assert "compensation-failed" in out["takeoverReason"]
assert "saga.manual_takeover" in _audit_actions(store)
assert "saga.step.compensation_failed" in _audit_actions(store)
# ------------------------------------------------------------
# ⑤ 门禁暂停(auto_approve=False)→ 中间态可见 → 恢复
# ------------------------------------------------------------
def test_wms_closure_saga_gate_pause_and_resume(tmp_path: Path):
store, ev = _fresh_closure(tmp_path, "gate")
coord = saga.get_coordinator(store, actor="tester")
record = saga.build_wms_shortage_saga(coord, ev, session_id="s-wms", actor="tester")
out = coord.run(record["id"], auto_approve=False)
assert out["status"] == "WAITING_HUMAN"
res = next(s for s in out["steps"] if s["name"] == "reschedule")
assert res["status"] == "WAITING_HUMAN"
assert res["result"] is None # 人工门禁前未执行
stage = next(s for s in out["steps"] if s["name"] == "stage_solution")
assert stage["status"] == "SUCCEEDED"
assert stage["result"]["confirmId"] # 中间态可见:方案卡已出
# 人工放行后自动续跑(auto_approve=True 代行确认)
out2 = coord.run(record["id"], auto_approve=True)
assert out2["status"] == "SUCCEEDED"
assert len(store.data["wmsReceipts"]) == 1
# ------------------------------------------------------------
# ⑥ run_shortage_saga 入口:事件物化 + 幂等去重
# ------------------------------------------------------------
def test_run_shortage_saga_entry_idempotent(tmp_path: Path):
store, ev = _fresh_closure(tmp_path, "entry")
r1 = wms_events.run_shortage_saga(store, ev, session_id="s-wms", actor="tester")
assert r1["status"] == "SUCCEEDED"
assert len(store.data["wmsReceipts"]) == 1
# 同一事件重复触发 → 复用已有 saga(不重复下发)
r2 = wms_events.run_shortage_saga(store, ev, session_id="s-wms", actor="tester")
assert r2["id"] == r1["id"]
assert len(store.data["sagas"]) == 1
assert len(store.data["wmsReceipts"]) == 1
assert len(store.data["wmsReceipts"]) == 1
# ------------------------------------------------------------
# Round 65:重排确认被拒绝时,绝不能把既有 DRAFT 当成本次新版本继续发布?
# ------------------------------------------------------------
def test_wms_closure_saga_rejected_reschedule_never_publishes_old_draft(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
):
store, ev = _fresh_closure(tmp_path, "rejected-reschedule")
versions_before = list(store.data["flexScheduleVersions"])
old_version_id = versions_before[-1]["id"]
import server.aps_domain.workflow as workflow_module
monkeypatch.setattr(
workflow_module,
"execute_confirmed",
lambda *_args, **_kwargs: "证据校验未通过,未执行任何变更:source drift",
)
coord = saga.get_coordinator(store, actor="tester")
record = saga.build_wms_shortage_saga(
coord, ev, session_id="s-wms", actor="tester"
)
out = coord.run(record["id"], auto_approve=True)
steps = {step["name"]: step for step in out["steps"]}
assert out["status"] == "COMPENSATED"
assert steps["reschedule"]["status"] == "FAILED"
assert "未创建唯一的新排产版本" in steps["reschedule"]["error"]
assert steps["stage_publish"]["status"] == "PENDING"
assert steps["publish"]["status"] == "PENDING"
assert steps["stage_dispatch"]["status"] == "PENDING"
assert steps["dispatch"]["status"] == "PENDING"
assert len(store.data["flexScheduleVersions"]) == len(versions_before)
assert store.data["flexScheduleVersions"][-1]["id"] == old_version_id
assert store.data["flexScheduleVersions"][-1]["status"] == "DRAFT"
assert not store.data.get("wmsReceipts")
assert get_mes_client().status()["woCount"] == 0