Compare commits

6 Commits

16 changed files with 1198 additions and 87 deletions

View File

@@ -93,12 +93,24 @@ Commands and events are separate concepts:
be written to the audit outbox before delivery.
`govoplan_audit.backend.outbox.SqlAuditOutbox` persists platform events in
`audit_outbox_events`. Dispatchers can later call
`dispatch_pending_platform_events()` to publish pending events and record retry
state. The outbox payload stores the full governed event envelope:
`audit_outbox_events` and one durable state row per allowlisted consumer in
`audit_outbox_deliveries`. Dispatchers supply stable consumer IDs and
idempotent handlers. Consumer work and its delivered marker share one database
transaction; retries reuse the stable `<event-id>:<consumer-id>` delivery key.
Bounded failures are quarantined instead of retried forever. The outbox payload
stores the full governed event envelope:
correlation/causation ids, actor, tenant, subject, resource, classification,
module id, event id, type, and payload.
Public and internal events may use an allowlisted subscription directly.
Confidential and restricted subscriptions additionally require a persisted
policy-decision reference. Operators can inspect delivery metrics at
`GET /api/v1/admin/audit/event-delivery/metrics` and replay a retrying or
quarantined delivery with a reason through
`POST /api/v1/admin/audit/event-deliveries/{event_id}/{consumer_id}/replay`.
Replay itself is written to the audit log. Successful envelopes are subject to
configured retention; quarantined evidence is not removed automatically.
Application code should enqueue or publish facts only after the state change
they describe is known. Long-running operators and installers should model
requested work as commands first, then emit facts as events as each step

View File

