from __future__ import annotations import unittest from unittest.mock import MagicMock, patch from govoplan_core.celery_app import celery, dispatch_dataflow_triggers class DataflowTriggerWorkerTests(unittest.TestCase): def test_worker_commits_dispatch_outcome(self) -> None: session = MagicMock() database = MagicMock() database.SessionLocal.return_value.__enter__.return_value = session provider = MagicMock() provider.dispatch_due.return_value = { "queued": 1, "processed": 1, "succeeded": 1, "failed": 0, "blocked": 0, "skipped": 0, } with ( patch( "govoplan_core.celery_app._dataflow_trigger_dispatcher", return_value=provider, ), patch( "govoplan_core.db.session.get_database", return_value=database, ), ): result = dispatch_dataflow_triggers.run(25) provider.dispatch_due.assert_called_once_with(session, limit=25) session.commit.assert_called_once_with() self.assertEqual(result["succeeded"], 1) def test_worker_route_and_periodic_dispatch_are_registered(self) -> None: self.assertEqual( celery.conf.task_routes["govoplan.dataflow.dispatch_triggers"], {"queue": "dataflow"}, ) schedule = celery.conf.beat_schedule[ "dataflow-triggers-every-minute" ] self.assertEqual( schedule["task"], "govoplan.dataflow.dispatch_triggers", ) self.assertEqual(schedule["schedule"], 60.0) if __name__ == "__main__": unittest.main()