diff --git a/src/govoplan_core/celery_app.py b/src/govoplan_core/celery_app.py index d373b84..42724ec 100644 --- a/src/govoplan_core/celery_app.py +++ b/src/govoplan_core/celery_app.py @@ -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.modules import ModuleContext 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.runtime import configure_runtime from govoplan_core.settings import settings @@ -47,6 +51,7 @@ celery.conf.update( "govoplan.dataflow.dispatch_runs": {"queue": "dataflow"}, "govoplan.dataflow.purge_runs": {"queue": "dataflow"}, "govoplan.dataflow.dispatch_triggers": {"queue": "dataflow"}, + "govoplan.workflow.reconcile": {"queue": "workflow"}, "govoplan.events.dispatch_outbox": {"queue": "events"}, "govoplan.events.purge_outbox": {"queue": "events"}, }, @@ -74,6 +79,11 @@ celery.conf.update( "schedule": 24 * 60 * 60.0, "args": (500,), }, + "workflow-reconcile-every-five-seconds": { + "task": "govoplan.workflow.reconcile", + "schedule": 5.0, + "args": (50,), + }, "platform-events-every-ten-seconds": { "task": "govoplan.events.dispatch_outbox", "schedule": 10.0, @@ -157,6 +167,20 @@ def _dataflow_run_worker( 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( registry: PlatformRegistry | None = None, ) -> PlatformEventOutbox | None: @@ -310,6 +334,30 @@ def purge_dataflow_runs(self, limit: int = 500): 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( name="govoplan.events.dispatch_outbox", bind=True, diff --git a/src/govoplan_core/core/workflows.py b/src/govoplan_core/core/workflows.py new file mode 100644 index 0000000..32c77f8 --- /dev/null +++ b/src/govoplan_core/core/workflows.py @@ -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", +] diff --git a/tests/test_workflow_runtime_worker.py b/tests/test_workflow_runtime_worker.py new file mode 100644 index 0000000..91fbe8e --- /dev/null +++ b/tests/test_workflow_runtime_worker.py @@ -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()