test_dag_and_gates.py 2.1 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768
  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_fifteen_strictly_ordered_steps() -> None:
  5. assert len(PIPELINE_STEPS) == 15
  6. assert PIPELINE_STEPS[0].key == "global_tree_sync"
  7. assert PIPELINE_STEPS[-1].key == "aigc_write_record"
  8. for previous, current in zip(PIPELINE_STEPS, PIPELINE_STEPS[1:]):
  9. assert current.dependencies == (previous.key,)
  10. assert current.critical is True
  11. classify = next(step for step in PIPELINE_STEPS if step.key == "demand_classify")
  12. assert classify.retryable is True
  13. assert classify.max_attempts == 2
  14. def test_aigc_gate_requires_durable_effect_record() -> None:
  15. failed = evaluate_step_gate(
  16. "aigc_write_record",
  17. {"success": True, "external_request_made": True},
  18. )
  19. assert failed.passed is False
  20. live_ok = evaluate_step_gate(
  21. "aigc_write_record",
  22. {
  23. "success": True,
  24. "effect_recorded": True,
  25. "outbox_ids": ["outbox-1"],
  26. "external_request_made": True,
  27. },
  28. )
  29. assert live_ok.passed is True
  30. no_candidates = evaluate_step_gate(
  31. "aigc_write_record",
  32. {
  33. "success": True,
  34. "effect_recorded": False,
  35. "external_request_made": False,
  36. "candidate_count": 0,
  37. },
  38. )
  39. assert no_candidates.passed is True
  40. def test_aigc_gate_rejects_failed_batches() -> None:
  41. decision = evaluate_step_gate(
  42. "aigc_write_record",
  43. {
  44. "success": True,
  45. "effect_recorded": True,
  46. "outbox_ids": ["outbox-1"],
  47. "external_request_made": True,
  48. "failed_batch_count": 2,
  49. },
  50. )
  51. assert decision.passed is False
  52. assert decision.error_code == "aigc_publish_failed"
  53. def test_explicit_step_failure_fails_closed() -> None:
  54. decision = evaluate_step_gate(
  55. "global_tree_sync",
  56. {"success": False, "error": "source missing"},
  57. )
  58. assert decision.passed is False
  59. assert decision.error_code == "step_reported_failure"