"""Unit tests for src.app (run_renew + main daemon loop).""" from typing import Any from unittest.mock import AsyncMock, MagicMock, patch import pytest from src import app, config, queries def _make_cap_row(**overrides: Any) -> dict[str, Any]: row: dict[str, Any] = { "dms_master_master_id": 1, "jobs_to_queue": 5, } row.update(overrides) return row def _make_job_row(**overrides: Any) -> dict[str, Any]: row: dict[str, Any] = { "encoding_queue_detail_id": 1, "dms_master_master_id": 1, "upc": 12345, "cd": 1, "track_id": 1, "clip_number": None, "encoder_id": 18, "encoding_order_type": "release", "meta_update": "N", "encoding_order_priority": 1, # jp.priority "store_priority": 1, # cmm.priority "r": 1, } row.update(overrides) return row def _mock_renew_connection( mock_dd_conn: MagicMock, caps: list[dict[str, Any]], jobs: list[dict[str, Any]], ) -> tuple[MagicMock, MagicMock, MagicMock, MagicMock]: """Wire mock_dd_conn so first two cursors yield caps/jobs, rest share status.""" mock_caps_cursor = MagicMock() mock_caps_cursor.fetchall.return_value = caps mock_jobs_cursor = MagicMock() mock_jobs_cursor.fetchall.return_value = jobs mock_status_cursor = MagicMock() head = iter([mock_caps_cursor, mock_jobs_cursor]) def enter_side_effect(*args: Any, **kwargs: Any) -> MagicMock: return next(head, mock_status_cursor) mock_conn = MagicMock() mock_conn.cursor.return_value.__enter__.side_effect = enter_side_effect mock_conn.cursor.return_value.__exit__.return_value = False mock_dd_conn.return_value.__enter__.return_value = mock_conn mock_dd_conn.return_value.__exit__.return_value = False return mock_conn, mock_caps_cursor, mock_jobs_cursor, mock_status_cursor @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_no_caps_skips_everything( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """When step 1 returns no rows, skip step 2, SQS, status, and metrics.""" mock_conn, *_ = _mock_renew_connection(mock_dd_conn, caps=[], jobs=[]) app.run_renew(MagicMock()) mock_fan_out.assert_not_called() mock_conn.commit.assert_not_called() mock_dd.assert_not_called() @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_caps_but_no_jobs_skips_sqs( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """Caps exist but step 2 returns nothing → no SQS, no commit, but still emit per-DMS metrics.""" mock_conn, *_ = _mock_renew_connection( mock_dd_conn, caps=[_make_cap_row(dms_master_master_id=7, jobs_to_queue=3)], jobs=[], ) app.run_renew(MagicMock()) mock_fan_out.assert_not_called() mock_conn.commit.assert_not_called() metrics_sent = [c[0][0] for c in mock_dd.call_args_list] assert metrics_sent.count("renew.attempt") == 1 assert metrics_sent.count("renew.run") == 1 tag_sets = [c[0][1] for c in mock_dd.call_args_list] assert all("dms_id:7" in tags for tags in tag_sets) @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_caps_query_uses_max_jobs_per_run( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """SQL_GET_STORE_CAPS must be executed with config.MAX_JOBS_PER_RUN.""" _, mock_caps_cursor, *_ = _mock_renew_connection(mock_dd_conn, caps=[], jobs=[]) app.run_renew(MagicMock()) mock_caps_cursor.execute.assert_called_once_with( queries.SQL_GET_STORE_CAPS, (config.MAX_JOBS_PER_RUN,), ) @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.get_queue_name_for_delivery_job") @patch("src.app.util.dd_connection") def test_run_renew_promotes_jobs_and_fans_out( mock_dd_conn: MagicMock, mock_get_queue: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock, ) -> None: """run_renew should DB-promote then hand the messages dict to fan_out_to_sqs.""" queue_name = "test-queue" mock_get_queue.return_value = queue_name mock_fan_out.return_value = [] job = _make_job_row() mock_conn, _, _, mock_status_cursor = _mock_renew_connection( mock_dd_conn, caps=[_make_cap_row(jobs_to_queue=1)], jobs=[job] ) app.run_renew(MagicMock()) mock_fan_out.assert_awaited_once() assert mock_fan_out.await_args is not None messages_by_queue, promoted_ids, _ = mock_fan_out.await_args[0] assert list(messages_by_queue.keys()) == [queue_name] assert len(messages_by_queue[queue_name]) == 1 pushed = messages_by_queue[queue_name][0] assert pushed["encoding_queue_detail_id"] == job["encoding_queue_detail_id"] assert pushed["meta_update"] is False assert pushed["priority"] == job["store_priority"] assert promoted_ids == [job["encoding_queue_detail_id"]] mock_status_cursor.execute.assert_called_once_with( queries.SQL_SET_JOBS_STATUS, ("ready_to_encode", (job["encoding_queue_detail_id"],), "new"), ) mock_conn.commit.assert_called_once() @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_builds_case_when_from_caps( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """Step 2 SQL should embed a CASE WHEN per store and pass caps as params.""" caps = [ _make_cap_row(dms_master_master_id=1, jobs_to_queue=2), _make_cap_row(dms_master_master_id=2, jobs_to_queue=3), ] _, _, mock_jobs_cursor, _ = _mock_renew_connection(mock_dd_conn, caps=caps, jobs=[]) app.run_renew(MagicMock()) sql, params = mock_jobs_cursor.execute.call_args[0] assert sql.count("WHEN %s THEN %s") == 2 assert params[0] == (1, 2) assert params[1:] == (1, 2, 2, 3) @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.get_queue_name_for_delivery_job") @patch("src.app.util.dd_connection") def test_run_renew_groups_jobs_by_queue( mock_dd_conn: MagicMock, mock_get_queue: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock, ) -> None: """Jobs with different priorities are grouped under distinct queue names.""" mock_get_queue.side_effect = lambda **kw: f"queue-{kw['encoding_order_priority']}" mock_fan_out.return_value = [] jobs = [ _make_job_row(encoding_queue_detail_id=1, encoding_order_priority=1), _make_job_row(encoding_queue_detail_id=2, encoding_order_priority=2), _make_job_row(encoding_queue_detail_id=3, encoding_order_priority=1), ] _mock_renew_connection( mock_dd_conn, caps=[_make_cap_row(jobs_to_queue=3)], jobs=jobs ) app.run_renew(MagicMock()) assert mock_fan_out.await_args is not None messages_by_queue = mock_fan_out.await_args[0][0] assert sorted(messages_by_queue.keys()) == ["queue-1", "queue-2"] assert len(messages_by_queue["queue-1"]) == 2 assert len(messages_by_queue["queue-2"]) == 1 @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_meta_update_y_becomes_true( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """meta_update='Y' should be converted to True in the SQS message dict.""" mock_fan_out.return_value = [] _mock_renew_connection( mock_dd_conn, caps=[_make_cap_row(jobs_to_queue=1)], jobs=[_make_job_row(meta_update="Y")], ) app.run_renew(MagicMock()) assert mock_fan_out.await_args is not None messages_by_queue = mock_fan_out.await_args[0][0] pushed = next(iter(messages_by_queue.values()))[0] assert pushed["meta_update"] is True @patch("src.app.config.DB_BATCH_SIZE", 2) @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_db_update_chunked( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """Phase-1 promotion runs in DB_BATCH_SIZE chunks; each chunk commits.""" mock_fan_out.return_value = [] jobs = [ _make_job_row(encoding_queue_detail_id=1), _make_job_row(encoding_queue_detail_id=2, dms_master_master_id=2), _make_job_row(encoding_queue_detail_id=3, dms_master_master_id=3), ] mock_conn, _, _, mock_status_cursor = _mock_renew_connection( mock_dd_conn, caps=[ _make_cap_row(dms_master_master_id=1, jobs_to_queue=1), _make_cap_row(dms_master_master_id=2, jobs_to_queue=1), _make_cap_row(dms_master_master_id=3, jobs_to_queue=1), ], jobs=jobs, ) app.run_renew(MagicMock()) executed_args = [c[0] for c in mock_status_cursor.execute.call_args_list] assert [args[1][1] for args in executed_args] == [(1, 2), (3,)] for args in executed_args: assert args[0] == queries.SQL_SET_JOBS_STATUS assert args[1][0] == "ready_to_encode" assert args[1][2] == "new" assert mock_conn.commit.call_count == 2 @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_resets_failed_ids( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """fan_out_to_sqs failures should drive a phase-3 reset to 'new'.""" mock_fan_out.return_value = [2] jobs = [ _make_job_row(encoding_queue_detail_id=1), _make_job_row(encoding_queue_detail_id=2), ] mock_conn, _, _, mock_status_cursor = _mock_renew_connection( mock_dd_conn, caps=[_make_cap_row(jobs_to_queue=2)], jobs=jobs ) app.run_renew(MagicMock()) # Two execute calls: promote + reset executed = mock_status_cursor.execute.call_args_list assert len(executed) == 2 assert executed[0][0][1] == ("ready_to_encode", (1, 2), "new") assert executed[1][0][1] == ("new", (2,), "ready_to_encode") assert mock_conn.commit.call_count == 2 @patch("src.app.config.DB_BATCH_SIZE", 2) @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_reset_is_chunked( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """Phase-3 reset runs in DB_BATCH_SIZE chunks, each committed.""" mock_fan_out.return_value = [1, 2, 3] jobs = [ _make_job_row(encoding_queue_detail_id=1), _make_job_row(encoding_queue_detail_id=2, dms_master_master_id=2), _make_job_row(encoding_queue_detail_id=3, dms_master_master_id=3), ] mock_conn, _, _, mock_status_cursor = _mock_renew_connection( mock_dd_conn, caps=[ _make_cap_row(dms_master_master_id=1, jobs_to_queue=1), _make_cap_row(dms_master_master_id=2, jobs_to_queue=1), _make_cap_row(dms_master_master_id=3, jobs_to_queue=1), ], jobs=jobs, ) app.run_renew(MagicMock()) executed = [c[0] for c in mock_status_cursor.execute.call_args_list] reset_calls = [args for args in executed if args[1][0] == "new"] # Three failed ids reset in chunks of 2 → (1, 2) then (3,). assert [args[1][1] for args in reset_calls] == [(1, 2), (3,)] for args in reset_calls: assert args[0] == queries.SQL_SET_JOBS_STATUS assert args[1][2] == "ready_to_encode" # 2 promote commits + 2 reset commits. assert mock_conn.commit.call_count == 4 @patch("src.app.call_datadog_with_metric") @patch("src.app.fan_out_to_sqs", new_callable=AsyncMock) @patch("src.app.util.dd_connection") def test_run_renew_emits_attempt_and_run_per_capped_dms( mock_dd_conn: MagicMock, mock_fan_out: AsyncMock, mock_dd: MagicMock ) -> None: """Attempt + run fire once per DMS from step 1's caps, not per job row.""" mock_fan_out.return_value = [] jobs = [ _make_job_row(encoding_queue_detail_id=1, dms_master_master_id=10), _make_job_row(encoding_queue_detail_id=2, dms_master_master_id=10), ] _mock_renew_connection( mock_dd_conn, caps=[ _make_cap_row(dms_master_master_id=10, jobs_to_queue=2), _make_cap_row(dms_master_master_id=20, jobs_to_queue=1), ], jobs=jobs, ) app.run_renew(MagicMock()) metrics_sent = [c[0][0] for c in mock_dd.call_args_list] assert metrics_sent.count("renew.attempt") == 2 assert metrics_sent.count("renew.run") == 2 tag_sets = [c[0][1] for c in mock_dd.call_args_list] assert sum("dms_id:10" in tags for tags in tag_sets) == 2 assert sum("dms_id:20" in tags for tags in tag_sets) == 2 @patch("src.task_protection._set_task_protection", return_value=True) @patch("src.app.call_datadog_with_metric") @patch("src.app.time.sleep") @patch("src.app.time.monotonic", return_value=0) @patch("src.app.run_renew") def test_main_runs_configured_iterations_with_sleeps_between( mock_run_renew: MagicMock, mock_monotonic: MagicMock, mock_sleep: MagicMock, mock_datadog: MagicMock, mock_set_protection: MagicMock, ) -> None: """main() should run ITERATIONS times and emit renew_manager.run per iteration.""" app.main() assert mock_run_renew.call_count == config.ITERATIONS assert mock_sleep.call_count == config.ITERATIONS - 1 mock_sleep.assert_called_with(config.INTERVAL_SECONDS) assert mock_datadog.call_count == config.ITERATIONS for call in mock_datadog.call_args_list: assert call[0][0] == "renew_manager.run" assert call[0][1] == [f"environment:{config.ENVIRONMENT}"] @patch("src.task_protection._set_task_protection", return_value=True) @patch("src.app.call_datadog_with_metric") @patch("src.app.time.sleep") @patch( "src.app.time.monotonic", side_effect=[0, config.INTERVAL_SECONDS + 10] * config.ITERATIONS, ) @patch("src.app.run_renew") def test_main_clamps_sleep_to_zero_when_iteration_overruns( mock_run_renew: MagicMock, mock_monotonic: MagicMock, mock_sleep: MagicMock, mock_datadog: MagicMock, mock_set_protection: MagicMock, ) -> None: """Sleep should be 0 when processing takes longer than the interval.""" app.main() mock_sleep.assert_called_with(0) @patch("src.task_protection._set_task_protection", return_value=True) @patch("src.app.call_datadog_with_metric") @patch("src.app.time.sleep") @patch("src.app.time.monotonic", return_value=0) @patch("src.app.run_renew") def test_main_toggles_task_protection_per_iteration( mock_run_renew: MagicMock, mock_monotonic: MagicMock, mock_sleep: MagicMock, mock_datadog: MagicMock, mock_set_protection: MagicMock, ) -> None: """Each iteration should enable protection before run_renew and disable after.""" app.main() assert mock_set_protection.call_count == 2 * config.ITERATIONS toggles = [call[0][1] for call in mock_set_protection.call_args_list] assert toggles == [True, False] * config.ITERATIONS @patch("src.task_protection._set_task_protection", return_value=True) @patch("src.app.call_datadog_with_metric") @patch("src.app.time.sleep") @patch("src.app.time.monotonic", return_value=0) @patch("src.app.run_renew", side_effect=RuntimeError("boom")) def test_main_disables_task_protection_on_failure( mock_run_renew: MagicMock, mock_monotonic: MagicMock, mock_sleep: MagicMock, mock_datadog: MagicMock, mock_set_protection: MagicMock, ) -> None: """If run_renew raises, protection should still be disabled for that iteration.""" with pytest.raises(RuntimeError): app.main() assert mock_set_protection.call_count == 2 assert mock_set_protection.call_args_list[-1][0][1] is False @patch("src.task_protection._set_task_protection", return_value=False) @patch("src.app.call_datadog_with_metric") @patch("src.app.time.sleep") @patch("src.app.time.monotonic", return_value=0) @patch("src.app.run_renew") def test_main_skips_iteration_when_protection_fails( mock_run_renew: MagicMock, mock_monotonic: MagicMock, mock_sleep: MagicMock, mock_datadog: MagicMock, mock_set_protection: MagicMock, ) -> None: """If enabling protection fails, the iteration's work is skipped (no run_renew, no metric).""" mock_run_renew.__name__ = "run_renew" app.main() mock_run_renew.assert_not_called() mock_datadog.assert_not_called() assert mock_set_protection.call_count == config.ITERATIONS assert all(call[0][1] is True for call in mock_set_protection.call_args_list) assert mock_sleep.call_count == config.ITERATIONS - 1