diff --git a/src/govoplan_core/celery_app.py b/src/govoplan_core/celery_app.py index 1c02e65..c8046f4 100644 --- a/src/govoplan_core/celery_app.py +++ b/src/govoplan_core/celery_app.py @@ -729,7 +729,7 @@ def dispatch_campaign_schedules( tenant_id: str | None = None, limit: int = 50, ): - """Prepare due campaign drafts; delivery always remains a separate action.""" + """Prepare manual drafts or governed autonomous Mail commands for due schedules.""" from govoplan_core.db.session import get_database @@ -738,11 +738,21 @@ def dispatch_campaign_schedules( defaults = { "selected": 0, "prepared": 0, + "autonomous_prepared": 0, "failed": 0, "completed": 0, "coalesced": 0, + "duplicates": 0, + "deferred": 0, "campaign_ids": [], "operator_actions": [], + "refreshed": { + "checked": 0, + "accepted": 0, + "uncertain": 0, + "failed": 0, + "skipped": 0, + }, } if not registry.has_capability(CAPABILITY_CAMPAIGNS_SCHEDULES): return defaults diff --git a/src/govoplan_core/core/campaigns.py b/src/govoplan_core/core/campaigns.py index f35bfe7..48d9e86 100644 --- a/src/govoplan_core/core/campaigns.py +++ b/src/govoplan_core/core/campaigns.py @@ -108,7 +108,7 @@ class CampaignDeliveryTaskProvider(Protocol): @runtime_checkable class CampaignScheduleProvider(Protocol): - """Durable boundary for preparing due recurring Campaign drafts.""" + """Durable boundary for due manual drafts and governed autonomous occurrences.""" def dispatch_due( self, diff --git a/tests/test_campaign_schedule_worker.py b/tests/test_campaign_schedule_worker.py new file mode 100644 index 0000000..13cd77f --- /dev/null +++ b/tests/test_campaign_schedule_worker.py @@ -0,0 +1,106 @@ +from __future__ import annotations + +import unittest +from contextlib import contextmanager +from types import SimpleNamespace +from unittest.mock import MagicMock, patch + +from govoplan_core.celery_app import celery, dispatch_campaign_schedules +from tests.worker_test_support import allowed_worker_admissions + + +class CampaignScheduleWorkerTests(unittest.TestCase): + def test_dispatch_preserves_autonomous_recovery_result(self) -> None: + session = MagicMock() + + @contextmanager + def session_scope(): + yield session + + provider = MagicMock() + provider.dispatch_due.return_value = { + "selected": 1, + "prepared": 0, + "autonomous_prepared": 1, + "failed": 0, + "completed": 0, + "coalesced": 0, + "duplicates": 0, + "deferred": 0, + "campaign_ids": ["campaign-1"], + "operator_actions": [], + "refreshed": { + "checked": 1, + "accepted": 1, + "uncertain": 0, + "failed": 0, + "skipped": 0, + }, + } + registry = MagicMock() + registry.has_capability.return_value = True + database = SimpleNamespace(SessionLocal=session_scope) + + with ( + patch("govoplan_core.celery_app._platform_registry", return_value=registry), + patch("govoplan_core.celery_app._campaign_schedules", return_value=provider), + patch( + "govoplan_core.celery_app._worker_admissions", + side_effect=allowed_worker_admissions, + ), + patch("govoplan_core.db.session.get_database", return_value=database), + ): + result = dispatch_campaign_schedules.run("tenant-1", 17) + + provider.dispatch_due.assert_called_once_with( + session, + tenant_id="tenant-1", + limit=17, + ) + session.commit.assert_called_once_with() + self.assertEqual(result["autonomous_prepared"], 1) + self.assertEqual(result["refreshed"]["accepted"], 1) + + def test_missing_provider_returns_complete_autonomous_defaults(self) -> None: + session = MagicMock() + + @contextmanager + def session_scope(): + yield session + + registry = MagicMock() + registry.has_capability.return_value = False + database = SimpleNamespace(SessionLocal=session_scope) + + with ( + patch("govoplan_core.celery_app._platform_registry", return_value=registry), + patch("govoplan_core.db.session.get_database", return_value=database), + ): + result = dispatch_campaign_schedules.run("tenant-1", 17) + + self.assertEqual(result["autonomous_prepared"], 0) + self.assertEqual(result["duplicates"], 0) + self.assertEqual(result["deferred"], 0) + self.assertEqual( + result["refreshed"], + { + "checked": 0, + "accepted": 0, + "uncertain": 0, + "failed": 0, + "skipped": 0, + }, + ) + + def test_worker_route_and_periodic_dispatch_are_registered(self) -> None: + self.assertEqual( + celery.conf.task_routes["govoplan.campaigns.dispatch_schedules"], + {"queue": "default"}, + ) + schedule = celery.conf.beat_schedule["campaign-schedules-every-minute"] + self.assertEqual(schedule["task"], "govoplan.campaigns.dispatch_schedules") + self.assertEqual(schedule["schedule"], 60.0) + + +if __name__ == "__main__": + unittest.main()