feat(postbox): reconcile lifecycle notifications

This commit is contained in:
2026-08-20 04:58:34 +02:00
parent 41ea8d8e23
commit 8a21876634
8 changed files with 765 additions and 9 deletions
+2
View File
@@ -353,6 +353,7 @@ class PostboxRouterTests(unittest.TestCase):
"idempotency_key": "campaign-1:recipient-1",
"subject": "Decision",
"body_text": "The decision is ready.",
"action_required": True,
},
)
self.assertEqual(201, delivery.status_code, delivery.text)
@@ -365,6 +366,7 @@ class PostboxRouterTests(unittest.TestCase):
self.assertEqual(200, messages.status_code, messages.text)
self.assertEqual(1, messages.json()["total"])
self.assertEqual(message_id, messages.json()["messages"][0]["id"])
self.assertTrue(messages.json()["messages"][0]["metadata"]["action_required"])
filtered = self.client.get(
"/api/v1/postbox/messages",
+216 -1
View File
@@ -509,6 +509,28 @@ class FakeNotificationDispatch:
return {"id": f"notification-{len(self.requests)}"}
class FailingOnceNotificationDispatch(FakeNotificationDispatch):
def __init__(self) -> None:
super().__init__()
self.fail_next = True
def enqueue_notification(
self,
session,
request,
*,
enqueue_delivery=True,
):
if self.fail_next:
self.fail_next = False
raise RuntimeError("notification provider unavailable")
return super().enqueue_notification(
session,
request,
enqueue_delivery=enqueue_delivery,
)
class PostboxServiceTests(unittest.TestCase):
def setUp(self) -> None:
self.engine = create_engine("sqlite:///:memory:")
@@ -602,6 +624,7 @@ class PostboxServiceTests(unittest.TestCase):
vacancy_escalation: bool = False,
structure_id: str = "structure-1",
max_depth: int = 2,
notifications: FakeNotificationDispatch | None = None,
) -> tuple[PostboxService, Postbox, PostboxTemplate]:
service = PostboxService(
identities=FakeIdentityDirectory(), # type: ignore[arg-type]
@@ -609,6 +632,7 @@ class PostboxServiceTests(unittest.TestCase):
incumbencies=self.idm, # type: ignore[arg-type]
organizations=self.organizations, # type: ignore[arg-type]
hierarchy=self.organizations, # type: ignore[arg-type]
notifications=notifications, # type: ignore[arg-type]
)
target_template = service.create_template(
session,
@@ -1229,6 +1253,180 @@ class PostboxServiceTests(unittest.TestCase):
[event.type for event in events],
)
def test_notification_lifecycle_reconciles_assignment_churn_and_retries(
self,
) -> None:
self.idm.assignments.append(self.assignment)
notifications = FailingOnceNotificationDispatch()
service = PostboxService(
identities=FakeIdentityDirectory(), # type: ignore[arg-type]
idm=self.idm, # type: ignore[arg-type]
incumbencies=self.idm, # type: ignore[arg-type]
organizations=self.organizations, # type: ignore[arg-type]
notifications=notifications, # type: ignore[arg-type]
)
event_bus = EventBus()
events = []
event_bus.subscribe("*", events.append)
with Session(self.engine) as session, event_bus_context(event_bus):
postbox = service.create_exact_postbox(
session,
tenant_id="tenant-1",
name="District North / Lifecycle notifications",
organization_unit_id="unit-1",
function_id="function-1",
address_key=None,
description=None,
classification="internal",
actor_id="admin-1",
)
session.commit()
baseline = service.reconcile_notification_lifecycle(
session,
tenant_id="tenant-1",
)
session.commit()
self.assertEqual(1, baseline["scanned"])
self.assertEqual(0, baseline["changed"])
baseline_event_count = len(events)
delegated = replace(
self.assignment,
id="assignment-delegated",
identity_id="identity-2",
account_id="account-2",
source="delegated",
delegated_from_assignment_id=self.assignment.id,
)
overlapping_delegation = replace(
delegated,
id="assignment-delegated-overlap",
)
self.idm.assignments.extend((delegated, overlapping_delegation))
added = service.reconcile_notification_lifecycle(
session,
tenant_id="tenant-1",
)
session.commit()
self.assertEqual(1, added["changed"])
self.assertEqual(3, added["events"])
self.assertEqual(0, added["notifications"])
self.assertEqual(1, added["notification_failures"])
retried = service.reconcile_notification_lifecycle(
session,
tenant_id="tenant-1",
)
session.commit()
self.assertEqual(0, retried["changed"])
self.assertEqual(0, retried["events"])
self.assertEqual(1, retried["notifications"])
self.assertEqual(1, len(notifications.requests))
request, enqueue_delivery = notifications.requests[0]
self.assertEqual("postbox.delegation.newly_visible.v1", request.event_kind)
self.assertEqual("account-2", request.recipient_id)
self.assertEqual(f"/postbox?postbox={postbox.id}", request.action_url)
self.assertFalse(enqueue_delivery)
self.idm.assignments.clear()
vacant = service.reconcile_notification_lifecycle(
session,
tenant_id="tenant-1",
)
session.commit()
self.assertEqual(1, vacant["changed"])
self.assertEqual(5, vacant["events"])
replacement = replace(
self.assignment,
id="assignment-replacement",
identity_id="identity-3",
account_id="account-3",
)
self.idm.assignments.append(replacement)
filled = service.reconcile_notification_lifecycle(
session,
tenant_id="tenant-1",
)
session.commit()
self.assertEqual(1, filled["changed"])
self.assertEqual(3, filled["events"])
self.assertEqual(1, filled["notifications"])
lifecycle_events = events[baseline_event_count:]
self.assertEqual(
{
"postbox.assignment.reassigned.v1",
"postbox.assignment.newly_visible.v1",
"postbox.assignment.visibility_withdrawn.v1",
"postbox.delegation.newly_visible.v1",
"postbox.vacancy.detected.v1",
"postbox.vacancy.resolved.v1",
},
{event.type for event in lifecycle_events},
)
for event in lifecycle_events:
serialized = repr(dict(event.payload)).casefold()
self.assertNotIn("body", serialized)
self.assertNotIn("attachment", serialized)
self.assertNotIn("subject", serialized)
def test_action_required_delivery_uses_metadata_only_notification(self) -> None:
self.idm.assignments.append(self.assignment)
notifications = FakeNotificationDispatch()
service = PostboxService(
identities=FakeIdentityDirectory(), # type: ignore[arg-type]
idm=self.idm, # type: ignore[arg-type]
incumbencies=self.idm, # type: ignore[arg-type]
organizations=self.organizations, # type: ignore[arg-type]
notifications=notifications, # type: ignore[arg-type]
)
events = []
event_bus = EventBus()
event_bus.subscribe("*", events.append)
with Session(self.engine) as session, event_bus_context(event_bus):
postbox = service.create_exact_postbox(
session,
tenant_id="tenant-1",
name="District North / Required action",
organization_unit_id="unit-1",
function_id="function-1",
address_key=None,
description=None,
classification="internal",
actor_id="admin-1",
)
delivered = service.deliver(
session,
PostboxDeliveryRequest(
tenant_id="tenant-1",
target=PostboxTargetRef(postbox_id=postbox.id),
producer_module="tests",
producer_resource_type="required_action",
producer_resource_id="action-1",
idempotency_key="required-action-1",
subject="Sensitive subject is not copied",
body_text="Sensitive body is not copied",
action_required=True,
),
)
session.commit()
self.assertIn(
"postbox.message.action_required.v1",
[event.type for event in events],
)
notification, _enqueue_delivery = notifications.requests[0]
self.assertEqual(
"postbox.message.action_required.v1",
notification.event_kind,
)
self.assertEqual(4, notification.priority)
self.assertEqual(delivered.message_id, notification.payload["message_id"])
self.assertNotIn("Sensitive", notification.body_text or "")
self.assertNotIn("Sensitive", notification.subject or "")
def test_withdrawn_and_expired_messages_keep_metadata_but_hide_content(
self,
) -> None:
@@ -1945,10 +2143,15 @@ class PostboxServiceTests(unittest.TestCase):
source="direct",
)
self.idm.assignments.append(top_assignment)
with Session(self.engine) as session:
notifications = FakeNotificationDispatch()
events = []
event_bus = EventBus()
event_bus.subscribe("*", events.append)
with Session(self.engine) as session, event_bus_context(event_bus):
service, source, _target_template = self._create_routing_source(
session,
vacancy_escalation=True,
notifications=notifications,
)
service.deliver(
session,
@@ -1983,6 +2186,18 @@ class PostboxServiceTests(unittest.TestCase):
self.assertEqual(frozen_target_id, pending.target_postbox_id)
self.assertIsNotNone(pending.target_message_id)
self.assertEqual(3, session.query(PostboxMessage).count())
self.assertIn(
"postbox.escalation.due.v1",
[event.type for event in events],
)
escalation_notification = next(
request
for request, _enqueue_delivery in notifications.requests
if request.event_kind == "postbox.escalation.due.v1"
)
self.assertEqual("account-top", escalation_notification.recipient_id)
self.assertEqual(4, escalation_notification.priority)
self.assertNotIn("Escalated decision", escalation_notification.body_text)
def test_vacancy_escalation_stops_when_previous_function_is_filled(
self,