feat: support autonomous campaign schedule results
This commit is contained in:
@@ -729,7 +729,7 @@ def dispatch_campaign_schedules(
|
|||||||
tenant_id: str | None = None,
|
tenant_id: str | None = None,
|
||||||
limit: int = 50,
|
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
|
from govoplan_core.db.session import get_database
|
||||||
|
|
||||||
@@ -738,11 +738,21 @@ def dispatch_campaign_schedules(
|
|||||||
defaults = {
|
defaults = {
|
||||||
"selected": 0,
|
"selected": 0,
|
||||||
"prepared": 0,
|
"prepared": 0,
|
||||||
|
"autonomous_prepared": 0,
|
||||||
"failed": 0,
|
"failed": 0,
|
||||||
"completed": 0,
|
"completed": 0,
|
||||||
"coalesced": 0,
|
"coalesced": 0,
|
||||||
|
"duplicates": 0,
|
||||||
|
"deferred": 0,
|
||||||
"campaign_ids": [],
|
"campaign_ids": [],
|
||||||
"operator_actions": [],
|
"operator_actions": [],
|
||||||
|
"refreshed": {
|
||||||
|
"checked": 0,
|
||||||
|
"accepted": 0,
|
||||||
|
"uncertain": 0,
|
||||||
|
"failed": 0,
|
||||||
|
"skipped": 0,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
if not registry.has_capability(CAPABILITY_CAMPAIGNS_SCHEDULES):
|
if not registry.has_capability(CAPABILITY_CAMPAIGNS_SCHEDULES):
|
||||||
return defaults
|
return defaults
|
||||||
|
|||||||
@@ -108,7 +108,7 @@ class CampaignDeliveryTaskProvider(Protocol):
|
|||||||
|
|
||||||
@runtime_checkable
|
@runtime_checkable
|
||||||
class CampaignScheduleProvider(Protocol):
|
class CampaignScheduleProvider(Protocol):
|
||||||
"""Durable boundary for preparing due recurring Campaign drafts."""
|
"""Durable boundary for due manual drafts and governed autonomous occurrences."""
|
||||||
|
|
||||||
def dispatch_due(
|
def dispatch_due(
|
||||||
self,
|
self,
|
||||||
|
|||||||
@@ -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()
|
||||||
Reference in New Issue
Block a user