Files
govoplan-audit/tests/test_audit_delivery.py
T

344 lines
12 KiB
Python

from __future__ import annotations
import unittest
from datetime import datetime, timedelta, timezone
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
from govoplan_audit.backend.commands import AuditCommand, CommandBus
from govoplan_audit.backend.db.models import (
AuditOutboxDelivery,
AuditOutboxEvent,
)
from govoplan_audit.backend.outbox import SqlAuditOutbox
from govoplan_core.core.events import (
DurableEventConsumer,
EventActorRef,
PlatformEvent,
)
from govoplan_core.core.institutional import (
GovernedContextEnvelope,
InstitutionalReference,
TemporalRevision,
)
from govoplan_core.db.base import Base
class AuditCommandBusTests(unittest.TestCase):
def test_command_bus_dispatches_commands_separately_from_events(self) -> None:
bus = CommandBus()
seen: list[AuditCommand] = []
wildcard: list[AuditCommand] = []
bus.subscribe("retention.run", seen.append)
bus.subscribe("*", wildcard.append)
command = AuditCommand(type="retention.run", module_id="policy", payload={"dry_run": True})
bus.dispatch(command)
self.assertEqual([command], seen)
self.assertEqual([command], wildcard)
self.assertEqual("retention.run", command.to_dict()["type"])
class AuditOutboxTests(unittest.TestCase):
def _database(self):
engine = create_engine("sqlite:///:memory:")
self.addCleanup(engine.dispose)
Base.metadata.create_all(
bind=engine,
tables=[
AuditOutboxEvent.__table__,
AuditOutboxDelivery.__table__,
],
)
return sessionmaker(bind=engine)
def test_outbox_enqueues_governed_event_and_dispatches_pending_rows(self) -> None:
Session = self._database()
outbox = SqlAuditOutbox()
seen: list[PlatformEvent] = []
observed: list[PlatformEvent] = []
with Session() as session:
event = PlatformEvent(
type="tenant.created",
module_id="tenancy",
payload={"tenant_id": "tenant-1"},
actor=EventActorRef(type="user", id="user-1"),
)
row = outbox.enqueue(session, event)
self.assertEqual("pending", row.status)
self.assertEqual("tenant.created", row.event_type)
self.assertEqual(event.event_id, row.event_id)
self.assertEqual(event.event_id, row.correlation_id)
self.assertEqual("user-1", row.payload["actor"]["id"])
counts = outbox.dispatch_pending(
session,
consumers=(
DurableEventConsumer(
consumer_id="tests.consumer.v1",
event_types=frozenset({"tenant.created"}),
handler=lambda delivered, _key: seen.append(
delivered
),
),
),
observer=observed.append,
)
self.assertEqual(
{
"selected": 1,
"delivered": 1,
"retrying": 0,
"quarantined": 0,
"dispatched": 1,
"observer_failed": 0,
},
counts,
)
self.assertEqual(1, len(seen))
self.assertEqual(1, len(observed))
self.assertEqual("tenant.created", seen[0].type)
self.assertEqual("dispatched", row.status)
self.assertEqual(1, row.attempts)
self.assertIsNotNone(row.dispatched_at)
delivery = session.query(AuditOutboxDelivery).one()
self.assertEqual("delivered", delivery.status)
self.assertEqual(
f"{event.event_id}:tests.consumer.v1",
delivery.delivery_key,
)
def test_outbox_preserves_institutional_context(self) -> None:
Session = self._database()
outbox = SqlAuditOutbox()
seen: list[PlatformEvent] = []
now = datetime.now(timezone.utc)
context = GovernedContextEnvelope(
tenant_id="tenant-1",
temporal=TemporalRevision(revision="decision:7", recorded_at=now),
decision_ref=InstitutionalReference(
kind="decision",
owner_module="committee",
object_id="decision-7",
tenant_id="tenant-1",
version="7",
valid_at=now,
),
approval_refs=(
InstitutionalReference(
kind="approval",
owner_module="workflow",
object_id="approval-3",
tenant_id="tenant-1",
version="3",
valid_at=now,
),
),
)
with Session() as session:
outbox.enqueue(
session,
PlatformEvent(
type="committee.decision.recorded",
module_id="committee",
institutional_context=context,
),
)
outbox.dispatch_pending(
session,
consumers=(
DurableEventConsumer(
consumer_id="tests.institutional-context.v1",
event_types=frozenset({"committee.decision.recorded"}),
handler=lambda delivered, _key: seen.append(delivered),
),
),
)
self.assertEqual("decision-7", seen[0].institutional_context.decision_ref.object_id)
self.assertEqual(
"approval-3",
seen[0].institutional_context.approval_refs[0].object_id,
)
def test_outbox_retries_then_quarantines_a_failed_consumer(self) -> None:
Session = self._database()
outbox = SqlAuditOutbox(max_attempts=2)
consumer = DurableEventConsumer(
consumer_id="tests.failing.v1",
handler=lambda _event, _key: (_ for _ in ()).throw(
RuntimeError("offline")
),
)
with Session() as session:
row = outbox.enqueue(session, PlatformEvent(type="demo.failed", module_id="audit"))
counts = outbox.dispatch_pending(
session,
consumers=(consumer,),
observer=None,
)
self.assertEqual(1, counts["retrying"])
self.assertEqual("retrying", row.status)
self.assertEqual(1, row.attempts)
self.assertEqual("offline", row.last_error)
self.assertIsNotNone(row.next_attempt_at)
delivery = session.query(AuditOutboxDelivery).one()
delivery.next_attempt_at = None
row.next_attempt_at = None
second = outbox.dispatch_pending(
session,
consumers=(consumer,),
observer=None,
)
self.assertEqual(1, second["quarantined"])
self.assertEqual("quarantined", row.status)
self.assertEqual("quarantined", delivery.status)
self.assertIsNotNone(delivery.quarantined_at)
self.assertIsNone(delivery.next_attempt_at)
def test_replay_keeps_a_stable_delivery_key_and_runs_once(self) -> None:
Session = self._database()
outbox = SqlAuditOutbox(max_attempts=1)
event = PlatformEvent(type="demo.replay", module_id="audit")
failing = DurableEventConsumer(
consumer_id="tests.replay.v1",
handler=lambda _event, _key: (_ for _ in ()).throw(
RuntimeError("offline")
),
)
delivered: list[str] = []
with Session() as session:
outbox.enqueue(session, event)
outbox.dispatch_pending(
session,
consumers=(failing,),
observer=None,
)
state = outbox.replay_delivery(
session,
event_id=event.event_id,
consumer_id=failing.consumer_id,
operator_id="operator-1",
reason="Dependency recovered",
)
self.assertEqual("pending", state["status"])
self.assertEqual(1, state["replay_count"])
expected_key = f"{event.event_id}:{failing.consumer_id}"
self.assertEqual(expected_key, state["delivery_key"])
healthy = DurableEventConsumer(
consumer_id=failing.consumer_id,
handler=lambda _event, key: delivered.append(key),
)
outbox.dispatch_pending(
session,
consumers=(healthy,),
observer=None,
)
replay = outbox.dispatch_pending(
session,
consumers=(healthy,),
observer=None,
)
self.assertEqual([expected_key], delivered)
self.assertEqual(0, replay["selected"])
def test_classified_subscription_requires_and_persists_policy_decision(self) -> None:
with self.assertRaisesRegex(ValueError, "policy decision"):
DurableEventConsumer(
consumer_id="tests.restricted.v1",
classifications=frozenset({"restricted"}),
handler=lambda _event, _key: None,
)
Session = self._database()
outbox = SqlAuditOutbox()
consumer = DurableEventConsumer(
consumer_id="tests.restricted.v1",
classifications=frozenset({"restricted"}),
policy_decision_ref="policy-decision:42",
handler=lambda _event, _key: None,
)
with Session() as session:
outbox.enqueue(
session,
PlatformEvent(
type="case.changed",
module_id="cases",
classification="restricted",
),
)
outbox.dispatch_pending(
session,
consumers=(consumer,),
observer=None,
)
delivery = session.query(AuditOutboxDelivery).one()
self.assertEqual(
"policy-decision:42",
delivery.policy_decision_ref,
)
def test_metrics_and_retention_keep_quarantined_evidence(self) -> None:
Session = self._database()
outbox = SqlAuditOutbox(max_attempts=1)
with Session() as session:
delivered_event = PlatformEvent(
type="demo.delivered",
module_id="audit",
)
failed_event = PlatformEvent(
type="demo.failed",
module_id="audit",
)
delivered_row = outbox.enqueue(session, delivered_event)
outbox.enqueue(session, failed_event)
consumer = DurableEventConsumer(
consumer_id="tests.metrics.v1",
handler=lambda event, _key: (
(_ for _ in ()).throw(RuntimeError("offline"))
if event.type == "demo.failed"
else None
),
)
outbox.dispatch_pending(
session,
consumers=(consumer,),
observer=None,
)
delivered_row.dispatched_at = (
datetime.now(timezone.utc) - timedelta(days=100)
)
metrics = outbox.delivery_metrics(session)
purged = outbox.purge_terminal(
session,
before=datetime.now(timezone.utc)
- timedelta(days=90),
)
self.assertEqual(1, metrics["events"]["dispatched"])
self.assertEqual(1, metrics["events"]["quarantined"])
self.assertEqual(1, purged["deleted"])
remaining = session.query(AuditOutboxEvent).one()
self.assertEqual("quarantined", remaining.status)
if __name__ == "__main__":
unittest.main()