| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768 |
- from __future__ import annotations
- from supply_infra.pipeline.dag import PIPELINE_STEPS
- from supply_infra.pipeline.gates import evaluate_step_gate
- def test_pipeline_has_fifteen_strictly_ordered_steps() -> None:
- assert len(PIPELINE_STEPS) == 15
- assert PIPELINE_STEPS[0].key == "global_tree_sync"
- assert PIPELINE_STEPS[-1].key == "aigc_write_record"
- for previous, current in zip(PIPELINE_STEPS, PIPELINE_STEPS[1:]):
- assert current.dependencies == (previous.key,)
- assert current.critical is True
- classify = next(step for step in PIPELINE_STEPS if step.key == "demand_classify")
- assert classify.retryable is True
- assert classify.max_attempts == 2
- def test_aigc_gate_requires_durable_effect_record() -> None:
- failed = evaluate_step_gate(
- "aigc_write_record",
- {"success": True, "external_request_made": True},
- )
- assert failed.passed is False
- live_ok = evaluate_step_gate(
- "aigc_write_record",
- {
- "success": True,
- "effect_recorded": True,
- "outbox_ids": ["outbox-1"],
- "external_request_made": True,
- },
- )
- assert live_ok.passed is True
- no_candidates = evaluate_step_gate(
- "aigc_write_record",
- {
- "success": True,
- "effect_recorded": False,
- "external_request_made": False,
- "candidate_count": 0,
- },
- )
- assert no_candidates.passed is True
- def test_aigc_gate_rejects_failed_batches() -> None:
- decision = evaluate_step_gate(
- "aigc_write_record",
- {
- "success": True,
- "effect_recorded": True,
- "outbox_ids": ["outbox-1"],
- "external_request_made": True,
- "failed_batch_count": 2,
- },
- )
- assert decision.passed is False
- assert decision.error_code == "aigc_publish_failed"
- def test_explicit_step_failure_fails_closed() -> None:
- decision = evaluate_step_gate(
- "global_tree_sync",
- {"success": False, "error": "source missing"},
- )
- assert decision.passed is False
- assert decision.error_code == "step_reported_failure"
|