from __future__ import annotations import unittest from unittest.mock import patch from govoplan_core.celery_app import celery, worker_acceptance_probe from govoplan_core.settings import settings class CeleryQueueContractTests(unittest.TestCase): def test_default_worker_queues_cover_every_routed_task(self) -> None: configured = { value.strip() for value in settings.celery_queues.split(",") if value.strip() } routes = celery.conf.task_routes self.assertIsInstance(routes, dict) routed = { str(route["queue"]) for route in routes.values() if isinstance(route, dict) and route.get("queue") } self.assertEqual(set(), routed - configured) def test_delivery_tasks_keep_worker_loss_protection(self) -> None: self.assertTrue(celery.conf.task_acks_late) self.assertTrue(celery.conf.task_reject_on_worker_lost) self.assertTrue(celery.conf.task_track_started) self.assertEqual(1, celery.conf.worker_prefetch_multiplier) self.assertEqual( settings.celery_visibility_timeout_seconds, celery.conf.broker_transport_options["visibility_timeout"], ) def test_worker_acceptance_probe_is_bounded_and_side_effect_free(self) -> None: with patch.object(worker_acceptance_probe, "update_state") as update_state: result = worker_acceptance_probe.run( "probe-1", mode="complete", delay_seconds=0, ) self.assertEqual( { "probe_id": "probe-1", "mode": "complete", "retries": 0, "redelivered": False, "delivery_count": None, }, result, ) update_state.assert_called_once() if __name__ == "__main__": unittest.main()