From 9e219bc4d3f559666ee670112c93118d0a3332d0 Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Thu, 30 Jul 2026 03:59:38 +0200 Subject: [PATCH] feat(core): dispatch bounded postbox routes --- src/govoplan_core/celery_app.py | 59 ++++++++++++ src/govoplan_core/core/organizations.py | 14 +++ src/govoplan_core/core/postbox.py | 24 +++++ tests/test_organization_hierarchy_contract.py | 4 + tests/test_postbox_routing_worker.py | 93 +++++++++++++++++++ 5 files changed, 194 insertions(+) create mode 100644 tests/test_postbox_routing_worker.py diff --git a/src/govoplan_core/celery_app.py b/src/govoplan_core/celery_app.py index 42724ec..de17d83 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.postbox import ( + CAPABILITY_POSTBOX_ROUTING, + PostboxRoutingProvider, +) from govoplan_core.core.workflows import ( CAPABILITY_WORKFLOW_RUNTIME_WORKER, WorkflowRuntimeWorker, @@ -52,6 +56,7 @@ celery.conf.update( "govoplan.dataflow.purge_runs": {"queue": "dataflow"}, "govoplan.dataflow.dispatch_triggers": {"queue": "dataflow"}, "govoplan.workflow.reconcile": {"queue": "workflow"}, + "govoplan.postbox.dispatch_routes": {"queue": "postbox"}, "govoplan.events.dispatch_outbox": {"queue": "events"}, "govoplan.events.purge_outbox": {"queue": "events"}, }, @@ -84,6 +89,11 @@ celery.conf.update( "schedule": 5.0, "args": (50,), }, + "postbox-routes-every-minute": { + "task": "govoplan.postbox.dispatch_routes", + "schedule": 60.0, + "args": (None, 50), + }, "platform-events-every-ten-seconds": { "task": "govoplan.events.dispatch_outbox", "schedule": 10.0, @@ -181,6 +191,18 @@ def _workflow_runtime_worker( return capability +def _postbox_routing_provider( + registry: PlatformRegistry | None = None, +) -> PostboxRoutingProvider | None: + registry = registry or _platform_registry() + if not registry.has_capability(CAPABILITY_POSTBOX_ROUTING): + return None + capability = registry.require_capability(CAPABILITY_POSTBOX_ROUTING) + if not isinstance(capability, PostboxRoutingProvider): + raise RuntimeError("Postbox routing capability is invalid") + return capability + + def _platform_event_outbox( registry: PlatformRegistry | None = None, ) -> PlatformEventOutbox | None: @@ -358,6 +380,43 @@ def reconcile_workflow_instances(self, limit: int = 50): return result +@celery.task( + name="govoplan.postbox.dispatch_routes", + bind=True, + max_retries=0, +) +def dispatch_postbox_routes( + self, + tenant_id: str | None = None, + limit: int = 50, +): + """Deliver due Postbox vacancy escalations from durable route rows.""" + + from govoplan_core.db.session import get_database + + with get_database().SessionLocal() as session: + provider = _postbox_routing_provider() + if provider is None: + return { + "selected": 0, + "delivered": 0, + "vacant": 0, + "rescheduled": 0, + "cancelled": 0, + "failed": 0, + "route_ids": [], + } + result = dict( + provider.dispatch_due_routes( + session, + tenant_id=tenant_id, + limit=limit, + ) + ) + session.commit() + return result + + @celery.task( name="govoplan.events.dispatch_outbox", bind=True, diff --git a/src/govoplan_core/core/organizations.py b/src/govoplan_core/core/organizations.py index 345d221..4ddb6e6 100644 --- a/src/govoplan_core/core/organizations.py +++ b/src/govoplan_core/core/organizations.py @@ -130,6 +130,13 @@ class OrganizationRelationTypeRef: status: OrganizationStatus = "active" +@dataclass(frozen=True, slots=True) +class OrganizationHierarchyCatalogRef: + tenant_id: str + structures: tuple[OrganizationStructureRef, ...] = () + relation_types: tuple[OrganizationRelationTypeRef, ...] = () + + @dataclass(frozen=True, slots=True) class OrganizationHierarchyEdgeRef: id: str @@ -231,6 +238,12 @@ class OrganizationDirectory(Protocol): @runtime_checkable class OrganizationHierarchyDirectory(Protocol): + def hierarchy_catalog( + self, + tenant_id: str, + ) -> OrganizationHierarchyCatalogRef: + ... + def get_unit_type( self, tenant_id: str, @@ -355,6 +368,7 @@ __all__ = [ "OrganizationFunctionTypeRef", "OrganizationFunctionTypeResolution", "OrganizationHierarchyDirection", + "OrganizationHierarchyCatalogRef", "OrganizationHierarchyDirectory", "OrganizationHierarchyEdgeRef", "OrganizationHierarchyMatchRef", diff --git a/src/govoplan_core/core/postbox.py b/src/govoplan_core/core/postbox.py index 264cc76..c930745 100644 --- a/src/govoplan_core/core/postbox.py +++ b/src/govoplan_core/core/postbox.py @@ -12,6 +12,7 @@ CAPABILITY_POSTBOX_ACCESS = "postbox.access" CAPABILITY_POSTBOX_MESSAGES = "postbox.messages" CAPABILITY_POSTBOX_DELIVERY = "postbox.delivery" CAPABILITY_POSTBOX_EVIDENCE = "postbox.evidence" +CAPABILITY_POSTBOX_ROUTING = "postbox.routing" PostboxAction = Literal["discover", "read", "send", "acknowledge", "administer"] PostboxMessageListState = Literal["all", "unread", "read", "acknowledged"] @@ -317,6 +318,18 @@ class PostboxEvidenceProvider(Protocol): ... +@runtime_checkable +class PostboxRoutingProvider(Protocol): + def dispatch_due_routes( + self, + session: object, + *, + tenant_id: str | None = None, + limit: int = 50, + ) -> Mapping[str, object]: + ... + + def _postbox_provider( registry: object | None, *, @@ -384,3 +397,14 @@ def postbox_evidence_provider( provider_type=PostboxEvidenceProvider, ) return provider if isinstance(provider, PostboxEvidenceProvider) else None + + +def postbox_routing_provider( + registry: object | None, +) -> PostboxRoutingProvider | None: + provider = _postbox_provider( + registry, + capability_name=CAPABILITY_POSTBOX_ROUTING, + provider_type=PostboxRoutingProvider, + ) + return provider if isinstance(provider, PostboxRoutingProvider) else None diff --git a/tests/test_organization_hierarchy_contract.py b/tests/test_organization_hierarchy_contract.py index 66227b0..1662d89 100644 --- a/tests/test_organization_hierarchy_contract.py +++ b/tests/test_organization_hierarchy_contract.py @@ -6,6 +6,7 @@ from govoplan_core.core.organizations import ( CAPABILITY_ORGANIZATION_HIERARCHY_DIRECTORY, OrganizationDirectory, OrganizationFunctionTypeRef, + OrganizationHierarchyCatalogRef, OrganizationHierarchyDirectory, OrganizationUnitTypeRef, organization_hierarchy_directory, @@ -35,6 +36,9 @@ class LegacyOrganizationDirectory: class HierarchyDirectory: + def hierarchy_catalog(self, tenant_id): + return OrganizationHierarchyCatalogRef(tenant_id=tenant_id) + def get_unit_type(self, _tenant_id, _unit_type_id): return None diff --git a/tests/test_postbox_routing_worker.py b/tests/test_postbox_routing_worker.py new file mode 100644 index 0000000..00dc447 --- /dev/null +++ b/tests/test_postbox_routing_worker.py @@ -0,0 +1,93 @@ +from __future__ import annotations + +import unittest +from unittest.mock import MagicMock, patch + +from govoplan_core.celery_app import celery, dispatch_postbox_routes +from govoplan_core.core.modules import ModuleContext, ModuleManifest +from govoplan_core.core.postbox import ( + CAPABILITY_POSTBOX_ROUTING, + PostboxRoutingProvider, + postbox_routing_provider, +) +from govoplan_core.core.registry import PlatformRegistry + + +class _RoutingProvider: + def dispatch_due_routes( + self, + session, + *, + tenant_id=None, + limit=50, + ): + del session, tenant_id, limit + return {"selected": 0} + + +class PostboxRoutingWorkerTests(unittest.TestCase): + def test_contract_is_runtime_checkable_and_resolved(self) -> None: + provider = _RoutingProvider() + self.assertIsInstance(provider, PostboxRoutingProvider) + registry = PlatformRegistry() + registry.register( + ModuleManifest( + id="postbox_routing_test", + name="Postbox routing test", + version="test", + capability_factories={ + CAPABILITY_POSTBOX_ROUTING: ( + lambda _context: provider + ) + }, + ) + ) + registry.configure_capability_context( + ModuleContext(registry=registry, settings=object()) + ) + + self.assertIs(provider, postbox_routing_provider(registry)) + self.assertIsNone(postbox_routing_provider(PlatformRegistry())) + + def test_worker_commits_provider_outcome(self) -> None: + session = MagicMock() + database = MagicMock() + database.SessionLocal.return_value.__enter__.return_value = session + provider = MagicMock() + provider.dispatch_due_routes.return_value = { + "selected": 1, + "delivered": 1, + } + + with ( + patch( + "govoplan_core.celery_app._postbox_routing_provider", + return_value=provider, + ), + patch( + "govoplan_core.db.session.get_database", + return_value=database, + ), + ): + result = dispatch_postbox_routes.run("tenant-1", 25) + + provider.dispatch_due_routes.assert_called_once_with( + session, + tenant_id="tenant-1", + limit=25, + ) + session.commit.assert_called_once_with() + self.assertEqual(1, result["delivered"]) + + def test_route_and_periodic_recovery_are_registered(self) -> None: + self.assertEqual( + {"queue": "postbox"}, + celery.conf.task_routes["govoplan.postbox.dispatch_routes"], + ) + schedule = celery.conf.beat_schedule["postbox-routes-every-minute"] + self.assertEqual("govoplan.postbox.dispatch_routes", schedule["task"]) + self.assertEqual(60.0, schedule["schedule"]) + + +if __name__ == "__main__": + unittest.main()