# ============================================================ # round-39 方向 P:Saga/补偿编排框架黄金测试 # 固化:状态机 / 幂等 / 重试 / 超时 / 补偿链 / 人工接管 / 持久化恢复 # 矩阵 78:任一步失败可重试或补偿;中间态可见;人工接管点和审计记录明确 # 矩阵 119:每个 P3 动作定义幂等键、超时、重试、补偿、人工接管和审计事件 # ============================================================ from __future__ import annotations import copy from pathlib import Path from fastapi.testclient import TestClient from server.aps_domain import saga from server.state.checkpoints import CheckpointStore from server.state.seed import seed_world class _MemStore: def __init__(self, data=None, checkpoint_path: Path | None = None): self.data = data if data is not None else seed_world() self.checkpoints = CheckpointStore(str(checkpoint_path or Path(__file__).parent / "_saga_cp.json")) def next_id(self, kind: str) -> int: key = f"_c_{kind}" self.data[key] = self.data.get(key, 1000) + 1 return self.data[key] def save(self) -> None: pass def make_action(name: str, log: list[str], *, fails: int = 0, timeout: bool = False): state = {"calls": 0} def fn(ctx) -> dict: state["calls"] += 1 log.append(name) if timeout: raise saga.SagaStepTimeout("simulated timeout") if state["calls"] <= fails: raise saga.SagaStepFailed("simulated failure") return {"done": name, "calls": state["calls"]} return fn def make_registry(log: list[str], *, s2_fails: int = 0, s3_fails: int = 0, timeout: bool = False, comp_fails: bool = False) -> dict: reg = { "s1": {"fn": make_action("s1", log), "compensation": "c1"}, "c1": {"fn": make_action("c1", log)}, "s2": {"fn": make_action("s2", log, fails=s2_fails), "compensation": "c2"}, "c2": {"fn": make_action("c2", log)}, "s3": {"fn": make_action("s3", log, fails=s3_fails, timeout=timeout), "compensation": "c3"}, "c3": {"fn": make_action("c3", log)}, } if comp_fails: def broken_c1(_ctx): log.append("c1") raise saga.SagaCompensationError("simulated compensation failure") reg["c1"]["fn"] = broken_c1 return reg def _demo_steps() -> list[saga.SagaStep]: return [ saga.SagaStep("a", "s1", idem_key="k1", compensation_action="c1"), saga.SagaStep("b", "s2", idem_key="k2", max_retries=2, compensation_action="c2"), saga.SagaStep("c", "s3", idem_key="k3", compensation_action="c3"), ] # ------------------------------------------------------------ # ① 创建:记录结构 + 随世界状态持久化 # ------------------------------------------------------------ def test_saga_record_created_with_steps_and_persisted(): store = _MemStore() coord = saga.SagaCoordinator(store, actor="test") rec = coord.create_saga("demo", _demo_steps(), context={"x": 1}) assert rec["status"] == "PENDING" assert rec["steps"][0]["status"] == "PENDING" assert rec["steps"][0]["idemKey"] == "k1" assert rec["steps"][0]["timeoutSec"] == 30.0 assert rec["steps"][0]["maxRetries"] == 2 assert rec["steps"][0]["compensationAction"] == "c1" assert rec["auditRefs"] == [f"saga:{rec['id']}"] # 随世界状态持久化 assert store.data["sagas"][0]["id"] == rec["id"] assert any(e["action"] == "saga.created" for e in store.data["auditEvents"]) def test_saga_dedupe_key_reuses_existing_record(): store = _MemStore() coord = saga.SagaCoordinator(store, actor="test") r1 = coord.create_saga("demo", _demo_steps(), dedupe_key="EV-1") r2 = coord.create_saga("demo", _demo_steps(), dedupe_key="EV-1") assert r2["id"] == r1["id"] assert len(store.data["sagas"]) == 1 # ------------------------------------------------------------ # ② 全链成功 + 审计(含 idemKey/status/compensation) # ------------------------------------------------------------ def test_saga_runs_all_steps_to_succeeded_with_audit(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log), actor="test") rec = coord.create_saga("demo", _demo_steps()) out = coord.run(rec["id"]) assert out["status"] == "SUCCEEDED" assert [s["status"] for s in out["steps"]] == ["SUCCEEDED"] * 3 assert log == ["s1", "s2", "s3"] # 幂等键登记 assert set(store.data["sagaIdem"]) == {"k1", "k2", "k3"} actions = [e["action"] for e in store.data["auditEvents"]] assert actions.count("saga.step.succeeded") == 3 assert "saga.succeeded" in actions # 每次 step 迁移审计含 idemKey/status/compensation step_audit = next(e for e in store.data["auditEvents"] if e["action"] == "saga.step.succeeded") assert step_audit["rationale"]["idemKey"] assert step_audit["rationale"]["status"] == "SUCCEEDED" assert "compensation" in step_audit["rationale"] assert f"saga.step:{rec['id']}:a" in step_audit["evidenceRefs"] # ------------------------------------------------------------ # ③ 重试:瞬时失败重试成功后 saga 完成 # ------------------------------------------------------------ def test_saga_step_retry_after_transient_failure(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log, s2_fails=2), actor="test") rec = coord.create_saga("demo", _demo_steps()) out = coord.run(rec["id"]) assert out["status"] == "SUCCEEDED" b = next(s for s in out["steps"] if s["name"] == "b") assert b["attempts"] == 3 # 1 次初始 + 2 次重试 assert b["error"] is None assert "saga.step.retry" in [e["action"] for e in store.data["auditEvents"]] assert "saga.step.timeout" not in [e["action"] for e in store.data["auditEvents"]] # ------------------------------------------------------------ # ④ 失败超限 → 自动补偿链(逆序) # ------------------------------------------------------------ def test_saga_exhausts_retries_then_compensates_in_reverse(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log, s3_fails=99), actor="test") rec = coord.create_saga("demo", _demo_steps()) out = coord.run(rec["id"]) assert out["status"] == "COMPENSATED" assert out["compensationReason"].startswith("step:c") c = next(s for s in out["steps"] if s["name"] == "c") assert c["status"] == "FAILED" assert c["compensation"]["status"] == "SUCCEEDED" # 逆序补偿:c3(失败步自身的部分效果)→ c2 → c1(s3 重试预算=2 → 3 次尝试) assert log == ["s1", "s2", "s3", "s3", "s3", "c3", "c2", "c1"] actions = [e["action"] for e in store.data["auditEvents"]] assert "saga.compensating" in actions and "saga.compensated" in actions assert "saga.step.compensated" in actions # ------------------------------------------------------------ # ⑤ 超时:动作主动超时 / 真实时钟超时 → 重试 → 补偿 # ------------------------------------------------------------ def test_saga_timeout_triggers_retry_then_compensation(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log, timeout=True), actor="test") rec = coord.create_saga("demo", _demo_steps()) out = coord.run(rec["id"]) assert out["status"] == "COMPENSATED" c = next(s for s in out["steps"] if s["name"] == "c") assert c["timeouts"] == 3 # 每次尝试都超时 assert c["error"].startswith("timeout:") assert "saga.step.timeout" in [e["action"] for e in store.data["auditEvents"]] class _FakeClock: def __init__(self) -> None: self.t = 0.0 def __call__(self) -> float: return self.t def test_saga_clock_timeout_detection(): store = _MemStore() log: list[str] = [] clock = _FakeClock() reg = { "s1": {"fn": make_action("s1", log), "compensation": None}, "slow": {"fn": _slow_action(log, clock), "compensation": "slow-comp"}, "slow-comp": {"fn": make_action("slow-comp", log)}, } coord = saga.SagaCoordinator(store, registry=reg, clock=clock, actor="test") rec = coord.create_saga("demo", [ saga.SagaStep("a", "s1", idem_key="k1"), saga.SagaStep("b", "slow", idem_key="k2", timeout_sec=1.0, max_retries=1, compensation_action="slow-comp"), ]) out = coord.run(rec["id"]) assert out["status"] == "COMPENSATED" b = next(s for s in out["steps"] if s["name"] == "b") assert b["timeouts"] == 2 assert b["error"].startswith("timeout:") def _slow_action(log: list[str], clock: _FakeClock): def fn(_ctx) -> dict: log.append("slow") clock.t += 100.0 # 模拟耗时超过 timeout_sec return {"done": "slow"} return fn # ------------------------------------------------------------ # ⑥ 补偿失败 → MANUAL_TAKEOVER(人工接管点) # ------------------------------------------------------------ def test_saga_compensation_failure_leads_to_manual_takeover(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log, s2_fails=99, comp_fails=True), actor="test") rec = coord.create_saga("demo", [ saga.SagaStep("a", "s1", idem_key="k1", compensation_action="c1"), saga.SagaStep("b", "s2", idem_key="k2", max_retries=1), # 无补偿定义 ]) out = coord.run(rec["id"]) assert out["status"] == "MANUAL_TAKEOVER" assert out["manualTakeover"] is True assert "compensation-failed" in out["takeoverReason"] a = next(s for s in out["steps"] if s["name"] == "a") assert a["compensation"]["status"] == "FAILED" actions = [e["action"] for e in store.data["auditEvents"]] assert "saga.manual_takeover" in actions assert "saga.step.compensation_failed" in actions def test_saga_takeover_blocks_auto_run_then_retry_resumes(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log), actor="test") rec = coord.create_saga("demo", _demo_steps()) out = coord.takeover(rec["id"], reason="operator pause") assert out["status"] == "MANUAL_TAKEOVER" assert out["manualTakeover"] is True assert out["takeoverReason"] == "operator pause" # 接管后自动 run 被阻断 out2 = coord.run(rec["id"]) assert out2["status"] == "MANUAL_TAKEOVER" assert log == [] # 没有任何步骤执行 # 人工接管后显式 retry → 恢复自动化 out3 = coord.retry(rec["id"]) assert out3["status"] == "SUCCEEDED" assert out3["manualTakeover"] is False assert log == ["s1", "s2", "s3"] assert "saga.retry" in [e["action"] for e in store.data["auditEvents"]] # ------------------------------------------------------------ # ⑦ 恢复队列:中断的 saga 可 resume/replay # ------------------------------------------------------------ def test_saga_persistence_resume_after_restart(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log), actor="test") rec = coord.create_saga("demo", _demo_steps()) # 模拟进程中断:步骤 b 执行中崩溃(status=RUNNING,无 finishedAt) rec["status"] = "RUNNING" rec["steps"][0]["status"] = "SUCCEEDED" rec["steps"][0]["result"] = {"done": "s1"} rec["steps"][1]["status"] = "RUNNING" rec["steps"][1]["startedAt"] = "2026-08-02 10:00:00" store.save() # 重启:世界数据(含 sagas/sagaIdem)整体重建 store2 = _MemStore(copy.deepcopy(store.data)) log2: list[str] = [] coord2 = saga.SagaCoordinator(store2, registry=make_registry(log2), actor="test") out = coord2.resume_all() assert len(out) == 1 assert out[0]["status"] == "SUCCEEDED" assert log2 == ["s2", "s3"] # 已成功步骤不重放 assert all(s["status"] == "SUCCEEDED" for s in out[0]["steps"]) # ------------------------------------------------------------ # ⑧ 幂等重放:同 idem_key 已成功 → 复用记录,不重复执行外部效果 # ------------------------------------------------------------ def test_saga_idempotent_replay_reuses_recorded_result(): store = _MemStore() store.data.setdefault("sagaIdem", {})["k1"] = { "sagaId": "saga-other", "step": "a", "result": {"replayed": True}, "at": "2026-08-02 10:00:00"} log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log), actor="test") rec = coord.create_saga("demo", _demo_steps()) out = coord.run(rec["id"]) assert out["status"] == "SUCCEEDED" a = next(s for s in out["steps"] if s["name"] == "a") assert a["result"]["replayed"] is True assert log == ["s2", "s3"] # s1 未重新执行 assert "saga.step.idempotent_replayed" in [e["action"] for e in store.data["auditEvents"]] # ------------------------------------------------------------ # ⑨ 中间态可见:gate 步骤停在 WAITING_HUMAN,恢复后继续 # ------------------------------------------------------------ def test_saga_gate_pauses_waits_human_and_resumes(): store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log), actor="test") rec = coord.create_saga("demo", [ saga.SagaStep("a", "s1", idem_key="k1", gate=True), saga.SagaStep("b", "s2", idem_key="k2"), ]) out = coord.run(rec["id"], auto_approve=False) assert out["status"] == "WAITING_HUMAN" assert out["steps"][0]["status"] == "WAITING_HUMAN" assert out["steps"][1]["status"] == "PENDING" assert "saga.step.waiting_human" in [e["action"] for e in store.data["auditEvents"]] # 中间态可见:详情查询可看到门禁点 detail = coord.get(rec["id"]) assert detail["status"] == "WAITING_HUMAN" # 人工放行后继续 out2 = coord.run(rec["id"], auto_approve=True) assert out2["status"] == "SUCCEEDED" assert log == ["s1", "s2"] # ------------------------------------------------------------ # ⑩ 网关 API:GET /api/saga、详情、retry/compensate/takeover # ------------------------------------------------------------ def test_saga_gateway_api_endpoints(monkeypatch): import server.aps_domain.saga as saga_module import server.gateway.app as gateway_module from tests.auth_provider import install_test_auth store = _MemStore() log: list[str] = [] coord = saga.SagaCoordinator(store, registry=make_registry(log), actor="web") monkeypatch.setattr(gateway_module, "get_store", lambda: store) monkeypatch.setattr(saga_module, "get_coordinator", lambda _store, **kw: coord) install_test_auth(monkeypatch, "tenant-saga-test") rec = coord.create_saga("api-demo", _demo_steps()) client = TestClient(gateway_module.create_app()) login = client.post("/api/auth/login", json={"username": "planner"}) assert login.status_code == 200 # 人工接管 r = client.post(f"/api/saga/{rec['id']}/takeover", json={"reason": "ops"}) assert r.status_code == 200 assert r.json()["status"] == "MANUAL_TAKEOVER" assert r.json()["manualTakeover"] is True # 重试恢复自动化 r = client.post(f"/api/saga/{rec['id']}/retry", json={}) assert r.status_code == 200 assert r.json()["status"] == "SUCCEEDED" # 列表(含中间态 steps 明细) r = client.get("/api/saga") assert r.status_code == 200 assert r.json()["count"] >= 1 assert any(s["id"] == rec["id"] for s in r.json()["sagas"]) # 详情 r = client.get(f"/api/saga/{rec['id']}") assert r.status_code == 200 assert r.json()["id"] == rec["id"] assert [s["name"] for s in r.json()["steps"]] == ["a", "b", "c"] # 手动补偿 r = client.post(f"/api/saga/{rec['id']}/compensate", json={"reason": "manual"}) assert r.status_code == 200 assert r.json()["status"] == "COMPENSATED" # 不存在 → error 语义 r = client.get("/api/saga/saga-nope") assert r.status_code == 200 assert "error" in r.json()