| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657 |
- """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")
|