"""Compatibility wrapper tests for the durable pipeline.""" from __future__ import annotations from unittest.mock import patch from supply_infra.pipeline.run_service import RunSubmission from supply_infra.scheduler.jobs.run_supply_pipeline import run_supply_pipeline @patch("supply_infra.scheduler.jobs.run_supply_pipeline.submit_pipeline_run") def test_pipeline_wrapper_only_submits_durable_run(mock_submit) -> None: mock_submit.return_value = RunSubmission( run_id="run-1", created=True, status="queued", biz_dt="20260721", ) result = run_supply_pipeline("20260721") assert result["run_id"] == "run-1" assert result["status"] == "queued" assert result["dry_run"] is False mock_submit.assert_called_once() @patch("supply_infra.scheduler.jobs.run_supply_pipeline.get_pipeline_run") @patch("supply_infra.scheduler.jobs.run_supply_pipeline.submit_pipeline_run") def test_wait_reads_terminal_state_from_mysql(mock_submit, mock_get) -> None: mock_submit.return_value = RunSubmission( run_id="run-2", created=True, status="queued", biz_dt="20260721", ) mock_get.return_value = { "run_id": "run-2", "status": "failed", "error_message": "strict gate failed", } result = run_supply_pipeline("20260721", wait=True, poll_seconds=0.01) assert result["success"] is False assert result["status"] == "failed" @patch("supply_infra.scheduler.jobs.run_supply_pipeline.submit_pipeline_run") def test_submit_validation_error_is_not_hidden(mock_submit) -> None: mock_submit.side_effect = ValueError("Invalid biz_dt") try: run_supply_pipeline("invalid-date") except ValueError as exc: assert "Invalid biz_dt" in str(exc) else: raise AssertionError("validation error must propagate to the CLI/API boundary")