From a192a2215fdb35dc6ee05fd15e8f736232a71855 Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Wed, 29 Jul 2026 17:34:52 +0200 Subject: [PATCH] Add durable platform event delivery contract --- docs/DEPLOYMENT_OPERATOR_GUIDE.md | 6 +- src/govoplan_core/celery_app.py | 76 +++++++++++-- src/govoplan_core/core/events.py | 89 ++++++++++++++- src/govoplan_core/core/install_config.py | 8 +- src/govoplan_core/settings.py | 19 +++- tests/test_core_events.py | 37 +++++++ tests/test_install_config.py | 24 +++++ tests/test_platform_event_worker.py | 131 +++++++++++++++++++++++ 8 files changed, 373 insertions(+), 17 deletions(-) create mode 100644 tests/test_platform_event_worker.py diff --git a/docs/DEPLOYMENT_OPERATOR_GUIDE.md b/docs/DEPLOYMENT_OPERATOR_GUIDE.md index b4e6ac2..7c00222 100644 --- a/docs/DEPLOYMENT_OPERATOR_GUIDE.md +++ b/docs/DEPLOYMENT_OPERATOR_GUIDE.md @@ -157,13 +157,15 @@ release evidence. | --- | --- | --- | | `REDIS_URL` | `redis://redis:6379/0` | Celery broker/result backend when async workers are enabled. | | `CELERY_ENABLED` | `false` | Local/dev can send synchronously. Production campaign delivery should run workers and set this to `true`. | -| `CELERY_QUEUES` | `send_email,append_sent,notifications,calendar,default` | Queue list expected by worker/process manager definitions. The Calendar queue drains durable external-calendar operations. | +| `CELERY_QUEUES` | `send_email,append_sent,notifications,calendar,dataflow,events,default` | Queue list expected by worker/process manager definitions. The `events` queue drains transactional platform events; `dataflow` drains trigger deliveries and schedules. | +| `PLATFORM_EVENT_OUTBOX_MAX_ATTEMPTS` | `8` | Failed durable consumer deliveries are quarantined after this many attempts. | +| `PLATFORM_EVENT_OUTBOX_TERMINAL_RETENTION_DAYS` | `90` | Successful event envelopes older than this are removed by the daily retention task. Quarantined evidence is retained. | Worker command: ```bash python -m celery -A govoplan_core.celery_app:celery worker \ - --queues send_email,append_sent,notifications,calendar,default \ + --queues send_email,append_sent,notifications,calendar,dataflow,events,default \ --loglevel INFO ``` diff --git a/src/govoplan_core/celery_app.py b/src/govoplan_core/celery_app.py index 4f4a44c..227bfd1 100644 --- a/src/govoplan_core/celery_app.py +++ b/src/govoplan_core/celery_app.py @@ -1,5 +1,7 @@ from __future__ import annotations +from datetime import datetime, timedelta, timezone + from celery import Celery from govoplan_core.core.campaigns import CAPABILITY_CAMPAIGNS_DELIVERY_TASKS, CampaignDeliveryTaskProvider @@ -10,6 +12,7 @@ from govoplan_core.core.dataflows import ( ) from govoplan_core.core.events import ( CAPABILITY_PLATFORM_EVENT_OUTBOX, + DurableEventConsumer, PlatformEvent, PlatformEventOutbox, publish_platform_event, @@ -41,6 +44,7 @@ celery.conf.update( "govoplan.calendar.dispatch_outbox": {"queue": "calendar"}, "govoplan.dataflow.dispatch_triggers": {"queue": "dataflow"}, "govoplan.events.dispatch_outbox": {"queue": "events"}, + "govoplan.events.purge_outbox": {"queue": "events"}, }, worker_prefetch_multiplier=1, task_acks_late=True, @@ -61,6 +65,11 @@ celery.conf.update( "schedule": 10.0, "args": (100,), }, + "platform-event-retention-daily": { + "task": "govoplan.events.purge_outbox", + "schedule": 24 * 60 * 60.0, + "args": (500,), + }, }, ) @@ -231,7 +240,7 @@ def dispatch_dataflow_triggers(self, limit: int = 100): max_retries=0, ) def dispatch_platform_events(self, limit: int = 100): - """Deliver committed platform events once across all worker processes.""" + """Deliver committed platform events through persistent consumer ledgers.""" from govoplan_core.db.session import get_database @@ -239,21 +248,68 @@ def dispatch_platform_events(self, limit: int = 100): registry = _platform_registry() outbox = _platform_event_outbox(registry) if outbox is None: - return {"selected": 0, "dispatched": 0, "failed": 0} + return { + "selected": 0, + "delivered": 0, + "retrying": 0, + "quarantined": 0, + "dispatched": 0, + "observer_failed": 0, + } dataflow_dispatcher = _dataflow_trigger_dispatcher(registry) + consumers = () + if dataflow_dispatcher is not None: + def deliver_to_dataflow( + event: PlatformEvent, + _delivery_key: str, + ) -> None: + dataflow_dispatcher.ingest_event( + session, + event=event, + ) - def dispatch(event: PlatformEvent) -> None: - if ( - dataflow_dispatcher is not None - and event.classification in {"public", "internal"} - ): - dataflow_dispatcher.ingest_event(session, event=event) - publish_platform_event(event) + consumers = ( + DurableEventConsumer( + consumer_id="dataflow.event-triggers.v1", + event_types=frozenset({"*"}), + classifications=frozenset({"public", "internal"}), + handler=deliver_to_dataflow, + ), + ) result = dict( outbox.dispatch_pending( session, - dispatcher=dispatch, + consumers=consumers, + observer=publish_platform_event, + limit=limit, + ) + ) + session.commit() + return result + + +@celery.task( + name="govoplan.events.purge_outbox", + bind=True, + max_retries=0, +) +def purge_platform_events(self, limit: int = 500): + """Remove old terminal event envelopes while retaining quarantine evidence.""" + + from govoplan_core.db.session import get_database + + with get_database().SessionLocal() as session: + outbox = _platform_event_outbox() + if outbox is None: + return {"deleted": 0} + before = datetime.now(timezone.utc) - timedelta( + days=settings.platform_event_outbox_terminal_retention_days + ) + result = dict( + outbox.purge_terminal( + session, + before=before, limit=limit, ) ) diff --git a/src/govoplan_core/core/events.py b/src/govoplan_core/core/events.py index f7b38e8..a3d9a19 100644 --- a/src/govoplan_core/core/events.py +++ b/src/govoplan_core/core/events.py @@ -1,7 +1,7 @@ from __future__ import annotations from collections import defaultdict -from collections.abc import Callable, Mapping +from collections.abc import Callable, Mapping, Sequence from contextlib import contextmanager from contextvars import ContextVar from dataclasses import dataclass, field @@ -15,6 +15,7 @@ from sqlalchemy.orm import Session _TRACE_ID_RE = re.compile(r"^[A-Za-z0-9_.:-]{1,128}$") +_CONSUMER_ID_RE = re.compile(r"^[a-z][a-z0-9_.:-]{0,127}$") _PENDING_EVENTS_KEY = "govoplan.pending_platform_events" CAPABILITY_PLATFORM_EVENT_OUTBOX = "platform.eventOutbox" @@ -105,6 +106,63 @@ class PlatformEvent: EventHandler = Callable[[PlatformEvent], None] +DurableEventHandler = Callable[[PlatformEvent, str], None] + + +@dataclass(frozen=True, slots=True) +class DurableEventConsumer: + """Allowlisted durable consumer with an explicit disclosure boundary.""" + + consumer_id: str + handler: DurableEventHandler + event_types: frozenset[str] = field( + default_factory=lambda: frozenset({"*"}) + ) + classifications: frozenset[EventClassification] = field( + default_factory=lambda: frozenset({"public", "internal"}) + ) + policy_decision_ref: str | None = None + + def __post_init__(self) -> None: + if not _CONSUMER_ID_RE.fullmatch(self.consumer_id): + raise ValueError("Durable event consumer id is invalid") + if not self.event_types or any( + item != "*" and not _TRACE_ID_RE.fullmatch(item) + for item in self.event_types + ): + raise ValueError( + "Durable event consumers require valid event-type allowlists" + ) + invalid_classifications = set(self.classifications) - { + "public", + "internal", + "confidential", + "restricted", + } + if not self.classifications or invalid_classifications: + raise ValueError( + "Durable event consumer classifications are invalid" + ) + if ( + self.classifications & {"confidential", "restricted"} + and not normalize_trace_id(self.policy_decision_ref) + ): + raise ValueError( + "Confidential or restricted event subscriptions require " + "an explicit policy decision reference" + ) + + def accepts(self, event: PlatformEvent) -> bool: + return ( + event.classification in self.classifications + and ( + "*" in self.event_types + or event.type in self.event_types + ) + ) + + def delivery_key(self, event: PlatformEvent) -> str: + return f"{event.event_id}:{self.consumer_id}" @runtime_checkable @@ -116,11 +174,38 @@ class PlatformEventOutbox(Protocol): self, session: object, *, - dispatcher: EventHandler, + consumers: Sequence[DurableEventConsumer] = (), + observer: EventHandler | None = None, limit: int = 100, ) -> Mapping[str, int]: ... + def replay_delivery( + self, + session: object, + *, + event_id: str, + consumer_id: str, + operator_id: str, + reason: str, + ) -> Mapping[str, object]: + ... + + def purge_terminal( + self, + session: object, + *, + before: datetime, + limit: int = 500, + ) -> Mapping[str, int]: + ... + + def delivery_metrics( + self, + session: object, + ) -> Mapping[str, object]: + ... + def current_event_trace() -> EventTrace | None: return _current_trace.get() diff --git a/src/govoplan_core/core/install_config.py b/src/govoplan_core/core/install_config.py index 8d51f81..7577c10 100644 --- a/src/govoplan_core/core/install_config.py +++ b/src/govoplan_core/core/install_config.py @@ -338,8 +338,10 @@ ENABLED_MODULES=tenancy,organizations,identity,access,admin,dashboard,policy,aud CELERY_ENABLED=true REDIS_URL=redis://127.0.0.1:6379/0 -CELERY_QUEUES=send_email,append_sent,notifications,calendar,default +CELERY_QUEUES=send_email,append_sent,notifications,calendar,dataflow,events,default CALENDAR_OUTBOX_TERMINAL_RETENTION_DAYS=90 +PLATFORM_EVENT_OUTBOX_MAX_ATTEMPTS=8 +PLATFORM_EVENT_OUTBOX_TERMINAL_RETENTION_DAYS=90 SCHEDULING_CANCELLATION_NOTICE_DAYS=30 # Deployment-wide connector egress policy. Enable private networks only when @@ -404,8 +406,10 @@ DATABASE_URL=postgresql+psycopg://govoplan:govoplan-dev@127.0.0.1:55433/govoplan GOVOPLAN_DATABASE_URL_PGTOOLS=postgresql://govoplan:govoplan-dev@127.0.0.1:55433/govoplan REDIS_URL=redis://127.0.0.1:56379/0 CELERY_ENABLED=true -CELERY_QUEUES=send_email,append_sent,notifications,calendar,default +CELERY_QUEUES=send_email,append_sent,notifications,calendar,dataflow,events,default CALENDAR_OUTBOX_TERMINAL_RETENTION_DAYS=90 +PLATFORM_EVENT_OUTBOX_MAX_ATTEMPTS=8 +PLATFORM_EVENT_OUTBOX_TERMINAL_RETENTION_DAYS=90 SCHEDULING_CANCELLATION_NOTICE_DAYS=30 GOVOPLAN_CONNECTOR_ALLOW_PRIVATE_NETWORKS=true diff --git a/src/govoplan_core/settings.py b/src/govoplan_core/settings.py index b9f3d8a..99fed2f 100644 --- a/src/govoplan_core/settings.py +++ b/src/govoplan_core/settings.py @@ -111,12 +111,29 @@ class Settings(BaseSettings): ) master_key_b64: str | None = Field(default=None, alias="MASTER_KEY_B64") - celery_queues: str = Field(default="send_email,append_sent,notifications,calendar,default", alias="CELERY_QUEUES") + celery_queues: str = Field( + default=( + "send_email,append_sent,notifications,calendar," + "dataflow,events,default" + ), + alias="CELERY_QUEUES", + ) calendar_outbox_terminal_retention_days: int = Field( default=90, ge=0, alias="CALENDAR_OUTBOX_TERMINAL_RETENTION_DAYS", ) + platform_event_outbox_max_attempts: int = Field( + default=8, + ge=1, + le=100, + alias="PLATFORM_EVENT_OUTBOX_MAX_ATTEMPTS", + ) + platform_event_outbox_terminal_retention_days: int = Field( + default=90, + ge=0, + alias="PLATFORM_EVENT_OUTBOX_TERMINAL_RETENTION_DAYS", + ) scheduling_cancellation_notice_days: int = Field( default=30, ge=1, diff --git a/tests/test_core_events.py b/tests/test_core_events.py index 1d8886d..b6f7f2c 100644 --- a/tests/test_core_events.py +++ b/tests/test_core_events.py @@ -13,6 +13,7 @@ from sqlalchemy.orm import Session from govoplan_core.audit.logging import audit_event, audit_operation_context from govoplan_core.core.events import ( + DurableEventConsumer, EventActorRef, EventBus, EventObjectRef, @@ -41,6 +42,42 @@ def _configure_audit_runtime() -> None: class CoreEventTests(unittest.TestCase): + def test_durable_consumer_requires_policy_for_classified_events(self) -> None: + with self.assertRaisesRegex(ValueError, "policy decision"): + DurableEventConsumer( + consumer_id="workflow.triggers.v1", + classifications=frozenset({"restricted"}), + handler=lambda _event, _delivery_key: None, + ) + + consumer = DurableEventConsumer( + consumer_id="workflow.triggers.v1", + event_types=frozenset({"case.changed"}), + classifications=frozenset( + {"internal", "confidential"} + ), + policy_decision_ref="policy:decision:1", + handler=lambda _event, _delivery_key: None, + ) + accepted = PlatformEvent( + type="case.changed", + module_id="cases", + classification="confidential", + event_id="event-1", + ) + rejected = PlatformEvent( + type="case.deleted", + module_id="cases", + classification="confidential", + ) + + self.assertTrue(consumer.accepts(accepted)) + self.assertFalse(consumer.accepts(rejected)) + self.assertEqual( + "event-1:workflow.triggers.v1", + consumer.delivery_key(accepted), + ) + def test_fallback_event_waits_for_outer_commit_after_savepoint(self) -> None: engine = create_engine("sqlite:///:memory:") seen: list[PlatformEvent] = [] diff --git a/tests/test_install_config.py b/tests/test_install_config.py index 7eca299..1be03ec 100644 --- a/tests/test_install_config.py +++ b/tests/test_install_config.py @@ -21,6 +21,30 @@ class InstallConfigTests(unittest.TestCase): with self.assertRaises(ValidationError): Settings(CALENDAR_OUTBOX_TERMINAL_RETENTION_DAYS="-1") + def test_platform_event_outbox_limits_are_validated(self) -> None: + defaults = Settings() + self.assertEqual(defaults.platform_event_outbox_max_attempts, 8) + self.assertEqual( + defaults.platform_event_outbox_terminal_retention_days, + 90, + ) + configured = Settings( + PLATFORM_EVENT_OUTBOX_MAX_ATTEMPTS="12", + PLATFORM_EVENT_OUTBOX_TERMINAL_RETENTION_DAYS="30", + ) + self.assertEqual(configured.platform_event_outbox_max_attempts, 12) + self.assertEqual( + configured.platform_event_outbox_terminal_retention_days, + 30, + ) + for values in ( + {"PLATFORM_EVENT_OUTBOX_MAX_ATTEMPTS": "0"}, + {"PLATFORM_EVENT_OUTBOX_MAX_ATTEMPTS": "101"}, + {"PLATFORM_EVENT_OUTBOX_TERMINAL_RETENTION_DAYS": "-1"}, + ): + with self.subTest(values=values), self.assertRaises(ValidationError): + Settings(**values) + def test_scheduling_cancellation_notice_setting_is_bounded(self) -> None: self.assertEqual(Settings().scheduling_cancellation_notice_days, 30) self.assertEqual( diff --git a/tests/test_platform_event_worker.py b/tests/test_platform_event_worker.py new file mode 100644 index 0000000..1ed9beb --- /dev/null +++ b/tests/test_platform_event_worker.py @@ -0,0 +1,131 @@ +from __future__ import annotations + +import unittest +from datetime import datetime, timezone +from unittest.mock import MagicMock, patch + +from govoplan_core.celery_app import ( + celery, + dispatch_platform_events, + purge_platform_events, +) +from govoplan_core.core.events import PlatformEvent + + +class PlatformEventWorkerTests(unittest.TestCase): + def test_dispatch_uses_a_durable_dataflow_consumer_and_commits(self) -> None: + session = MagicMock() + database = MagicMock() + database.SessionLocal.return_value.__enter__.return_value = session + outbox = MagicMock() + outbox.dispatch_pending.return_value = { + "selected": 1, + "delivered": 1, + "retrying": 0, + "quarantined": 0, + "dispatched": 1, + "observer_failed": 0, + } + dataflow = MagicMock() + registry = MagicMock() + + with ( + patch( + "govoplan_core.celery_app._platform_registry", + return_value=registry, + ), + patch( + "govoplan_core.celery_app._platform_event_outbox", + return_value=outbox, + ), + patch( + "govoplan_core.celery_app._dataflow_trigger_dispatcher", + return_value=dataflow, + ), + patch( + "govoplan_core.db.session.get_database", + return_value=database, + ), + ): + result = dispatch_platform_events.run(25) + + call = outbox.dispatch_pending.call_args + self.assertEqual(session, call.args[0]) + self.assertEqual(25, call.kwargs["limit"]) + consumer = call.kwargs["consumers"][0] + self.assertEqual( + "dataflow.event-triggers.v1", + consumer.consumer_id, + ) + self.assertEqual( + frozenset({"public", "internal"}), + consumer.classifications, + ) + event = PlatformEvent( + type="files.file.created", + module_id="files", + ) + consumer.handler(event, consumer.delivery_key(event)) + dataflow.ingest_event.assert_called_once_with( + session, + event=event, + ) + session.commit.assert_called_once_with() + self.assertEqual(1, result["delivered"]) + + def test_retention_task_uses_the_configured_terminal_window(self) -> None: + session = MagicMock() + database = MagicMock() + database.SessionLocal.return_value.__enter__.return_value = session + outbox = MagicMock() + outbox.purge_terminal.return_value = {"deleted": 2} + + with ( + patch( + "govoplan_core.celery_app._platform_event_outbox", + return_value=outbox, + ), + patch( + "govoplan_core.db.session.get_database", + return_value=database, + ), + patch( + "govoplan_core.celery_app.settings." + "platform_event_outbox_terminal_retention_days", + 30, + ), + ): + result = purge_platform_events.run(75) + + call = outbox.purge_terminal.call_args + self.assertEqual(session, call.args[0]) + self.assertEqual(75, call.kwargs["limit"]) + before = call.kwargs["before"] + self.assertIsInstance(before, datetime) + self.assertEqual(timezone.utc, before.tzinfo) + session.commit.assert_called_once_with() + self.assertEqual({"deleted": 2}, result) + + def test_worker_routes_and_periodic_tasks_are_registered(self) -> None: + self.assertEqual( + {"queue": "events"}, + celery.conf.task_routes[ + "govoplan.events.dispatch_outbox" + ], + ) + self.assertEqual( + {"queue": "events"}, + celery.conf.task_routes[ + "govoplan.events.purge_outbox" + ], + ) + self.assertEqual( + "govoplan.events.purge_outbox", + celery.conf.beat_schedule[ + "platform-event-retention-daily" + ]["task"], + ) + + +if __name__ == "__main__": + unittest.main()