feat(core): dispatch durable dataflow runs
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
|
||||
@@ -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__":
|
||||
|
||||
91
tests/test_dataflow_run_worker.py
Normal file
91
tests/test_dataflow_run_worker.py
Normal file
@@ -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()
|
||||
Reference in New Issue
Block a user