Add durable platform event delivery contract

This commit is contained in:
2026-07-29 17:34:52 +02:00
parent e8fed6d25a
commit a192a2215f
8 changed files with 373 additions and 17 deletions

View File

@@ -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
```

View File

@@ -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,
)
)

View File

@@ -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()

View File

@@ -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

View File

@@ -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,

View File

@@ -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] = []

View File

@@ -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(

View File

@@ -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()