From f0eff0ac92dbd7d91894ac5c0af50b2f1b3725ec Mon Sep 17 00:00:00 2001 From: tegwick Date: Mon, 28 Sep 2026 14:00:26 +0200 Subject: [PATCH] feat: atomically journal and deliver authorized native operation outcomes Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e747-8f27-7242-8df8-8bc44f88c929 --- docs/access-profile-v1.md | 4 +- docs/operation-outcome-audit.md | 90 ++++++ docs/owner-access-integration.md | 5 +- .../versions/0006_outcome_outbox.py | 31 +++ hub_core/runtime/app.py | 20 ++ hub_core/runtime/postgres_store.py | 77 +++++- hub_core/runtime/tables.py | 12 + hub_core/security/audit.py | 40 +-- hub_core/security/boundary.py | 15 +- hub_core/security/context.py | 8 + tests/test_audit_core_owner_contract.py | 32 +++ tests/test_outcome_audit.py | 258 ++++++++++++++++++ ...WP-0012-netkingdom-platform-root-access.md | 19 ++ 13 files changed, 583 insertions(+), 28 deletions(-) create mode 100644 docs/operation-outcome-audit.md create mode 100644 hub_core/migrations/versions/0006_outcome_outbox.py create mode 100644 hub_core/security/context.py create mode 100644 tests/test_outcome_audit.py diff --git a/docs/access-profile-v1.md b/docs/access-profile-v1.md index 4bdec80..4f7b291 100644 --- a/docs/access-profile-v1.md +++ b/docs/access-profile-v1.md @@ -56,7 +56,9 @@ A host composes `create_app(access_controller=AccessController(...))` with: only after an operational-custody probe and explicit durable acceptance. Every allow must reach this sink before handler execution; a failed sink blocks reads as well as writes. Authorization receipts say `authorized`, not “operation completed.” Domain commit/outcome audit remains a - separate requirement; this source seam does not claim transactional audit. + separate concern. Native durable mutations now have a + [transaction-linked outcome outbox](operation-outcome-audit.md); compatibility + and external mutations remain outside that slice. - The root's existing immutable issuer and subject, supplied after owner resolution. No username, email, first-login promotion or generic role establishes root. diff --git a/docs/operation-outcome-audit.md b/docs/operation-outcome-audit.md new file mode 100644 index 0000000..754e5e4 --- /dev/null +++ b/docs/operation-outcome-audit.md @@ -0,0 +1,90 @@ +# Transaction-linked native operation outcomes + +HUB-WP-0012 source candidate, 2026-09-28. Live sender admission, database rollout +and owner acceptance remain open. + +## Coverage and transaction boundary + +The durable native registry, message, progress-event and interaction-event writes +already append to `runtime_audit_ledger` in their business transaction. Enforced +requests now bind those ledger rows to the verified issuer/subject, actor/target +tenant, principal type, policy decision/version/caller, request digest and +server-generated authorization correlation ID. This context comes from the +access boundary, never command payloads or caller-supplied correlation headers. +It is scoped to handler execution and reset in `finally`. + +The same transaction inserts one immutable `hub.operation.committed` envelope +into `runtime_outcome_outbox`, keyed by its ledger ID. Failure to insert either +record rolls back the business mutation. A rolled-back transaction has no +committed outcome; the earlier authorization receipt remains an authorization +attempt. Registry duplicate acknowledgements have `registry.duplicate` operation +records and do not pretend a new registration was created. + +The envelope contains business object references and a payload hash, not message +bodies, event payloads, credentials or tokens. Business correlation IDs are +retained separately from verified authorization correlation IDs. That verified +ID and decision ID join the outcome to the pre-execution signed decision already +held by Audit Core. The outcome does not duplicate the signed envelope. + +This slice covers the native durable mutation ledger. It does not claim outcome +coverage for compatibility/embedded-host mutations, background projection refresh, +external services, or in-memory development stores. Non-enforced maintenance +writes retain their existing local ledger behavior without fabricating an actor +or emitting an attributed remote outcome. + +## Delivery, recovery and readiness + +When a durable store and outcome-capable audit sink are composed, the runtime +starts a dispatcher and cancels it before closing clients or stores. It selects +up to 25 eligible rows per batch, one database transaction per row, and waits +one second between batches. PostgreSQL uses `FOR UPDATE SKIP LOCKED` so workers +can select different rows. Each receiver call is bounded to three seconds. + +`AuditCoreSink.append_outcome` reuses the operational-custody probe, rotating +sender credential and exact accepted/duplicate receipt checks. All retries use +the same event ID, occurrence timestamp and body; the ID is also the idempotency +key. Only an accepted receipt marks the row delivered. Failed attempts retain +the row with persisted exponential retry delay (one second up to sixty seconds). +Exceptions are not stored as diagnostic text. Cancellation or a crash before the +local delivery acknowledgement leaves the row eligible for replay. Delivery is +at least once; receiver deduplication resolves a lost receipt or duplicate send. +It does not change business mutation idempotency or make a client retry safe. + +The authorization audit still requires synchronous remote custody before the +handler. Queuing a committed outcome does not weaken that check. An outcome +outage after authorization cannot undo an already committed operation. + +Protected readiness adds `outcome_delivery`: missing delivery composition is +`unavailable`, a pending row older than sixty seconds is `stale`, and a young or +empty backlog is `ok`. Database readiness checks that the new table is accessible. +Pending rows are never silently discarded. Delivered rows remain for inspection; +retention/pruning needs a separately admitted operational policy. + +## Schema and rollout + +Migration `0006_outcome_outbox` adds the table, ledger foreign key and pending-row +index. Apply it through the existing owner migration lane before deploying this +candidate and admit runtime SELECT/INSERT/UPDATE privileges. Automatic downgrade +requires an online check and refuses a nonempty pending queue. Drain pending +outcomes and preserve required custody records before schema rollback. + +The Audit Core sender must admit `hub.operation.committed` for source `hub-core` +and tenant `tenant:platform`, in addition to existing authorization record classes. +This document neither issues that grant nor claims production delivery. + +## Local evidence + +Tests use SQLite-backed durable stores to prove atomic rollback, handler rejection, +all native mutation families, actor/decision attribution, concurrent request +isolation, persisted retry delay, reopen/replay, cancellation, stale-backlog +readiness, dispatcher drain/shutdown and migration shape/downgrade refusal. +The real Audit Core receiver/storage source accepts the eight-field envelope and +returns a duplicate receipt on replay. Its operational-readiness classification +is explicitly a test fixture; this is not live custody evidence. + +PostgreSQL multiworker lock scheduling, production grants, retention and deployed +failure-detection/receiver acceptance still require integration receipts. + +Validation on 2026-09-28: 366 tests pass with the opt-in owner-source suite enabled; +inventory drift, package build and isolated installed-wheel checks pass. The +wheel includes the current migration, authority-context and durable-store code. diff --git a/docs/owner-access-integration.md b/docs/owner-access-integration.md index 75242bc..2929986 100644 --- a/docs/owner-access-integration.md +++ b/docs/owner-access-integration.md @@ -46,8 +46,9 @@ at three seconds, uses TLS, and never follows redirects. Allow is blocked until the archive accepts the authorization record. A lost receipt blocks the business operation even if the attempt reached storage. This is a pre-execution authorization journal, not proof that an operation committed. -There is no local success buffer or silent redaction. Domain transaction/outcome -atomicity and failure detection remain separate T03/T04 acceptance gates. +Authorization cannot use a local success buffer or silent redaction. Native +transaction outcomes now use a separate [durable outbox](operation-outcome-audit.md); +its production delivery and broader mutation coverage remain T03/T04 gates. The exact verified signed decision is retained under `data.signed_decision` as serialized JSON so another serialization of the archive cannot reorder its Go diff --git a/hub_core/migrations/versions/0006_outcome_outbox.py b/hub_core/migrations/versions/0006_outcome_outbox.py new file mode 100644 index 0000000..2b674c4 --- /dev/null +++ b/hub_core/migrations/versions/0006_outcome_outbox.py @@ -0,0 +1,31 @@ +"""Transaction-linked committed outcome delivery. + +Revision ID: 0006_outcome_outbox +Revises: 0005_message_identity_aliases +""" +from alembic import op +import sqlalchemy as sa + +revision = '0006_outcome_outbox' +down_revision = '0005_message_identity_aliases' +branch_labels = depends_on = None + + +def upgrade(): + op.create_table('runtime_outcome_outbox', + 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), + sa.Column('next_attempt', sa.Float(), nullable=False), + sa.Column('delivered_at', sa.Float(), nullable=True)) + op.create_index('ix_runtime_outcome_pending','runtime_outcome_outbox',['delivered_at','next_attempt']) + + +def downgrade(): + pending = op.get_bind().execute(sa.text( + 'SELECT id FROM runtime_outcome_outbox WHERE delivered_at IS NULL LIMIT 1')) + if pending is None or pending.first() is not None: + raise RuntimeError('outcome downgrade requires an online check and a drained outbox') + op.drop_index('ix_runtime_outcome_pending',table_name='runtime_outcome_outbox') + op.drop_table('runtime_outcome_outbox') diff --git a/hub_core/runtime/app.py b/hub_core/runtime/app.py index 2afd63b..412b878 100644 --- a/hub_core/runtime/app.py +++ b/hub_core/runtime/app.py @@ -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, diff --git a/hub_core/runtime/postgres_store.py b/hub_core/runtime/postgres_store.py index cd000f6..35a6ee8 100644 --- a/hub_core/runtime/postgres_store.py +++ b/hub_core/runtime/postgres_store.py @@ -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( diff --git a/hub_core/runtime/tables.py b/hub_core/runtime/tables.py index bbd78da..97aa997 100644 --- a/hub_core/runtime/tables.py +++ b/hub_core/runtime/tables.py @@ -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"), +) diff --git a/hub_core/security/audit.py b/hub_core/security/audit.py index 2b69e13..74c564f 100644 --- a/hub_core/security/audit.py +++ b/hub_core/security/audit.py @@ -16,7 +16,7 @@ class AuditCoreSink: """No allow returns before durable remote custody acknowledges its record. A lost receipt blocks execution, even if the archive already stored the - attempt. This is an authorization-attempt journal, not a mutation outbox. + attempt. Committed outcomes arrive separately from the transactional outbox. """ def __init__(self, *, base_url: str, token_file: Path, client: httpx.AsyncClient): @@ -34,30 +34,32 @@ class AuditCoreSink: raise ValueError("operational audit custody is required") async def append(self, record: dict) -> None: + correlation, outcome = record.get("correlation_id"), record.get("outcome") + if not isinstance(correlation, str) or not correlation or outcome not in {"authorized", "denied", "refused"}: + raise AccessFailure(503, "audit_unavailable") + await self._deliver({ + "id": str(uuid4()), "type": "hub.access." + outcome, "source": "hub-core", + "subject": "hub-access:" + correlation, "tenant": "tenant:platform", + "correlation_id": correlation, "occurred_at": datetime.now(timezone.utc).isoformat(), + "data": record, + }) + + async def append_outcome(self, event: dict) -> None: + # The transaction supplies the immutable ID, timestamp and full envelope. + if (event.get("type") != "hub.operation.committed" or event.get("source") != "hub-core" + or event.get("tenant") != "tenant:platform" + or not isinstance(event.get("id"), str) or not event["id"] + or set(event) != {"id", "type", "source", "subject", "tenant", "correlation_id", "occurred_at", "data"}): + raise AccessFailure(503, "audit_unavailable") + await self._deliver(event) + + async def _deliver(self, event: dict) -> None: try: async with asyncio.timeout(3): - # Probe each time, so a receiver's development fallback cannot - # be mistaken for admitted custody through a cached readiness. await self.readiness() token = self.token_file.read_text().strip() if not token or not token.isascii() or any(c.isspace() for c in token): raise ValueError("invalid sender credential") - correlation = record.get("correlation_id") - outcome = record.get("outcome") - if not isinstance(correlation, str) or not correlation or outcome not in { - "authorized", "denied", "refused", - }: - raise ValueError("invalid authorization audit record") - event = { - "id": str(uuid4()), - "type": "hub.access." + outcome, - "source": "hub-core", - "subject": "hub-access:" + correlation, - "tenant": "tenant:platform", - "correlation_id": correlation, - "occurred_at": datetime.now(timezone.utc).isoformat(), - "data": record, - } raw = json.dumps(event, ensure_ascii=False, allow_nan=False).encode() if len(raw) > 256 * 1024: raise ValueError("audit envelope exceeds receiver limit") diff --git a/hub_core/security/boundary.py b/hub_core/security/boundary.py index a3a61ef..de26e04 100644 --- a/hub_core/security/boundary.py +++ b/hub_core/security/boundary.py @@ -10,7 +10,7 @@ import hashlib import json import math import time -from dataclasses import dataclass +from dataclasses import dataclass, replace from importlib.resources import files from typing import Protocol from uuid import uuid4 @@ -20,6 +20,7 @@ from starlette.responses import JSONResponse from starlette.routing import Match from hub_core.security.identity import AccessFailure, Actor +from hub_core.security.context import current_authorization PROFILE = "hub-core.access/1.0.0" @@ -59,6 +60,9 @@ class Authorization: facts: LiveFacts correlation_id: str request_digest: str + decision_id: str = "" + policy_version: str = "" + policy_caller: str = "" @dataclass(frozen=True) @@ -156,7 +160,8 @@ class AccessController: raise AccessFailure(401, "expired_access_token") if time.time() - facts.checked_at > 5: raise AccessFailure(503, "facts_expired_during_authorization") - return context + return replace(context, decision_id=decision.decision_id, + policy_version=decision.policy_version, policy_caller=decision.caller) def route_key(route, method: str) -> str: @@ -312,4 +317,8 @@ class AccessBoundary: (b"x-correlation-id", correlation.encode())]) await send(message) - await self.app(scope, replay, protected_send) + binding = current_authorization.set(context) + try: + await self.app(scope, replay, protected_send) + finally: + current_authorization.reset(binding) diff --git a/hub_core/security/context.py b/hub_core/security/context.py new file mode 100644 index 0000000..3db44c2 --- /dev/null +++ b/hub_core/security/context.py @@ -0,0 +1,8 @@ +"""Request-scoped, verified authority for transaction attribution (never tokens).""" +from contextvars import ContextVar +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from hub_core.security.boundary import Authorization + +current_authorization: ContextVar['Authorization | None'] = ContextVar('hub_authorization', default=None) diff --git a/tests/test_audit_core_owner_contract.py b/tests/test_audit_core_owner_contract.py index bea6e15..9a26aca 100644 --- a/tests/test_audit_core_owner_contract.py +++ b/tests/test_audit_core_owner_contract.py @@ -210,3 +210,35 @@ def test_root_gate_retains_verifiable_decision_before_native_write(tmp_path): asyncio.run(client.aclose()) native.close() backend.close() + + +def test_owner_receiver_deduplicates_immutable_committed_outcome(tmp_path): + from uuid import uuid4 + from datetime import datetime, timezone + backend = OperationalReceiptFixture(str(tmp_path/'outcome-owner.db')) + event_id = str(uuid4()) + event = {'id':event_id,'type':'hub.operation.committed','source':'hub-core', + 'subject':'hub-outcome:'+event_id,'tenant':'tenant:platform', + 'correlation_id':str(uuid4()),'occurred_at':datetime.now(timezone.utc).isoformat(), + 'data':{'outcome':'committed','operation':'message.accepted','payload_hash':'a'*64, + 'authorization':{'subject':'fixture-root','decision_id':'decision:fixture'}}} + credential = tmp_path/'sender' + credential.write_text('fixture-only') + codes = [] + try: + with httpx.Client(transport=httpx.WSGITransport(receiver(backend))) as owner: + def handle(request): + response = owner.request(request.method,str(request.url),content=request.content,headers=request.headers) + if request.method == 'POST': + codes.append(response.status_code) + return httpx.Response(response.status_code,content=response.content) + async def run(): + async with httpx.AsyncClient(transport=httpx.MockTransport(handle)) as client: + sink = AuditCoreSink(base_url='https://audit.fixture',token_file=credential,client=client) + await sink.append_outcome(event) + await sink.append_outcome(event) + asyncio.run(run()) + assert codes == [202,200] + assert backend.get(event_id)['details']['data'] == event['data'] + finally: + backend.close() diff --git a/tests/test_outcome_audit.py b/tests/test_outcome_audit.py new file mode 100644 index 0000000..3033dad --- /dev/null +++ b/tests/test_outcome_audit.py @@ -0,0 +1,258 @@ +import asyncio +import json +import time +from dataclasses import replace +from uuid import uuid4 + +import httpx +import pytest +import sqlalchemy as sa +from fastapi.testclient import TestClient + +from hub_core.runtime.app import create_app +from hub_core.runtime.config import RuntimeSettings +from hub_core.runtime.postgres_store import PostgresPortStore +from hub_core.runtime.tables import runtime_metadata, runtime_messages, runtime_audit_ledger, runtime_outcome_outbox +from hub_core.security.audit import AuditCoreSink +from hub_core.security.context import current_authorization +from test_access_boundary import Owners, HEADERS +from test_access_audit import READY + + +@pytest.fixture +def target(tmp_path): + url = f'sqlite+aiosqlite:///{tmp_path / "outcomes.db"}' + store = PostgresPortStore.from_url(url) + async def schema(): + async with store.engine.begin() as connection: + await connection.run_sync(runtime_metadata.create_all) + asyncio.run(schema()) + owners = Owners() + async def facts(actor, resource): + return replace(owners.facts, checked_at=time.time(), subject=actor.subject) + owners.resolve = facts + app = create_app(settings=RuntimeSettings(environment='test',access_mode='enforce',backend='postgresql',database_url=url), + port_store=store,access_controller=owners.controller()) + with TestClient(app,raise_server_exceptions=False) as client: + yield client, store, owners, url + asyncio.run(store.aclose()) + + +def message(): + return {'schema_version':'0.1.0','correlation_id':str(uuid4()),'from_address':'agent:root', + 'to_addresses':['agent:reader'],'body':'private body never archived'} + + +def rows(store, table): + async def read(): + async with store.sessions() as session: + return [dict(r) for r in (await session.execute(sa.select(table))).mappings()] + return asyncio.run(read()) + + +def test_committed_outcome_is_atomic_and_joins_authorization(target): + client,store,owners,_ = target + body = message() + response = client.post('/ports/messaging/messages',headers=HEADERS,json=body) + assert response.status_code == 202 + ledger, = rows(store,runtime_audit_ledger) + pending, = rows(store,runtime_outcome_outbox) + envelope = pending['envelope'] + authorization = envelope['data']['authorization'] + assert envelope['id'] == ledger['id'] + assert envelope['correlation_id'] == owners.records[0]['correlation_id'] == response.headers['x-correlation-id'] + assert authorization == ledger['detail']['authorization'] + assert authorization['decision_id'] == owners.records[0]['decision_id'] + assert authorization['subject'] == 'immutable-root' + assert envelope['data']['business_correlation_id'] == body['correlation_id'] + assert envelope['data']['subject_id'] == response.json()['id'] + assert 'private body' not in json.dumps(envelope) + assert 'verified-root' not in json.dumps(envelope) + assert current_authorization.get() is None + + +def test_outbox_insert_failure_rolls_back_business_and_ledger(target): + client,store,owners,_ = target + def reject(connection,cursor,statement,parameters,context,many): + if statement.startswith('INSERT INTO runtime_outcome_outbox'): + raise RuntimeError('simulated storage failure') + sa.event.listen(store.engine.sync_engine,'before_cursor_execute',reject) + try: + assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 500 + finally: + sa.event.remove(store.engine.sync_engine,'before_cursor_execute',reject) + assert rows(store,runtime_messages) == rows(store,runtime_audit_ledger) == rows(store,runtime_outcome_outbox) == [] + assert owners.records[0]['outcome'] == 'authorized' # Never claims commit. + assert current_authorization.get() is None + + +def test_denial_and_invalid_body_never_queue_committed_outcome(target): + client,store,owners,_ = target + assert client.post('/ports/messaging/messages',headers=HEADERS,json={}).status_code == 422 + owners.allow = False + assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 403 + assert not rows(store,runtime_outcome_outbox) + assert not rows(store,runtime_messages) + + +def test_lost_receipt_retry_reopen_and_duplicate_delivery(target,tmp_path): + client,store,owners,url = target + assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202 + credential = tmp_path/'sender' + credential.write_text('fixture-only') + archived, requests = {}, [] + def receiver(request): + if request.url.path == '/readyz': + return httpx.Response(200,json=READY) + envelope = json.loads(request.content) + requests.append(bytes(request.content)) + assert request.headers['idempotency-key'] == envelope['id'] + if envelope['id'] not in archived: + archived[envelope['id']] = envelope + raise httpx.ReadTimeout('lost receipt after custody') + assert archived[envelope['id']] == envelope + return httpx.Response(200,json={'status':'duplicate','reference':'audit:'+envelope['id']}) + async def run(): + async with httpx.AsyncClient(transport=httpx.MockTransport(receiver)) as http: + sink = AuditCoreSink(base_url='https://audit.example',token_file=credential,client=http) + assert await store.deliver_outcomes(sink) == 0 + assert await store.deliver_outcomes(sink) == 0 # Backoff, no immediate retry. + assert len(requests) == 1 + await store.aclose() + reopened = PostgresPortStore.from_url(url) + try: + async with reopened.sessions.begin() as session: + await session.execute(runtime_outcome_outbox.update().values(next_attempt=0)) + assert await reopened.deliver_outcomes(sink) == 1 + assert await reopened.deliver_outcomes(sink) == 0 + finally: + await reopened.aclose() + asyncio.run(run()) + assert len(archived) == 1 and len(requests) == 2 and requests[0] == requests[1] + pending, = rows(store,runtime_outcome_outbox) + assert pending['attempts'] == 2 and pending['delivered_at'] is not None + + +def test_cancellation_after_receipt_leaves_replayable_row(target): + client,store,_,_ = target + assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202 + class Crash: + async def append_outcome(self,event): + raise asyncio.CancelledError() + async def run(): + with pytest.raises(asyncio.CancelledError): + await store.deliver_outcomes(Crash()) + asyncio.run(run()) + pending, = rows(store,runtime_outcome_outbox) + assert pending['delivered_at'] is None and pending['attempts'] == 0 + + +def test_stale_backlog_is_visible_and_clears_after_receipt(target): + client,store,_,_ = target + assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202 + class Receipt: + async def append_outcome(self,event): + pass + async def run(): + assert await store.outcome_readiness() == 'ok' + async with store.sessions.begin() as session: + await session.execute(runtime_outcome_outbox.update().values(created_at=time.time()-61)) + assert await store.outcome_readiness() == 'stale' + await store.deliver_outcomes(Receipt()) + assert await store.outcome_readiness() == 'ok' + asyncio.run(run()) + + +def test_all_native_mutation_families_keep_attributed_outcomes(target): + from hub_core.conformance import ConformanceHarness + client,store,owners,_ = target + owners.facts = replace(owners.facts,producer_addresses=frozenset({'hub:ops-hub'})) + # Composition is frozen when the app is created; run the twelve business + # checks except dependency readiness, which separately reports no dispatcher. + harness = ConformanceHarness(client) + client.headers.update(HEADERS) + report = harness.run() + assert all(c.status == 'pass' for c in report.checks if c.check_id != 'C9'), report.to_dict() + assert next(c for c in report.checks if c.check_id == 'C9').status == 'fail' + assert client.get('/readyz').json()['checks']['outcome_delivery'] == 'unavailable' + pending = rows(store,runtime_outcome_outbox) + operations = {r['envelope']['data']['operation'] for r in pending} + assert {'message.accepted','event.progress.accepted','event.interaction.accepted'} <= operations + assert any(op.startswith('registry.') for op in operations) + assert all(r['envelope']['data']['authorization']['subject'] == 'immutable-root' for r in pending) + assert len(pending) == len(rows(store,runtime_audit_ledger)) + + +def test_concurrent_requests_keep_separate_authorization_contexts(target): + from hub_core.security.identity import AccessFailure + client,store,owners,_ = target + async def authenticate(token): + if token not in {'workload:a','workload:b'}: + raise AccessFailure(401,'bad_fixture') + await asyncio.sleep(0) + return replace(owners.actor,subject=token,principal_type='service') + owners.authenticate = authenticate + async def run(): + async with httpx.AsyncClient(transport=httpx.ASGITransport(app=client.app),base_url='https://hub.example') as http: + async def send(subject): + result = await http.post('/ports/messaging/messages',headers={'Authorization':'Bearer '+subject},json=message()) + assert result.status_code == 202 + return result.json()['id'], subject + return dict(await asyncio.gather(send('workload:a'),send('workload:b'))) + subjects = asyncio.run(run()) + for row in rows(store,runtime_outcome_outbox): + data = row['envelope']['data'] + assert data['authorization']['subject'] == subjects[data['subject_id']] + assert current_authorization.get() is None + + +def test_outcome_migration_creates_index_and_foreign_key(tmp_path): + import importlib + from alembic.migration import MigrationContext + from alembic.operations import Operations + migration = importlib.import_module('hub_core.migrations.versions.0006_outcome_outbox') + engine = sa.create_engine('sqlite:///' + str(tmp_path/'migration.db')) + try: + runtime_audit_ledger.create(engine) + with engine.begin() as connection: + with Operations.context(MigrationContext.configure(connection)): + migration.upgrade() + inspect = sa.inspect(connection) + assert {x['name'] for x in inspect.get_columns('runtime_outcome_outbox')} == set(runtime_outcome_outbox.c.keys()) + assert inspect.get_foreign_keys('runtime_outcome_outbox')[0]['referred_table'] == 'runtime_audit_ledger' + assert inspect.get_indexes('runtime_outcome_outbox')[0]['name'] == 'ix_runtime_outcome_pending' + connection.execute(runtime_outcome_outbox.insert().values(id='fixture',envelope={},created_at=0,attempts=0,next_attempt=0)) + with Operations.context(MigrationContext.configure(connection)): + with pytest.raises(RuntimeError,match='drained outbox'): + migration.downgrade() + connection.execute(runtime_outcome_outbox.delete()) + with Operations.context(MigrationContext.configure(connection)): + migration.downgrade() + assert 'runtime_outcome_outbox' not in sa.inspect(connection).get_table_names() + finally: + engine.dispose() + + +def test_runtime_dispatcher_drains_pending_outcome_and_stops(target): + from hub_core.runtime.app import _deliver_outcomes + client,store,_,_ = target + assert client.post('/ports/messaging/messages',headers=HEADERS,json=message()).status_code == 202 + async def run(): + class Receipt: + async def append_outcome(self,envelope): + pass + task = asyncio.create_task(_deliver_outcomes(store,Receipt())) + try: + async with asyncio.timeout(3): + while True: + async with store.sessions() as session: + count = (await session.execute(sa.select(sa.func.count()).select_from(runtime_outcome_outbox) + .where(runtime_outcome_outbox.c.delivered_at.is_not(None)))).scalar_one() + if count == 1: + break + await asyncio.sleep(.01) + finally: + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + asyncio.run(run()) diff --git a/workplans/HUB-WP-0012-netkingdom-platform-root-access.md b/workplans/HUB-WP-0012-netkingdom-platform-root-access.md index deb351e..420634f 100644 --- a/workplans/HUB-WP-0012-netkingdom-platform-root-access.md +++ b/workplans/HUB-WP-0012-netkingdom-platform-root-access.md @@ -370,6 +370,25 @@ forces a fresh environment and refreshes the Hub package, and the rerun passed. No credential was issued or retrieved and no consumer was switched. T04 remains `progress`; host authentication/exchange, real caller admission and T05 remain open. +## Transaction outcome continuation — 2026-09-28 + +Native durable registry/message/progress/interaction mutations now atomically +store verified authorization attribution, their existing ledger row and an +immutable committed-outcome envelope. Added migration `0006_outcome_outbox`, +retry/backoff delivery with stable receiver idempotency keys, lifecycle-owned +dispatch and protected backlog readiness. Rollback creates no outcome; a lost +receipt can replay after restart. Downgrade refuses undelivered rows. + +[Outcome contract and evidence scope](../docs/operation-outcome-audit.md). +Local SQLite tests and real Audit Core receiver-source checks cover this slice; +production PostgreSQL lock scheduling, grants, receiver admission and retention +remain open. Compatibility/embedded/external writes and background projection +refreshes are not claimed as covered. T03/T04 remain `progress`. +Validation: **366 tests passed**, including the opt-in Audit Core owner-source +suite. Inventory drift, distribution builds, isolated installed-wheel checks +and current migration/context/store wheel contents pass. +No database migration or deployment was applied to a live service. + ## Acceptance checkpoints - [x] Architecture/source/runtime review captured; new implementation owner is hub-core