test_dag_and_gates.py 1.1 KB

1234567891011121314151617181920212223242526272829303132
  1. from __future__ import annotations
  2. from supply_infra.pipeline.dag import PIPELINE_STEPS
  3. from supply_infra.pipeline.gates import evaluate_step_gate
  4. def test_pipeline_has_eleven_strictly_ordered_steps() -> None:
  5. assert len(PIPELINE_STEPS) == 11
  6. assert PIPELINE_STEPS[0].key == "global_tree_sync"
  7. assert PIPELINE_STEPS[-1].key == "video_discovery"
  8. assert "aigc_write_record" not in {step.key for step in PIPELINE_STEPS}
  9. for previous, current in zip(PIPELINE_STEPS, PIPELINE_STEPS[1:]):
  10. assert current.dependencies == (previous.key,)
  11. assert current.critical is True
  12. def test_explicit_step_failure_fails_closed() -> None:
  13. decision = evaluate_step_gate(
  14. "global_tree_sync",
  15. {"success": False, "error": "source missing"},
  16. )
  17. assert decision.passed is False
  18. assert decision.error_code == "step_reported_failure"
  19. def test_video_discovery_gate_rejects_business_failures() -> None:
  20. decision = evaluate_step_gate(
  21. "video_discovery",
  22. {"success": True, "failed": 2},
  23. )
  24. assert decision.passed is False
  25. assert decision.error_code == "business_failures"