@@ -18,7 +18,7 @@
"LICENSE"
],
"peerDependencies": {
"@govoplan/core-webui": "^0.1.8",
"@govoplan/core-webui": "^0.1.9",
"lucide-react": "^1.23.0",
"react": "^19.0.0",
"react-dom": "^19.0.0",

View File

@@ -10,7 +10,12 @@ from sqlalchemy.orm import Session
from govoplan_core.auth import ApiPrincipal, has_scope, require_any_scope
from govoplan_audit.backend.db.models import AuditLog
from govoplan_core.audit.logging import AUDIT_MODULE_ID, AUDIT_SYSTEM_EVENTS_COLLECTION, AUDIT_TENANT_EVENTS_COLLECTION
from govoplan_core.audit.logging import (
AUDIT_MODULE_ID,
AUDIT_SYSTEM_EVENTS_COLLECTION,
AUDIT_TENANT_EVENTS_COLLECTION,
audit_from_principal,
)
from govoplan_core.core.access import CAPABILITY_ACCESS_ADMINISTRATION, AccessAdministration
from govoplan_core.core.change_sequence import decode_sequence_watermark, encode_sequence_watermark, max_sequence_id, sequence_entries_since, sequence_watermark_is_expired
from govoplan_core.core.pagination import KeysetCursorError, decode_keyset_cursor, encode_keyset_cursor, keyset_query_fingerprint
@@ -18,7 +23,18 @@ from govoplan_core.core.runtime import get_registry
from govoplan_core.db.session import get_session
from govoplan_core.tenancy.scope import Tenant
from .schemas import AuditAdminDeltaResponse, AuditAdminItem, AuditAdminListResponse, AuditLogItemResponse, AuditLogListResponse
from govoplan_core.core.events import platform_event_outbox
from .schemas import (
AuditAdminDeltaResponse,
AuditAdminItem,
AuditAdminListResponse,
AuditLogItemResponse,
AuditLogListResponse,
EventDeliveryMetricsResponse,
EventDeliveryReplayRequest,
EventDeliveryReplayResponse,
)
router = APIRouter(tags=["audit"])
@@ -196,6 +212,36 @@ def _audit_delta_response_watermark(
return encode_sequence_watermark(entries[-1].id) if has_more and entries else _audit_delta_watermark(session, effective_scope=effective_scope, tenant_id=tenant_id)
def _full_audit_delta_response(
session: Session,
*,
context: AuditAdminQueryContext,
page_query: Any,
start_cursor: str | None,
sort_by: str,
sort_direction: str,
) -> AuditAdminDeltaResponse:
rows_plus_one = page_query.order_by(context.order, AuditLog.id.desc()).limit(context.effective_page_size + 1).all()
rows = rows_plus_one[:context.effective_page_size]
next_cursor = (
_audit_cursor_for_row(rows[-1], sort_by=sort_by, sort_direction=sort_direction, fingerprint=context.fingerprint)
if len(rows_plus_one) > context.effective_page_size and rows else None
)
return AuditAdminDeltaResponse(
total=context.total,
page=1,
page_size=context.effective_page_size,
pages=context.pages,
cursor=start_cursor,
next_cursor=next_cursor,
items=_audit_items(session, rows, context.access_admin),
deleted=[],
watermark=_audit_delta_watermark(session, effective_scope=context.effective_scope, tenant_id=context.resolved_tenant_id),
has_more=False,
full=True,
)
def _audit_items(session: Session, rows: list[AuditLog], access_admin: AccessAdministration) -> list[AuditAdminItem]:
actor_email_by_user_id = access_admin.actor_email_by_user_id(session, {row.user_id for row in rows if row.user_id})
return [
@@ -530,24 +576,13 @@ def list_admin_audit_delta(
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=str(exc)) from exc
if since is None:
rows_plus_one = page_query.order_by(context.order, AuditLog.id.desc()).limit(context.effective_page_size + 1).all()
rows = rows_plus_one[:context.effective_page_size]
next_cursor = (
_audit_cursor_for_row(rows[-1], sort_by=sort_by, sort_direction=sort_direction, fingerprint=context.fingerprint)
if len(rows_plus_one) > context.effective_page_size and rows else None
)
return AuditAdminDeltaResponse(
total=context.total,
page=1,
page_size=context.effective_page_size,
pages=context.pages,
cursor=start_cursor,
next_cursor=next_cursor,
items=_audit_items(session, rows, context.access_admin),
deleted=[],
watermark=_audit_delta_watermark(session, effective_scope=context.effective_scope, tenant_id=context.resolved_tenant_id),
has_more=False,
full=True,
return _full_audit_delta_response(
session,
context=context,
page_query=page_query,
start_cursor=start_cursor,
sort_by=sort_by,
sort_direction=sort_direction,
)
entries, has_more = _audit_delta_entries(
@@ -558,24 +593,13 @@ def list_admin_audit_delta(
limit=context.effective_page_size,
)
if entries is None:
rows_plus_one = page_query.order_by(context.order, AuditLog.id.desc()).limit(context.effective_page_size + 1).all()
rows = rows_plus_one[:context.effective_page_size]
next_cursor = (
_audit_cursor_for_row(rows[-1], sort_by=sort_by, sort_direction=sort_direction, fingerprint=context.fingerprint)
if len(rows_plus_one) > context.effective_page_size and rows else None
)
return AuditAdminDeltaResponse(
total=context.total,
page=1,
page_size=context.effective_page_size,
pages=context.pages,
cursor=start_cursor,
next_cursor=next_cursor,
items=_audit_items(session, rows, context.access_admin),
deleted=[],
watermark=_audit_delta_watermark(session, effective_scope=context.effective_scope, tenant_id=context.resolved_tenant_id),
has_more=False,
full=True,
return _full_audit_delta_response(
session,
context=context,
page_query=page_query,
start_cursor=start_cursor,
sort_by=sort_by,
sort_direction=sort_direction,
)
changed_ids = [entry.resource_id for entry in entries if entry.resource_type == "audit_log"]
@@ -625,3 +649,78 @@ def list_audit_log(
query = query.filter(AuditLog.object_id == object_id)
items = query.order_by(AuditLog.created_at.desc()).offset(offset).limit(limit).all()
return AuditLogListResponse(items=[AuditLogItemResponse.model_validate(item) for item in items])
@router.get(
"/admin/audit/event-delivery/metrics",
response_model=EventDeliveryMetricsResponse,
)
def event_delivery_metrics(
session: Session = Depends(get_session),
_principal: ApiPrincipal = Depends(
require_any_scope("system:audit:read")
),
):
outbox = platform_event_outbox(get_registry())
if outbox is None:
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Durable platform event delivery is not configured",
)
return EventDeliveryMetricsResponse.model_validate(
outbox.delivery_metrics(session)
)
@router.post(
"/admin/audit/event-deliveries/{event_id}/{consumer_id}/replay",
response_model=EventDeliveryReplayResponse,
)
def replay_event_delivery(
event_id: str,
consumer_id: str,
payload: EventDeliveryReplayRequest,
session: Session = Depends(get_session),
principal: ApiPrincipal = Depends(
require_any_scope("system:governance:write")
),
):
outbox = platform_event_outbox(get_registry())
if outbox is None:
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail="Durable platform event delivery is not configured",
)
try:
result = outbox.replay_delivery(
session,
event_id=event_id,
consumer_id=consumer_id,
operator_id=principal.account_id,
reason=payload.reason,
)
except LookupError as exc:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND,
detail=str(exc),
) from exc
except ValueError as exc:
raise HTTPException(
status_code=status.HTTP_409_CONFLICT,
detail=str(exc),
) from exc
audit_from_principal(
session,
principal,
action="platform_event.delivery_replayed",
scope="system",
object_type="platform_event_delivery",
object_id=f"{event_id}:{consumer_id}",
details={
"event_id": event_id,
"consumer_id": consumer_id,
"reason": payload.reason,
},
)
session.commit()
return EventDeliveryReplayResponse.model_validate(result)

View File

@@ -53,3 +53,29 @@ class AuditLogItemResponse(BaseModel):
class AuditLogListResponse(BaseModel):
items: list[AuditLogItemResponse]
class EventDeliveryMetricsResponse(BaseModel):
events: dict[str, int] = Field(default_factory=dict)
deliveries: dict[str, int] = Field(default_factory=dict)
consumers: dict[str, dict[str, int]] = Field(default_factory=dict)
oldest_due_at: datetime | None = None
class EventDeliveryReplayRequest(BaseModel):
model_config = ConfigDict(extra="forbid")
reason: str = Field(min_length=1, max_length=2000)
class EventDeliveryReplayResponse(BaseModel):
event_id: str
consumer_id: str
delivery_key: str
status: str
attempts: int
replay_count: int
last_replayed_at: datetime | None = None
last_replayed_by: str | None = None
last_replay_reason: str | None = None
last_error: str | None = None

View File

@@ -57,4 +57,101 @@ class AuditOutboxEvent(Base, TimestampMixin):
last_error: Mapped[str | None] = mapped_column(Text, nullable=True)
__all__ = ["AuditLog", "AuditOutboxEvent", "new_uuid"]
class AuditOutboxDelivery(Base, TimestampMixin):
__tablename__ = "audit_outbox_deliveries"
__table_args__ = (
UniqueConstraint(
"outbox_event_id",
"consumer_id",
name="uq_audit_outbox_delivery_consumer",
),
UniqueConstraint(
"delivery_key",
name="uq_audit_outbox_delivery_key",
),
Index(
"ix_audit_outbox_delivery_status_next_attempt_at",
"status",
"next_attempt_at",
),
Index(
"ix_audit_outbox_delivery_consumer_status",
"consumer_id",
"status",
),
)
id: Mapped[str] = mapped_column(
String(36),
primary_key=True,
default=new_uuid,
)
outbox_event_id: Mapped[str] = mapped_column(
ForeignKey("audit_outbox_events.id", ondelete="CASCADE"),
nullable=False,
index=True,
)
consumer_id: Mapped[str] = mapped_column(
String(128),
nullable=False,
)
delivery_key: Mapped[str] = mapped_column(
String(300),
nullable=False,
)
policy_decision_ref: Mapped[str | None] = mapped_column(
String(128),
nullable=True,
)
status: Mapped[str] = mapped_column(
String(20),
nullable=False,
default="pending",
index=True,
)
attempts: Mapped[int] = mapped_column(
Integer,
nullable=False,
default=0,
)
next_attempt_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True),
nullable=True,
)
delivered_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True),
nullable=True,
)
quarantined_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True),
nullable=True,
)
replay_count: Mapped[int] = mapped_column(
Integer,
nullable=False,
default=0,
)
last_replayed_at: Mapped[datetime | None] = mapped_column(
DateTime(timezone=True),
nullable=True,
)
last_replayed_by: Mapped[str | None] = mapped_column(
String(128),
nullable=True,
)
last_replay_reason: Mapped[str | None] = mapped_column(
Text,
nullable=True,
)
last_error: Mapped[str | None] = mapped_column(
Text,
nullable=True,
)
__all__ = [
"AuditLog",
"AuditOutboxDelivery",
"AuditOutboxEvent",
"new_uuid",
]

