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"