diff --git a/src/govoplan_core/celery_app.py b/src/govoplan_core/celery_app.py index 227bfd1..d373b84 100644 --- a/src/govoplan_core/celery_app.py +++ b/src/govoplan_core/celery_app.py @@ -7,7 +7,9 @@ from celery import Celery from govoplan_core.core.campaigns import CAPABILITY_CAMPAIGNS_DELIVERY_TASKS, CampaignDeliveryTaskProvider from govoplan_core.core.calendar import CAPABILITY_CALENDAR_OUTBOX, CalendarOutboxProvider from govoplan_core.core.dataflows import ( + CAPABILITY_DATAFLOW_RUN_WORKER, CAPABILITY_DATAFLOW_TRIGGER_DISPATCHER, + DataflowRunWorker, DataflowTriggerDispatcher, ) from govoplan_core.core.events import ( @@ -42,6 +44,8 @@ celery.conf.update( "govoplan.notifications.deliver": {"queue": "notifications"}, "govoplan.notifications.deliver_pending": {"queue": "notifications"}, "govoplan.calendar.dispatch_outbox": {"queue": "calendar"}, + "govoplan.dataflow.dispatch_runs": {"queue": "dataflow"}, + "govoplan.dataflow.purge_runs": {"queue": "dataflow"}, "govoplan.dataflow.dispatch_triggers": {"queue": "dataflow"}, "govoplan.events.dispatch_outbox": {"queue": "events"}, "govoplan.events.purge_outbox": {"queue": "events"}, @@ -60,6 +64,16 @@ celery.conf.update( "schedule": 60.0, "args": (100,), }, + "dataflow-runs-every-five-seconds": { + "task": "govoplan.dataflow.dispatch_runs", + "schedule": 5.0, + "args": (10,), + }, + "dataflow-run-retention-daily": { + "task": "govoplan.dataflow.purge_runs", + "schedule": 24 * 60 * 60.0, + "args": (500,), + }, "platform-events-every-ten-seconds": { "task": "govoplan.events.dispatch_outbox", "schedule": 10.0, @@ -131,6 +145,18 @@ def _dataflow_trigger_dispatcher( return capability +def _dataflow_run_worker( + registry: PlatformRegistry | None = None, +) -> DataflowRunWorker | None: + registry = registry or _platform_registry() + if not registry.has_capability(CAPABILITY_DATAFLOW_RUN_WORKER): + return None + capability = registry.require_capability(CAPABILITY_DATAFLOW_RUN_WORKER) + if not isinstance(capability, DataflowRunWorker): + raise RuntimeError("Dataflow run worker capability is invalid") + return capability + + def _platform_event_outbox( registry: PlatformRegistry | None = None, ) -> PlatformEventOutbox | None: @@ -234,6 +260,56 @@ def dispatch_dataflow_triggers(self, limit: int = 100): return result +@celery.task( + name="govoplan.dataflow.dispatch_runs", + bind=True, + max_retries=0, +) +def dispatch_dataflow_runs(self, limit: int = 10): + """Claim and execute durable Dataflow runs outside the API process.""" + + from govoplan_core.db.session import get_database + + with get_database().SessionLocal() as session: + provider = _dataflow_run_worker() + if provider is None: + return { + "claimed": 0, + "succeeded": 0, + "retrying": 0, + "failed": 0, + "cancelled": 0, + } + result = dict( + provider.dispatch_pending( + session, + limit=limit, + worker_id=getattr(self.request, "hostname", None), + ) + ) + session.commit() + return result + + +@celery.task( + name="govoplan.dataflow.purge_runs", + bind=True, + max_retries=0, +) +def purge_dataflow_runs(self, limit: int = 500): + """Apply Dataflow evidence-retention policy without deleting run records.""" + + from govoplan_core.db.session import get_database + + with get_database().SessionLocal() as session: + provider = _dataflow_run_worker() + if provider is None: + return {"purged": 0} + result = dict(provider.purge_expired(session, limit=limit)) + session.commit() + return result + + @celery.task( name="govoplan.events.dispatch_outbox", bind=True, diff --git a/src/govoplan_core/core/dataflows.py b/src/govoplan_core/core/dataflows.py index 2997dec..3619d76 100644 --- a/src/govoplan_core/core/dataflows.py +++ b/src/govoplan_core/core/dataflows.py @@ -10,6 +10,7 @@ from govoplan_core.core.events import PlatformEvent CAPABILITY_DATAFLOW_RUN_LIFECYCLE = "dataflow.runLifecycle" +CAPABILITY_DATAFLOW_RUN_WORKER = "dataflow.runWorker" CAPABILITY_DATAFLOW_TRIGGER_DISPATCHER = "dataflow.triggerDispatcher" @@ -47,6 +48,10 @@ class DataflowRunRequest: revision: int idempotency_key: str row_limit: int = 500 + execution_backend: str = "auto" + environment: str = "development" + max_attempts: int = 3 + retention_days: int = 30 publication: DataflowPublicationTarget | None = None invocation: AutomationInvocation = field( default_factory=AutomationInvocation @@ -126,6 +131,34 @@ class DataflowTriggerDispatcher(Protocol): ... +@runtime_checkable +class DataflowRunWorker(Protocol): + """Durable worker boundary for queued Dataflow execution. + + Implementations own claim transaction boundaries so a lease is committed + before potentially long-running execution starts. + """ + + def dispatch_pending( + self, + session: object, + *, + now: datetime | None = None, + limit: int = 10, + worker_id: str | None = None, + ) -> Mapping[str, object]: + ... + + def purge_expired( + self, + session: object, + *, + now: datetime | None = None, + limit: int = 500, + ) -> Mapping[str, object]: + ... + + def dataflow_run_lifecycle( registry: object | None, ) -> DataflowRunLifecycleProvider | None: @@ -133,6 +166,13 @@ def dataflow_run_lifecycle( return capability if isinstance(capability, DataflowRunLifecycleProvider) else None +def dataflow_run_worker( + registry: object | None, +) -> DataflowRunWorker | None: + capability = _capability(registry, CAPABILITY_DATAFLOW_RUN_WORKER) + return capability if isinstance(capability, DataflowRunWorker) else None + + def dataflow_trigger_dispatcher( registry: object | None, ) -> DataflowTriggerDispatcher | None: @@ -157,6 +197,7 @@ def _capability(registry: object | None, name: str) -> object | None: __all__ = [ "CAPABILITY_DATAFLOW_RUN_LIFECYCLE", + "CAPABILITY_DATAFLOW_RUN_WORKER", "CAPABILITY_DATAFLOW_TRIGGER_DISPATCHER", "DataflowPublicationTarget", "DataflowRunConflictError", @@ -166,7 +207,9 @@ __all__ = [ "DataflowRunNotFoundError", "DataflowRunRequest", "DataflowRunUnavailableError", + "DataflowRunWorker", "DataflowTriggerDispatcher", "dataflow_run_lifecycle", + "dataflow_run_worker", "dataflow_trigger_dispatcher", ] diff --git a/tests/test_dataflow_contract.py b/tests/test_dataflow_contract.py index d63d4ff..17f6c61 100644 --- a/tests/test_dataflow_contract.py +++ b/tests/test_dataflow_contract.py @@ -4,9 +4,12 @@ import unittest from govoplan_core.core.dataflows import ( CAPABILITY_DATAFLOW_RUN_LIFECYCLE, + CAPABILITY_DATAFLOW_RUN_WORKER, DataflowRunDescriptor, DataflowRunLifecycleProvider, + DataflowRunWorker, dataflow_run_lifecycle, + dataflow_run_worker, ) from govoplan_core.core.modules import ModuleContext, ModuleManifest from govoplan_core.core.registry import PlatformRegistry @@ -49,10 +52,27 @@ class _Provider: ) +class _Worker: + def dispatch_pending( + self, + session, + *, + now=None, + limit=10, + worker_id=None, + ): + return {"claimed": 0} + + def purge_expired(self, session, *, now=None, limit=500): + return {"purged": 0} + + class DataflowContractTests(unittest.TestCase): def test_run_lifecycle_is_runtime_checkable_and_resolved(self) -> None: provider = _Provider() + worker = _Worker() self.assertIsInstance(provider, DataflowRunLifecycleProvider) + self.assertIsInstance(worker, DataflowRunWorker) registry = PlatformRegistry() registry.register( ModuleManifest( @@ -61,6 +81,7 @@ class DataflowContractTests(unittest.TestCase): version="test", capability_factories={ CAPABILITY_DATAFLOW_RUN_LIFECYCLE: lambda context: provider, + CAPABILITY_DATAFLOW_RUN_WORKER: lambda context: worker, }, ) ) @@ -69,7 +90,9 @@ class DataflowContractTests(unittest.TestCase): ) self.assertIs(provider, dataflow_run_lifecycle(registry)) + self.assertIs(worker, dataflow_run_worker(registry)) self.assertIsNone(dataflow_run_lifecycle(PlatformRegistry())) + self.assertIsNone(dataflow_run_worker(PlatformRegistry())) if __name__ == "__main__": diff --git a/tests/test_dataflow_run_worker.py b/tests/test_dataflow_run_worker.py new file mode 100644 index 0000000..328c4a6 --- /dev/null +++ b/tests/test_dataflow_run_worker.py @@ -0,0 +1,91 @@ +from __future__ import annotations + +import unittest +from unittest.mock import ANY, MagicMock, patch + +from govoplan_core.celery_app import ( + celery, + dispatch_dataflow_runs, + purge_dataflow_runs, +) + + +class DataflowRunWorkerTests(unittest.TestCase): + def test_dispatch_commits_worker_outcome(self) -> None: + session = MagicMock() + database = MagicMock() + database.SessionLocal.return_value.__enter__.return_value = session + provider = MagicMock() + provider.dispatch_pending.return_value = { + "claimed": 1, + "succeeded": 1, + } + + with ( + patch( + "govoplan_core.celery_app._dataflow_run_worker", + return_value=provider, + ), + patch( + "govoplan_core.db.session.get_database", + return_value=database, + ), + ): + result = dispatch_dataflow_runs.run(7) + + provider.dispatch_pending.assert_called_once_with( + session, + limit=7, + worker_id=ANY, + ) + session.commit.assert_called_once_with() + self.assertEqual(1, result["succeeded"]) + + def test_retention_commits_worker_outcome(self) -> None: + session = MagicMock() + database = MagicMock() + database.SessionLocal.return_value.__enter__.return_value = session + provider = MagicMock() + provider.purge_expired.return_value = {"purged": 2} + + with ( + patch( + "govoplan_core.celery_app._dataflow_run_worker", + return_value=provider, + ), + patch( + "govoplan_core.db.session.get_database", + return_value=database, + ), + ): + result = purge_dataflow_runs.run(25) + + provider.purge_expired.assert_called_once_with(session, limit=25) + session.commit.assert_called_once_with() + self.assertEqual(2, result["purged"]) + + def test_routes_and_periodic_jobs_are_registered(self) -> None: + self.assertEqual( + celery.conf.task_routes["govoplan.dataflow.dispatch_runs"], + {"queue": "dataflow"}, + ) + self.assertEqual( + celery.conf.task_routes["govoplan.dataflow.purge_runs"], + {"queue": "dataflow"}, + ) + self.assertEqual( + "govoplan.dataflow.dispatch_runs", + celery.conf.beat_schedule[ + "dataflow-runs-every-five-seconds" + ]["task"], + ) + self.assertEqual( + "govoplan.dataflow.purge_runs", + celery.conf.beat_schedule[ + "dataflow-run-retention-daily" + ]["task"], + ) + + +if __name__ == "__main__": + unittest.main()