View File

@@ -1,5 +1,7 @@
from __future__ import annotations
from pathlib import Path
from govoplan_audit.backend.db import models as audit_models # noqa: F401 - populate Audit ORM metadata
from govoplan_core.core.access import (
CAPABILITY_AUDIT_RECORDER,
@@ -9,6 +11,8 @@ from govoplan_core.core.access import (
)
from govoplan_core.core.module_guards import drop_table_retirement_provider, persistent_table_uninstall_guard
from govoplan_core.core.modules import FrontendModule, MigrationSpec, ModuleContext, ModuleManifest
from govoplan_core.core.events import CAPABILITY_PLATFORM_EVENT_OUTBOX
from govoplan_core.core.views import ViewSurface
from govoplan_core.db.base import Base
@@ -33,6 +37,18 @@ def _audit_retention(context: ModuleContext):
return SqlAuditRetentionProvider()
def _event_outbox(context: ModuleContext):
from govoplan_audit.backend.outbox import SqlAuditOutbox
return SqlAuditOutbox(
max_attempts=getattr(
context.settings,
"platform_event_outbox_max_attempts",
8,
)
)
manifest = ModuleManifest(
id="audit",
name="Audit",
@@ -42,20 +58,38 @@ manifest = ModuleManifest(
frontend=FrontendModule(
module_id="audit",
package_name="@govoplan/audit-webui",
view_surfaces=(
ViewSurface(id="audit.admin.system", module_id="audit", kind="section", label="System audit", order=90),
ViewSurface(id="audit.admin.tenant", module_id="audit", kind="section", label="Tenant audit", order=100),
),
),
migration_spec=MigrationSpec(
module_id="audit",
metadata=Base.metadata,
script_location=str(
Path(__file__).with_name("migrations") / "versions"
),
retirement_supported=True,
retirement_provider=drop_table_retirement_provider(audit_models.AuditLog, audit_models.AuditOutboxEvent, label="Audit"),
retirement_provider=drop_table_retirement_provider(
audit_models.AuditLog,
audit_models.AuditOutboxDelivery,
audit_models.AuditOutboxEvent,
label="Audit",
),
retirement_notes="Destructive retirement drops audit-owned database tables after the installer captures a database snapshot.",
),
uninstall_guard_providers=(
persistent_table_uninstall_guard(audit_models.AuditLog, audit_models.AuditOutboxEvent, label="Audit"),
persistent_table_uninstall_guard(
audit_models.AuditLog,
audit_models.AuditOutboxDelivery,
audit_models.AuditOutboxEvent,
label="Audit",
),
),
capability_factories={
CAPABILITY_AUDIT_RECORDER: _audit_recorder,
CAPABILITY_AUDIT_RETENTION: _audit_retention,
CAPABILITY_PLATFORM_EVENT_OUTBOX: _event_outbox,
},
)

View File

@@ -0,0 +1 @@
"""Audit module database migrations."""

View File

@@ -0,0 +1 @@
"""Development-track Audit migrations."""

View File

@@ -0,0 +1,92 @@
"""durable platform event delivery ledger
Revision ID: a8d1e4f7b2c5
Revises: None
Create Date: 2026-07-29 00:00:00.000000
"""
from __future__ import annotations
from alembic import op
import sqlalchemy as sa
revision = "a8d1e4f7b2c5"
down_revision = None
branch_labels = None
depends_on = "c91f0a72be34"
def upgrade() -> None:
inspector = sa.inspect(op.get_bind())
if "audit_outbox_deliveries" in inspector.get_table_names():
return
op.create_table(
"audit_outbox_deliveries",
sa.Column("id", sa.String(length=36), nullable=False),
sa.Column("outbox_event_id", sa.String(length=36), nullable=False),
sa.Column("consumer_id", sa.String(length=128), nullable=False),
sa.Column("delivery_key", sa.String(length=300), nullable=False),
sa.Column("policy_decision_ref", sa.String(length=128), nullable=True),
sa.Column("status", sa.String(length=20), nullable=False),
sa.Column("attempts", sa.Integer(), nullable=False),
sa.Column("next_attempt_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("delivered_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("quarantined_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("replay_count", sa.Integer(), nullable=False),
sa.Column("last_replayed_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("last_replayed_by", sa.String(length=128), nullable=True),
sa.Column("last_replay_reason", sa.Text(), nullable=True),
sa.Column("last_error", sa.Text(), nullable=True),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
sa.ForeignKeyConstraint(
["outbox_event_id"],
["audit_outbox_events.id"],
name=op.f(
"fk_audit_outbox_deliveries_outbox_event_id_"
"audit_outbox_events"
),
ondelete="CASCADE",
),
sa.PrimaryKeyConstraint(
"id",
name=op.f("pk_audit_outbox_deliveries"),
),
sa.UniqueConstraint(
"delivery_key",
name="uq_audit_outbox_delivery_key",
),
sa.UniqueConstraint(
"outbox_event_id",
"consumer_id",
name="uq_audit_outbox_delivery_consumer",
),
)
op.create_index(
"ix_audit_outbox_deliveries_outbox_event_id",
"audit_outbox_deliveries",
["outbox_event_id"],
)
op.create_index(
"ix_audit_outbox_deliveries_status",
"audit_outbox_deliveries",
["status"],
)
op.create_index(
"ix_audit_outbox_delivery_status_next_attempt_at",
"audit_outbox_deliveries",
["status", "next_attempt_at"],
)
op.create_index(
"ix_audit_outbox_delivery_consumer_status",
"audit_outbox_deliveries",
["consumer_id", "status"],
)
def downgrade() -> None:
if (
"audit_outbox_deliveries"
in sa.inspect(op.get_bind()).get_table_names()
):
op.drop_table("audit_outbox_deliveries")

View File

@@ -0,0 +1 @@
"""Release-track Audit migrations."""

View File

@@ -0,0 +1,92 @@
"""durable platform event delivery ledger
Revision ID: a8d1e4f7b2c5
Revises: None
Create Date: 2026-07-29 00:00:00.000000
"""
from __future__ import annotations
from alembic import op
import sqlalchemy as sa
revision = "a8d1e4f7b2c5"
down_revision = None
branch_labels = None
depends_on = "c91f0a72be34"
def upgrade() -> None:
inspector = sa.inspect(op.get_bind())
if "audit_outbox_deliveries" in inspector.get_table_names():
return
op.create_table(
"audit_outbox_deliveries",
sa.Column("id", sa.String(length=36), nullable=False),
sa.Column("outbox_event_id", sa.String(length=36), nullable=False),
sa.Column("consumer_id", sa.String(length=128), nullable=False),
sa.Column("delivery_key", sa.String(length=300), nullable=False),
sa.Column("policy_decision_ref", sa.String(length=128), nullable=True),
sa.Column("status", sa.String(length=20), nullable=False),
sa.Column("attempts", sa.Integer(), nullable=False),
sa.Column("next_attempt_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("delivered_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("quarantined_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("replay_count", sa.Integer(), nullable=False),
sa.Column("last_replayed_at", sa.DateTime(timezone=True), nullable=True),
sa.Column("last_replayed_by", sa.String(length=128), nullable=True),
sa.Column("last_replay_reason", sa.Text(), nullable=True),
sa.Column("last_error", sa.Text(), nullable=True),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False),
sa.ForeignKeyConstraint(
["outbox_event_id"],
["audit_outbox_events.id"],
name=op.f(
"fk_audit_outbox_deliveries_outbox_event_id_"
"audit_outbox_events"
),
ondelete="CASCADE",
),
sa.PrimaryKeyConstraint(
"id",
name=op.f("pk_audit_outbox_deliveries"),
),
sa.UniqueConstraint(
"delivery_key",
name="uq_audit_outbox_delivery_key",
),
sa.UniqueConstraint(
"outbox_event_id",
"consumer_id",
name="uq_audit_outbox_delivery_consumer",
),
)
op.create_index(
"ix_audit_outbox_deliveries_outbox_event_id",
"audit_outbox_deliveries",
["outbox_event_id"],
)
op.create_index(
"ix_audit_outbox_deliveries_status",
"audit_outbox_deliveries",
["status"],
)
op.create_index(
"ix_audit_outbox_delivery_status_next_attempt_at",
"audit_outbox_deliveries",
["status", "next_attempt_at"],
)
op.create_index(
"ix_audit_outbox_delivery_consumer_status",
"audit_outbox_deliveries",
["consumer_id", "status"],
)
def downgrade() -> None:
if (
"audit_outbox_deliveries"
in sa.inspect(op.get_bind()).get_table_names()
):
op.drop_table("audit_outbox_deliveries")

View File

@@ -1,14 +1,18 @@
from __future__ import annotations
from collections.abc import Callable, Mapping
from collections.abc import Callable, Mapping, Sequence
from datetime import datetime, timedelta, timezone
from typing import Any, cast
from sqlalchemy import or_
from sqlalchemy import delete, func, or_, select
from sqlalchemy.orm import Session
from govoplan_audit.backend.db.models import AuditOutboxEvent
from govoplan_audit.backend.db.models import (
AuditOutboxDelivery,
AuditOutboxEvent,
)
from govoplan_core.core.events import (
DurableEventConsumer,
EventActorRef,
EventClassification,
EventObjectRef,
@@ -22,9 +26,23 @@ EventDispatcher = Callable[[PlatformEvent], None]
class SqlAuditOutbox:
def __init__(self, *, max_attempts: int = 8) -> None:
self._max_attempts = max(1, min(int(max_attempts), 100))
def enqueue(self, session: object, event: PlatformEvent) -> AuditOutboxEvent:
db = _session(session)
traced = ensure_event_trace(event)
existing = db.scalar(
select(AuditOutboxEvent).where(
AuditOutboxEvent.event_id == traced.event_id
)
)
if existing is not None:
if existing.payload != traced.to_dict():
raise ValueError(
"A different platform event already uses this event id"
)
return existing
item = AuditOutboxEvent(
event_id=traced.event_id,
event_type=traced.type,
@@ -36,48 +54,470 @@ class SqlAuditOutbox:
status="pending",
)
db.add(item)
db.flush()
return item
def dispatch_pending(
self,
session: object,
*,
dispatcher: EventDispatcher = publish_platform_event,
consumers: Sequence[DurableEventConsumer] = (),
observer: EventDispatcher | None = publish_platform_event,
dispatcher: EventDispatcher | None = None,
limit: int = 100,
) -> dict[str, int]:
db = _session(session)
now = datetime.now(timezone.utc)
consumers_by_id = _consumer_map(consumers)
rows = (
db.query(AuditOutboxEvent)
.filter(
AuditOutboxEvent.status.in_(("pending", "failed")),
AuditOutboxEvent.status.in_(
("pending", "failed", "retrying")
),
or_(AuditOutboxEvent.next_attempt_at.is_(None), AuditOutboxEvent.next_attempt_at <= now),
)
.order_by(AuditOutboxEvent.created_at.asc(), AuditOutboxEvent.id.asc())
.with_for_update(skip_locked=True)
.limit(max(1, min(int(limit), 500)))
.all()
)
counts = {"selected": len(rows), "dispatched": 0, "failed": 0}
counts = {
"selected": len(rows),
"delivered": 0,
"retrying": 0,
"quarantined": 0,
"dispatched": 0,
"observer_failed": 0,
}
effective_observer = dispatcher or observer
for row in rows:
try:
dispatcher(_event_from_payload(row.payload))
except Exception as exc: # noqa: BLE001 - dispatcher errors must be retained for retry/diagnostics.
row.status = "failed"
row.attempts += 1
row.last_error = str(exc)
row.next_attempt_at = now + _retry_delay(row.attempts)
counts["failed"] += 1
continue
row.status = "dispatched"
row.attempts += 1
row.dispatched_at = now
row.next_attempt_at = None
row.last_error = None
counts["dispatched"] += 1
event = _event_from_payload(row.payload)
deliveries = _event_deliveries(
db,
row=row,
event=event,
consumers=consumers_by_id.values(),
)
_dispatch_event_deliveries(
event,
deliveries=deliveries,
consumers_by_id=consumers_by_id,
now=now,
max_attempts=self._max_attempts,
counts=counts,
)
_finish_event_dispatch(
row,
deliveries=deliveries,
observer=effective_observer,
event=event,
now=now,
counts=counts,
)
db.flush()
return counts
def replay_delivery(
self,
session: object,
*,
event_id: str,
consumer_id: str,
operator_id: str,
reason: str,
) -> dict[str, object]:
db = _session(session)
clean_reason = reason.strip()
clean_operator_id = operator_id.strip()
if not clean_reason or len(clean_reason) > 2000:
raise ValueError(
"Replay reason must contain between 1 and 2000 characters"
)
if not clean_operator_id or len(clean_operator_id) > 128:
raise ValueError("Replay operator id is invalid")
row = db.scalar(
select(AuditOutboxDelivery)
.join(
AuditOutboxEvent,
AuditOutboxEvent.id
== AuditOutboxDelivery.outbox_event_id,
)
.where(
AuditOutboxEvent.event_id == event_id,
AuditOutboxDelivery.consumer_id == consumer_id,
)
.with_for_update()
)
if row is None:
raise LookupError("Platform event delivery was not found")
if row.status not in {"retrying", "quarantined"}:
raise ValueError(
"Only retrying or quarantined deliveries can be replayed"
)
now = datetime.now(timezone.utc)
row.status = "pending"
row.attempts = 0
row.next_attempt_at = now
row.quarantined_at = None
row.last_error = None
row.replay_count += 1
row.last_replayed_at = now
row.last_replayed_by = clean_operator_id
row.last_replay_reason = clean_reason
event_row = db.get(AuditOutboxEvent, row.outbox_event_id)
if event_row is None:
raise LookupError("Platform event envelope was not found")
event_row.status = "pending"
event_row.next_attempt_at = now
event_row.last_error = None
event_row.dispatched_at = None
db.flush()
return _delivery_state(row, event_id=event_row.event_id)
def purge_terminal(
self,
session: object,
*,
before: datetime,
limit: int = 500,
) -> dict[str, int]:
db = _session(session)
ids = tuple(
db.scalars(
select(AuditOutboxEvent.id)
.where(
AuditOutboxEvent.status == "dispatched",
AuditOutboxEvent.dispatched_at.is_not(None),
AuditOutboxEvent.dispatched_at < before,
)
.order_by(
AuditOutboxEvent.dispatched_at,
AuditOutboxEvent.id,
)
.limit(max(1, min(int(limit), 5000)))
)
)
if ids:
db.execute(
delete(AuditOutboxEvent).where(
AuditOutboxEvent.id.in_(ids)
)
)
db.flush()
return {"deleted": len(ids)}
def delivery_metrics(
self,
session: object,
) -> dict[str, object]:
db = _session(session)
event_counts = {
str(status): int(count)
for status, count in db.execute(
select(
AuditOutboxEvent.status,
func.count(AuditOutboxEvent.id),
).group_by(AuditOutboxEvent.status)
)
}
delivery_counts = {
str(status): int(count)
for status, count in db.execute(
select(
AuditOutboxDelivery.status,
func.count(AuditOutboxDelivery.id),
).group_by(AuditOutboxDelivery.status)
)
}
consumer_counts = {
str(consumer_id): {
str(status): int(count)
for status, count in values
}
for consumer_id, values in _consumer_delivery_counts(db).items()
}
oldest_due = db.scalar(
select(func.min(AuditOutboxDelivery.created_at)).where(
AuditOutboxDelivery.status.in_(
("pending", "retrying")
)
)
)
return {
"events": event_counts,
"deliveries": delivery_counts,
"consumers": consumer_counts,
"oldest_due_at": (
oldest_due.isoformat()
if isinstance(oldest_due, datetime)
else None
),
}
def _consumer_map(
consumers: Sequence[DurableEventConsumer],
) -> dict[str, DurableEventConsumer]:
result: dict[str, DurableEventConsumer] = {}
for consumer in consumers:
if consumer.consumer_id in result:
raise ValueError(
f"Duplicate durable event consumer: {consumer.consumer_id}"
)
result[consumer.consumer_id] = consumer
return result
def _event_deliveries(
session: Session,
*,
row: AuditOutboxEvent,
event: PlatformEvent,
consumers: Sequence[DurableEventConsumer],
) -> list[AuditOutboxDelivery]:
existing = {
delivery.consumer_id: delivery
for delivery in session.scalars(
select(AuditOutboxDelivery).where(
AuditOutboxDelivery.outbox_event_id == row.id
)
)
}
for consumer in consumers:
if (
consumer.consumer_id in existing
or not consumer.accepts(event)
):
continue
delivery = AuditOutboxDelivery(
outbox_event_id=row.id,
consumer_id=consumer.consumer_id,
delivery_key=consumer.delivery_key(event),
policy_decision_ref=consumer.policy_decision_ref,
status="pending",
)
session.add(delivery)
existing[consumer.consumer_id] = delivery
session.flush()
return sorted(
existing.values(),
key=lambda item: (item.created_at, item.id),
)
def _dispatch_event_deliveries(
event: PlatformEvent,
*,
deliveries: Sequence[AuditOutboxDelivery],
consumers_by_id: Mapping[str, DurableEventConsumer],
now: datetime,
max_attempts: int,
counts: dict[str, int],
) -> None:
for delivery in deliveries:
if not _delivery_is_due(delivery, now=now):
continue
consumer = consumers_by_id.get(delivery.consumer_id)
if consumer is None:
_record_delivery_failure(
delivery,
error="Durable event consumer is not registered",
now=now,
max_attempts=max_attempts,
counts=counts,
)
continue
if not consumer.accepts(event):
_quarantine_delivery(
delivery,
error=(
"The current durable subscription no longer permits "
"this event"
),
now=now,
counts=counts,
)
continue
if (
event.classification in {"confidential", "restricted"}
and delivery.policy_decision_ref
!= consumer.policy_decision_ref
):
_quarantine_delivery(
delivery,
error=(
"The policy decision for this classified event "
"subscription changed"
),
now=now,
counts=counts,
)
continue
try:
consumer.handler(event, delivery.delivery_key)
except Exception as exc: # noqa: BLE001 - failures must be persisted.
_record_delivery_failure(
delivery,
error=str(exc),
now=now,
max_attempts=max_attempts,
counts=counts,
)
continue
delivery.status = "delivered"
delivery.attempts += 1
delivery.delivered_at = now
delivery.next_attempt_at = None
delivery.quarantined_at = None
delivery.last_error = None
counts["delivered"] += 1
def _delivery_is_due(
delivery: AuditOutboxDelivery,
*,
now: datetime,
) -> bool:
if delivery.status not in {"pending", "retrying"}:
return False
if delivery.next_attempt_at is None:
return True
return _as_utc(delivery.next_attempt_at) <= now
def _record_delivery_failure(
delivery: AuditOutboxDelivery,
*,
error: str,
now: datetime,
max_attempts: int,
counts: dict[str, int],
) -> None:
delivery.attempts += 1
delivery.last_error = _bounded_error(error)
if delivery.attempts >= max_attempts:
_quarantine_delivery(
delivery,
error=delivery.last_error,
now=now,
counts=counts,
)
return
delivery.status = "retrying"
delivery.next_attempt_at = now + _retry_delay(delivery.attempts)
counts["retrying"] += 1
def _quarantine_delivery(
delivery: AuditOutboxDelivery,
*,
error: str,
now: datetime,
counts: dict[str, int],
) -> None:
delivery.status = "quarantined"
delivery.quarantined_at = now
delivery.next_attempt_at = None
delivery.last_error = _bounded_error(error)
counts["quarantined"] += 1
def _finish_event_dispatch(
row: AuditOutboxEvent,
*,
deliveries: Sequence[AuditOutboxDelivery],
observer: EventDispatcher | None,
event: PlatformEvent,
now: datetime,
counts: dict[str, int],
) -> None:
row.attempts += 1
quarantined = [
item for item in deliveries
if item.status == "quarantined"
]
outstanding = [
item for item in deliveries
if item.status in {"pending", "retrying"}
]
if quarantined:
row.status = "quarantined"
row.next_attempt_at = None
row.last_error = quarantined[0].last_error
return
if outstanding:
row.status = "retrying"
due_times = [
item.next_attempt_at
for item in outstanding
if item.next_attempt_at is not None
]
row.next_attempt_at = min(due_times) if due_times else now
row.last_error = next(
(
item.last_error
for item in outstanding
if item.last_error
),
None,
)
return
if observer is not None:
try:
observer(event)
except Exception as exc: # noqa: BLE001 - observers are non-durable.
counts["observer_failed"] += 1
row.last_error = _bounded_error(
f"Non-durable observer failed: {exc}"
)
else:
row.last_error = None
else:
row.last_error = None
row.status = "dispatched"
row.dispatched_at = now
row.next_attempt_at = None
counts["dispatched"] += 1
def _delivery_state(
delivery: AuditOutboxDelivery,
*,
event_id: str,
) -> dict[str, object]:
return {
"event_id": event_id,
"consumer_id": delivery.consumer_id,
"delivery_key": delivery.delivery_key,
"status": delivery.status,
"attempts": delivery.attempts,
"replay_count": delivery.replay_count,
"last_replayed_at": delivery.last_replayed_at,
"last_replayed_by": delivery.last_replayed_by,
"last_replay_reason": delivery.last_replay_reason,
"last_error": delivery.last_error,
}
def _consumer_delivery_counts(
session: Session,
) -> dict[str, list[tuple[str, int]]]:
result: dict[str, list[tuple[str, int]]] = {}
for consumer_id, status, count in session.execute(
select(
AuditOutboxDelivery.consumer_id,
AuditOutboxDelivery.status,
func.count(AuditOutboxDelivery.id),
).group_by(
AuditOutboxDelivery.consumer_id,
AuditOutboxDelivery.status,
)
):
result.setdefault(str(consumer_id), []).append(
(str(status), int(count))
)
return result
def enqueue_platform_event(session: object, event: PlatformEvent) -> AuditOutboxEvent:
return SqlAuditOutbox().enqueue(session, event)
@@ -97,6 +537,17 @@ def _retry_delay(attempts: int) -> timedelta:
return timedelta(seconds=seconds)
def _as_utc(value: datetime) -> datetime:
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
def _bounded_error(value: str) -> str:
clean = value.strip() or "Unknown durable event delivery failure"
return clean[:4000]
def _event_from_payload(payload: Mapping[str, Any]) -> PlatformEvent:
return PlatformEvent(
type=str(payload["type"]),

View File

@@ -1,14 +1,22 @@
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 AuditOutboxEvent
from govoplan_audit.backend.db.models import (
AuditOutboxDelivery,
AuditOutboxEvent,
)
from govoplan_audit.backend.outbox import SqlAuditOutbox
from govoplan_core.core.events import EventActorRef, PlatformEvent
from govoplan_core.core.events import (
DurableEventConsumer,
EventActorRef,
PlatformEvent,
)
from govoplan_core.db.base import Base
@@ -30,13 +38,23 @@ class AuditCommandBusTests(unittest.TestCase):
class AuditOutboxTests(unittest.TestCase):
def test_outbox_enqueues_governed_event_and_dispatches_pending_rows(self) -> None:
def _database(self):
engine = create_engine("sqlite:///:memory:")
self.addCleanup(engine.dispose)
Base.metadata.create_all(bind=engine, tables=[AuditOutboxEvent.__table__])
Session = sessionmaker(bind=engine)
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(
@@ -53,32 +71,213 @@ class AuditOutboxTests(unittest.TestCase):
self.assertEqual(event.event_id, row.correlation_id)
self.assertEqual("user-1", row.payload["actor"]["id"])
counts = outbox.dispatch_pending(session, dispatcher=seen.append)
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, "dispatched": 1, "failed": 0}, counts)
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_records_failed_dispatch_for_retry(self) -> None:
engine = create_engine("sqlite:///:memory:")
self.addCleanup(engine.dispose)
Base.metadata.create_all(bind=engine, tables=[AuditOutboxEvent.__table__])
Session = sessionmaker(bind=engine)
outbox = SqlAuditOutbox()
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, dispatcher=lambda event: (_ for _ in ()).throw(RuntimeError("offline")))
counts = outbox.dispatch_pending(
session,
consumers=(consumer,),
observer=None,
)
self.assertEqual({"selected": 1, "dispatched": 0, "failed": 1}, counts)
self.assertEqual("failed", row.status)
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__":

View File

@@ -13,11 +13,11 @@
}
},
"peerDependencies": {
"@govoplan/core-webui": "^0.1.8",
"@govoplan/core-webui": "^0.1.9",
"lucide-react": "^1.23.0",
"react": "^19.0.0",
"react-dom": "^19.0.0",
"react-router-dom": "^7.1.1"
"react-router-dom": ">=7.18.2 <8"
},
"peerDependenciesMeta": {
"@govoplan/core-webui": {

View File

@@ -1,7 +1,6 @@
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { Search } from "lucide-react";
import {
AdminIconButton,
AdminPageLayout,
adminErrorMessage,
Button,
@@ -9,6 +8,7 @@ import {
Dialog,
formatAdminDateTime as formatDateTime,
mergeDeltaRows,
TableActionGroup,
useDeltaWatermarks,
type ApiSettings,
type AuthInfo,
@@ -118,7 +118,7 @@ export default function AdminAuditPanel({ settings, auth, systemMode = false }:
{ id: "action", header: "Action", width: 250, minWidth: 170, maxWidth: 420, resizable: true, sortable: true, filterable: true, value: (row) => row.action },
{ id: "object", header: "Object", width: 300, minWidth: 180, maxWidth: 640, resizable: true, fill: true, sortable: true, filterable: true, value: (row) => `${row.object_type || "-"} ${row.object_id || ""}`.trim() },
...(systemMode ? [{ id: "tenant", header: "Tenant context", width: 190, minWidth: 150, maxWidth: 300, resizable: true, sortable: true, filterable: true, value: (row: AuditAdminItem) => row.tenant_id || "-" }] : []),
{ id: "actions", header: "Actions", width: 70, sticky: "end", resizable: false, align: "right", render: (row) => <div className="admin-icon-actions"><AdminIconButton label="Inspect audit event" icon={<Search />} onClick={() => setSelected(row)} /></div> }
{ id: "actions", header: "Actions", width: 70, sticky: "end", resizable: false, align: "right", render: (row) => <TableActionGroup actions={[{ id: "inspect", label: "Inspect audit event", icon: <Search />, onClick: () => setSelected(row) }]} /> }
], [systemMode]);
const firstShown = total === 0 ? 0 : (page - 1) * pageSize + 1;

View File

@@ -7,6 +7,7 @@ const auditAdminSections: AdminSectionsUiCapability = {
sections: [
{
id: "system-audit",
surfaceId: "audit.admin.system",
label: "Audit",
group: "SYSTEM",
order: 90,
@@ -19,6 +20,7 @@ const auditAdminSections: AdminSectionsUiCapability = {
},
{
id: "tenant-audit",
surfaceId: "audit.admin.tenant",
label: "Audit",
group: "TENANT",
order: 100,
@@ -37,6 +39,10 @@ export const auditModule: PlatformWebModule = {
label: "Audit",
version: "0.1.6",
dependencies: ["access", "admin"],
viewSurfaces: [
{ id: "audit.admin.system", moduleId: "audit", kind: "section", label: "System audit", order: 90 },
{ id: "audit.admin.tenant", moduleId: "audit", kind: "section", label: "Tenant audit", order: 100 }
],
uiCapabilities: {
"admin.sections": auditAdminSections
}