feat(core): reconcile durable workflow instances
This commit is contained in:
@@ -22,6 +22,10 @@ from govoplan_core.core.events import (
|
|||||||
from govoplan_core.core.module_management import load_startup_enabled_modules, startup_candidate_module_ids
|
from govoplan_core.core.module_management import load_startup_enabled_modules, startup_candidate_module_ids
|
||||||
from govoplan_core.core.modules import ModuleContext
|
from govoplan_core.core.modules import ModuleContext
|
||||||
from govoplan_core.core.notifications import CAPABILITY_NOTIFICATIONS_DISPATCH, NotificationDispatchProvider
|
from govoplan_core.core.notifications import CAPABILITY_NOTIFICATIONS_DISPATCH, NotificationDispatchProvider
|
||||||
|
from govoplan_core.core.workflows import (
|
||||||
|
CAPABILITY_WORKFLOW_RUNTIME_WORKER,
|
||||||
|
WorkflowRuntimeWorker,
|
||||||
|
)
|
||||||
from govoplan_core.core.registry import PlatformRegistry
|
from govoplan_core.core.registry import PlatformRegistry
|
||||||
from govoplan_core.core.runtime import configure_runtime
|
from govoplan_core.core.runtime import configure_runtime
|
||||||
from govoplan_core.settings import settings
|
from govoplan_core.settings import settings
|
||||||
@@ -47,6 +51,7 @@ celery.conf.update(
|
|||||||
"govoplan.dataflow.dispatch_runs": {"queue": "dataflow"},
|
"govoplan.dataflow.dispatch_runs": {"queue": "dataflow"},
|
||||||
"govoplan.dataflow.purge_runs": {"queue": "dataflow"},
|
"govoplan.dataflow.purge_runs": {"queue": "dataflow"},
|
||||||
"govoplan.dataflow.dispatch_triggers": {"queue": "dataflow"},
|
"govoplan.dataflow.dispatch_triggers": {"queue": "dataflow"},
|
||||||
|
"govoplan.workflow.reconcile": {"queue": "workflow"},
|
||||||
"govoplan.events.dispatch_outbox": {"queue": "events"},
|
"govoplan.events.dispatch_outbox": {"queue": "events"},
|
||||||
"govoplan.events.purge_outbox": {"queue": "events"},
|
"govoplan.events.purge_outbox": {"queue": "events"},
|
||||||
},
|
},
|
||||||
@@ -74,6 +79,11 @@ celery.conf.update(
|
|||||||
"schedule": 24 * 60 * 60.0,
|
"schedule": 24 * 60 * 60.0,
|
||||||
"args": (500,),
|
"args": (500,),
|
||||||
},
|
},
|
||||||
|
"workflow-reconcile-every-five-seconds": {
|
||||||
|
"task": "govoplan.workflow.reconcile",
|
||||||
|
"schedule": 5.0,
|
||||||
|
"args": (50,),
|
||||||
|
},
|
||||||
"platform-events-every-ten-seconds": {
|
"platform-events-every-ten-seconds": {
|
||||||
"task": "govoplan.events.dispatch_outbox",
|
"task": "govoplan.events.dispatch_outbox",
|
||||||
"schedule": 10.0,
|
"schedule": 10.0,
|
||||||
@@ -157,6 +167,20 @@ def _dataflow_run_worker(
|
|||||||
return capability
|
return capability
|
||||||
|
|
||||||
|
|
||||||
|
def _workflow_runtime_worker(
|
||||||
|
registry: PlatformRegistry | None = None,
|
||||||
|
) -> WorkflowRuntimeWorker | None:
|
||||||
|
registry = registry or _platform_registry()
|
||||||
|
if not registry.has_capability(CAPABILITY_WORKFLOW_RUNTIME_WORKER):
|
||||||
|
return None
|
||||||
|
capability = registry.require_capability(
|
||||||
|
CAPABILITY_WORKFLOW_RUNTIME_WORKER
|
||||||
|
)
|
||||||
|
if not isinstance(capability, WorkflowRuntimeWorker):
|
||||||
|
raise RuntimeError("Workflow runtime worker capability is invalid")
|
||||||
|
return capability
|
||||||
|
|
||||||
|
|
||||||
def _platform_event_outbox(
|
def _platform_event_outbox(
|
||||||
registry: PlatformRegistry | None = None,
|
registry: PlatformRegistry | None = None,
|
||||||
) -> PlatformEventOutbox | None:
|
) -> PlatformEventOutbox | None:
|
||||||
@@ -310,6 +334,30 @@ def purge_dataflow_runs(self, limit: int = 500):
|
|||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
||||||
|
@celery.task(
|
||||||
|
name="govoplan.workflow.reconcile",
|
||||||
|
bind=True,
|
||||||
|
max_retries=0,
|
||||||
|
)
|
||||||
|
def reconcile_workflow_instances(self, limit: int = 50):
|
||||||
|
"""Resume asynchronous Workflow steps from durable provider state."""
|
||||||
|
|
||||||
|
from govoplan_core.db.session import get_database
|
||||||
|
|
||||||
|
with get_database().SessionLocal() as session:
|
||||||
|
provider = _workflow_runtime_worker()
|
||||||
|
if provider is None:
|
||||||
|
return {
|
||||||
|
"inspected": 0,
|
||||||
|
"advanced": 0,
|
||||||
|
"waiting": 0,
|
||||||
|
"failed": 0,
|
||||||
|
}
|
||||||
|
result = dict(provider.reconcile_pending(session, limit=limit))
|
||||||
|
session.commit()
|
||||||
|
return result
|
||||||
|
|
||||||
|
|
||||||
@celery.task(
|
@celery.task(
|
||||||
name="govoplan.events.dispatch_outbox",
|
name="govoplan.events.dispatch_outbox",
|
||||||
bind=True,
|
bind=True,
|
||||||
|
|||||||
44
src/govoplan_core/core/workflows.py
Normal file
44
src/govoplan_core/core/workflows.py
Normal file
@@ -0,0 +1,44 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from collections.abc import Mapping
|
||||||
|
from datetime import datetime
|
||||||
|
from typing import Protocol, runtime_checkable
|
||||||
|
|
||||||
|
|
||||||
|
CAPABILITY_WORKFLOW_RUNTIME_WORKER = "workflow.runtimeWorker"
|
||||||
|
|
||||||
|
|
||||||
|
@runtime_checkable
|
||||||
|
class WorkflowRuntimeWorker(Protocol):
|
||||||
|
def reconcile_pending(
|
||||||
|
self,
|
||||||
|
session: object,
|
||||||
|
*,
|
||||||
|
now: datetime | None = None,
|
||||||
|
limit: int = 50,
|
||||||
|
) -> Mapping[str, object]:
|
||||||
|
...
|
||||||
|
|
||||||
|
|
||||||
|
def workflow_runtime_worker(
|
||||||
|
registry: object | None,
|
||||||
|
) -> WorkflowRuntimeWorker | None:
|
||||||
|
if (
|
||||||
|
registry is None
|
||||||
|
or not hasattr(registry, "has_capability")
|
||||||
|
or not registry.has_capability(CAPABILITY_WORKFLOW_RUNTIME_WORKER)
|
||||||
|
):
|
||||||
|
return None
|
||||||
|
capability = registry.capability(CAPABILITY_WORKFLOW_RUNTIME_WORKER)
|
||||||
|
return (
|
||||||
|
capability
|
||||||
|
if isinstance(capability, WorkflowRuntimeWorker)
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"CAPABILITY_WORKFLOW_RUNTIME_WORKER",
|
||||||
|
"WorkflowRuntimeWorker",
|
||||||
|
"workflow_runtime_worker",
|
||||||
|
]
|
||||||
87
tests/test_workflow_runtime_worker.py
Normal file
87
tests/test_workflow_runtime_worker.py
Normal file
@@ -0,0 +1,87 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import unittest
|
||||||
|
from unittest.mock import MagicMock, patch
|
||||||
|
|
||||||
|
from govoplan_core.celery_app import (
|
||||||
|
celery,
|
||||||
|
reconcile_workflow_instances,
|
||||||
|
)
|
||||||
|
from govoplan_core.core.modules import ModuleContext, ModuleManifest
|
||||||
|
from govoplan_core.core.registry import PlatformRegistry
|
||||||
|
from govoplan_core.core.workflows import (
|
||||||
|
CAPABILITY_WORKFLOW_RUNTIME_WORKER,
|
||||||
|
WorkflowRuntimeWorker,
|
||||||
|
workflow_runtime_worker,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class _Worker:
|
||||||
|
def reconcile_pending(self, session, *, now=None, limit=50):
|
||||||
|
return {"inspected": 0, "advanced": 0}
|
||||||
|
|
||||||
|
|
||||||
|
class WorkflowRuntimeWorkerTests(unittest.TestCase):
|
||||||
|
def test_contract_is_runtime_checkable_and_resolved(self) -> None:
|
||||||
|
worker = _Worker()
|
||||||
|
self.assertIsInstance(worker, WorkflowRuntimeWorker)
|
||||||
|
registry = PlatformRegistry()
|
||||||
|
registry.register(
|
||||||
|
ModuleManifest(
|
||||||
|
id="workflow_runtime_test",
|
||||||
|
name="Workflow runtime test",
|
||||||
|
version="test",
|
||||||
|
capability_factories={
|
||||||
|
CAPABILITY_WORKFLOW_RUNTIME_WORKER: (
|
||||||
|
lambda _context: worker
|
||||||
|
),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
)
|
||||||
|
registry.configure_capability_context(
|
||||||
|
ModuleContext(registry=registry, settings=object())
|
||||||
|
)
|
||||||
|
|
||||||
|
self.assertIs(worker, workflow_runtime_worker(registry))
|
||||||
|
self.assertIsNone(workflow_runtime_worker(PlatformRegistry()))
|
||||||
|
|
||||||
|
def test_celery_task_commits_reconciliation(self) -> None:
|
||||||
|
session = MagicMock()
|
||||||
|
database = MagicMock()
|
||||||
|
database.SessionLocal.return_value.__enter__.return_value = session
|
||||||
|
worker = MagicMock()
|
||||||
|
worker.reconcile_pending.return_value = {
|
||||||
|
"inspected": 2,
|
||||||
|
"advanced": 1,
|
||||||
|
}
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch(
|
||||||
|
"govoplan_core.celery_app._workflow_runtime_worker",
|
||||||
|
return_value=worker,
|
||||||
|
),
|
||||||
|
patch(
|
||||||
|
"govoplan_core.db.session.get_database",
|
||||||
|
return_value=database,
|
||||||
|
),
|
||||||
|
):
|
||||||
|
result = reconcile_workflow_instances.run(25)
|
||||||
|
|
||||||
|
worker.reconcile_pending.assert_called_once_with(session, limit=25)
|
||||||
|
session.commit.assert_called_once_with()
|
||||||
|
self.assertEqual(1, result["advanced"])
|
||||||
|
|
||||||
|
def test_route_and_periodic_job_are_registered(self) -> None:
|
||||||
|
self.assertEqual(
|
||||||
|
{"queue": "workflow"},
|
||||||
|
celery.conf.task_routes["govoplan.workflow.reconcile"],
|
||||||
|
)
|
||||||
|
schedule = celery.conf.beat_schedule[
|
||||||
|
"workflow-reconcile-every-five-seconds"
|
||||||
|
]
|
||||||
|
self.assertEqual("govoplan.workflow.reconcile", schedule["task"])
|
||||||
|
self.assertEqual(5.0, schedule["schedule"])
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
Reference in New Issue
Block a user