feat: atomically journal and deliver authorized native operation outcomes
Some checks failed
CI Smoke / host-smoke (push) Successful in 1s
CI Smoke / pytest-smoke (push) Failing after 4s

Assistant: codex
Assistant-Model: gpt-6-astra
Assistant-Session: 01a0e747-8f27-7242-8df8-8bc44f88c929
This commit is contained in:
tegwick 2026-09-28 14:00:26 +02:00
parent 3c5cbfbafe
commit f0eff0ac92
13 changed files with 583 additions and 28 deletions

View file

@ -68,6 +68,9 @@ def create_app(
browser = BrowserSessions(settings=browser_settings, controller=access_controller) if browser_settings else None
resolved_store = port_store or _create_store(resolved_settings)
owns_store = port_store is None
outcome_delivery = (access_controller is not None
and callable(getattr(access_controller.audit, "append_outcome", None))
and callable(getattr(resolved_store, "deliver_outcomes", None)))
resolved_repo_projection_client = repo_projection_client
owns_repo_projection_client = False
if (
@ -111,9 +114,14 @@ def create_app(
await workload_projection.refresh()
except WorkloadProjectionRejected:
pass
outcome_task = asyncio.create_task(_deliver_outcomes(resolved_store, access_controller.audit)) if outcome_delivery else None
try:
yield
finally:
if outcome_task is not None:
outcome_task.cancel()
with suppress(asyncio.CancelledError):
await outcome_task
if browser is not None:
browser.clear()
if refresh_task is not None:
@ -166,6 +174,8 @@ def create_app(
}
if resolved_settings.enforce_access:
dependency_checks["access_profile"] = "ok" if access_controller else "unavailable"
if callable(getattr(resolved_store, "deliver_outcomes", None)):
dependency_checks["outcome_delivery"] = (await resolved_store.outcome_readiness()) if outcome_delivery else "unavailable"
ready = resolved_settings.is_ready(resolved_store.backend_name) and all(
value in {"ok", "not_applicable"} for value in dependency_checks.values()
)
@ -191,6 +201,16 @@ def create_app(
return app
async def _deliver_outcomes(store, sink) -> None:
while True:
try:
await store.deliver_outcomes(sink)
except Exception:
# DB outages leave rows durable; readiness exposes missing/stale data.
pass
await asyncio.sleep(1)
async def _refresh_repository_projection(
service: RepositoryNavigationService,
interval_seconds: float,

View file

@ -1,5 +1,7 @@
from __future__ import annotations
import asyncio
import time
import hashlib
import json
from collections.abc import Mapping
@ -12,6 +14,7 @@ import sqlalchemy as sa
from sqlalchemy.ext.asyncio import AsyncEngine, async_sessionmaker, create_async_engine
from hub_core.contracts import CONTRACT_VERSION
from hub_core.security.context import current_authorization
from hub_core.runtime.models import (
EventCommand,
MessageCommand,
@ -28,6 +31,7 @@ from hub_core.runtime.tables import (
compat_api_keys,
compat_hubs,
runtime_audit_ledger,
runtime_outcome_outbox,
runtime_interaction_events,
runtime_messages,
runtime_progress_events,
@ -62,6 +66,7 @@ class PostgresPortStore:
# tables that protected traffic actually depends on.
for table in (
runtime_audit_ledger,
runtime_outcome_outbox,
compat_hubs,
compat_api_keys,
runtime_repository_navigation_state,
@ -531,19 +536,85 @@ class PostgresPortStore:
correlation_id: UUID | None,
value: Mapping[str, Any],
) -> None:
ledger_id = str(uuid4())
recorded_at = _now()
context = current_authorization.get()
detail = {"schema_version": CONTRACT_VERSION}
if context is not None:
if not context.decision_id:
raise ValueError("verified decision required for outcome attribution")
detail["authorization"] = {
"correlation_id": context.correlation_id, "issuer": context.actor.issuer,
"subject": context.actor.subject, "principal_type": context.actor.principal_type,
"actor_tenant": context.actor.tenant, "target_tenant": context.facts.target_tenant,
"decision_id": context.decision_id, "policy_version": context.policy_version,
"policy_caller": context.policy_caller, "action": context.action,
"request_digest": context.request_digest,
}
await session.execute(
runtime_audit_ledger.insert().values(
id=str(uuid4()),
id=ledger_id,
action=action,
subject_type=subject_type,
subject_id=subject_id,
correlation_id=str(correlation_id) if correlation_id else None,
payload_hash=_hash(value),
detail={"schema_version": CONTRACT_VERSION},
recorded_at=_now(),
detail=detail,
recorded_at=recorded_at,
)
)
if context is not None:
envelope = {
"id": ledger_id, "type": "hub.operation.committed", "source": "hub-core",
"subject": "hub-outcome:" + ledger_id, "tenant": context.facts.target_tenant,
"correlation_id": context.correlation_id, "occurred_at": recorded_at.isoformat(),
"data": {"outcome": "committed", "operation": action,
"subject_type": subject_type, "subject_id": subject_id,
"business_correlation_id": str(correlation_id) if correlation_id else None,
"payload_hash": _hash(value), "authorization": detail["authorization"]},
}
await session.execute(runtime_outcome_outbox.insert().values(
id=ledger_id, envelope=envelope, created_at=time.time(), attempts=0, next_attempt=0))
async def deliver_outcomes(self, sink, *, limit: int = 25) -> int:
"""At-least-once delivery; PostgreSQL workers lock distinct pending rows.
Locks span one bounded receiver call. A crash/lost receipt leaves the same
immutable envelope eligible for replay, including its idempotency key.
"""
delivered = 0
for _ in range(limit):
async with self.sessions.begin() as session:
row = (await session.execute(sa.select(runtime_outcome_outbox).where(
runtime_outcome_outbox.c.delivered_at.is_(None),
runtime_outcome_outbox.c.next_attempt <= time.time(),
).order_by(runtime_outcome_outbox.c.created_at, runtime_outcome_outbox.c.id)
.limit(1).with_for_update(skip_locked=True))).mappings().first()
if row is None:
break
try:
async with asyncio.timeout(3):
await sink.append_outcome(row["envelope"])
except Exception:
values = {"attempts": row["attempts"] + 1,
"next_attempt": time.time() + min(60, 2 ** min(row["attempts"], 6))}
else:
values = {"attempts": row["attempts"] + 1, "delivered_at": time.time()}
delivered += 1
await session.execute(runtime_outcome_outbox.update().where(
runtime_outcome_outbox.c.id == row["id"]).values(**values))
return delivered
async def outcome_readiness(self) -> str:
try:
async with self.sessions() as session:
oldest = (await session.execute(sa.select(sa.func.min(runtime_outcome_outbox.c.created_at))
.where(runtime_outcome_outbox.c.delivered_at.is_(None)))).scalar_one()
return "ok" if oldest is None or time.time() - oldest <= 60 else "stale"
except Exception:
return "unavailable"
def _record(self, kind: str, value: dict[str, Any]) -> PortRecord:
record_id = str(value.get("id") or kind)
return PortRecord(

View file

@ -235,3 +235,15 @@ runtime_import_runs = sa.Table(
sa.Column("status", sa.String(40), nullable=False),
sa.Column("created_at", sa.DateTime(timezone=True), nullable=False),
)
runtime_outcome_outbox = sa.Table(
"runtime_outcome_outbox", runtime_metadata,
sa.Column("id", sa.String(36), sa.ForeignKey("runtime_audit_ledger.id"), primary_key=True),
sa.Column("envelope", sa.JSON(), nullable=False),
sa.Column("created_at", sa.Float(), nullable=False),
sa.Column("attempts", sa.Integer(), nullable=False, default=0),
sa.Column("next_attempt", sa.Float(), nullable=False, default=0),
sa.Column("delivered_at", sa.Float(), nullable=True),
sa.Index("ix_runtime_outcome_pending", "delivered_at", "next_attempt"),
)