From 3ae6c330b1f958de7c3d6ea7a5a8da1db015da6c Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Fri, 18 Sep 2026 05:04:25 +0000 Subject: [PATCH 1/3] fix(reporting): cache schema proofs and retain safe receipt diagnostics Refs #1167 --- .github/workflows/ci.yml | 14 +- .github/workflows/pr-title-check.yml | 2 +- docs/reporting-receipt-ingress.md | 43 ++ docs/reporting-webhook-activity.md | 41 ++ src/adcp/reporting/outbox/_schema.py | 7 + src/adcp/reporting/outbox/status_schema.py | 7 + src/adcp/reporting/outbox/support.py | 148 ++++- src/adcp/reporting/receipts/_diagnostics.py | 71 +++ src/adcp/reporting/receipts/handler.py | 6 +- src/adcp/reporting/receipts/pg.py | 41 +- .../reporting/_hardening_installed.py | 157 ++++++ .../reporting/_hardening_packaging.py | 115 ++++ .../test_reporting_activity_migration.py | 2 + .../test_reporting_activity_schema_proof.py | 522 ++++++++++++++++++ ...test_reporting_feed_hardening_installed.py | 191 +++++++ .../test_reporting_feed_packaging.py | 3 + .../test_reporting_receipt_diagnostics.py | 452 +++++++++++++++ 17 files changed, 1787 insertions(+), 35 deletions(-) create mode 100644 src/adcp/reporting/receipts/_diagnostics.py create mode 100644 tests/conformance/reporting/_hardening_installed.py create mode 100644 tests/conformance/reporting/_hardening_packaging.py create mode 100644 tests/conformance/reporting/test_reporting_activity_schema_proof.py create mode 100644 tests/conformance/reporting/test_reporting_feed_hardening_installed.py create mode 100644 tests/conformance/reporting/test_reporting_receipt_diagnostics.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 32267a6e6..7a501fbac 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -11,6 +11,7 @@ on: - conductor/1167b1-materializer-contracts - conductor/1167b2-durable-managed-reporting - conductor/reporting-receipt-ingress-b22 + - conductor/reporting-frozen-account-feed-b23 # Default @adcp/sdk runner alias for storyboard jobs. Tracks the current # stable @adcp/sdk release via the ``latest`` npm dist-tag. @@ -370,6 +371,7 @@ jobs: tests/conformance/reporting/test_reporting_receipt_graph.py \ tests/conformance/reporting/test_reporting_receipt_transactions.py \ tests/conformance/reporting/test_reporting_receipt_transports.py \ + tests/conformance/reporting/test_reporting_receipt_diagnostics.py \ tests/conformance/reporting/test_reporting_receipt_materializer.py \ tests/conformance/reporting/test_reporting_receipt_migration.py \ tests/conformance/reporting/test_reporting_receipt_process.py \ @@ -555,7 +557,7 @@ jobs: pg-reporting-feed-installed: name: Installed frozen feed (Python 3.10 VCS and sdist) runs-on: ubuntu-latest - timeout-minutes: 25 + timeout-minutes: 50 services: postgres: image: postgres:16 @@ -572,6 +574,8 @@ jobs: --health-retries 10 steps: - uses: actions/checkout@v6 + - name: Fetch the exact approved B2.3 hardening comparison artifact + run: git fetch --no-tags origin 50e35f0ae3540f19b40e8fc460f5870dfe018bf9 - uses: actions/setup-python@v6 id: feed-python310 with: @@ -585,21 +589,25 @@ jobs: run: pip install -e ".[dev,pg]" - name: Run installed base and PostgreSQL restart cells shell: bash - timeout-minutes: 20 + timeout-minutes: 40 env: ADCP_PG_TEST_URL: postgresql://postgres@localhost:5432/adcp_feed_installed_test ADCP_PYTHON310: ${{ steps.feed-python310.outputs.python-path }} + ADCP_HARDENING_EVIDENCE: ${{ runner.temp }}/hardening-installed-evidence run: | python scripts/reporting_test_harness.py pytest \ tests/conformance/reporting/test_reporting_feed_packaging.py \ tests/conformance/reporting/test_reporting_feed_installed_pg.py \ + tests/conformance/reporting/test_reporting_feed_hardening_installed.py \ -v -s -ra | tee pg-reporting-feed-installed-evidence.log - name: Preserve installed origins, SQL, strict adopter and cold page evidence if: always() uses: actions/upload-artifact@v7 with: name: pg-reporting-feed-installed-evidence-${{ github.run_attempt }} - path: pg-reporting-feed-installed-evidence.log + path: | + pg-reporting-feed-installed-evidence.log + ${{ runner.temp }}/hardening-installed-evidence if-no-files-found: error conventional-commits: diff --git a/.github/workflows/pr-title-check.yml b/.github/workflows/pr-title-check.yml index c989316c2..de0669486 100644 --- a/.github/workflows/pr-title-check.yml +++ b/.github/workflows/pr-title-check.yml @@ -3,7 +3,7 @@ name: PR Title Check on: pull_request: types: [opened, edited, synchronize, reopened] - branches: [main, conductor/1167b2-durable-managed-reporting, conductor/reporting-receipt-ingress-b22] + branches: [main, conductor/1167b2-durable-managed-reporting, conductor/reporting-receipt-ingress-b22, conductor/reporting-frozen-account-feed-b23] permissions: contents: read diff --git a/docs/reporting-receipt-ingress.md b/docs/reporting-receipt-ingress.md index 27c807c53..b6d844d3b 100644 --- a/docs/reporting-receipt-ingress.md +++ b/docs/reporting-receipt-ingress.md @@ -46,6 +46,49 @@ Absent, anonymous, inactive or conflicting identities fail closed. Two consumers in one account and one consumer in two accounts have independent batch keys and receipt visibility. +## Private operator diagnostics + +Unexpected storage/driver failures and unexpected account-resolver or custom-store +failures emit one structured ERROR on `adcp.reporting.receipts`. The private helper +records the original failure at the PostgreSQL `_storage_errors` boundary before +translation, or at the receipt handler boundary for an unexpected adopter failure. +The handler does not log an already translated `ReportingReceiptError` again. + +The static message is `Receipt storage is unavailable`. Its only diagnostic +fields are: + +- `code`: always `RECEIPT_STORAGE_UNAVAILABLE`. +- `boundary`: `handler`, `store.create_schema`, `store.receipt_ingestion_ready`, + `store.ingest_receipt_batch`, or `store.read_receipt_boundaries`. +- `exception_type`: the exception class name. +- `origin_module`, `origin_function`, `origin_line`: the deepest original source + coordinates, with only bounded ASCII identifiers/dotted module names accepted; + invalid names become `unknown`, and an absent traceback has line `0`. + +No exception, traceback object, message, arguments, chain, frame locals or raw +traceback path is attached to the record. No request/batch/provider body, account +or consumer identity, receipt/idempotency/continuation identifier, SQL, bound +parameter, DSN, authentication value or financial data is a diagnostic field. +The helper creates a plain record without the ambient record factory or dynamic +task/thread/process names, so serializing the entire emitted `LogRecord` retains +this boundary. Operator handlers/filters must preserve that contract rather than +adding request context or raw exceptions. +If an operator logging sink raises, the buyer still receives the same safe +error; the SDK does not retry through another logger or expose the sink failure. + +Expected `INVALID_REQUEST`, `UNAUTHORIZED`, `IDEMPOTENCY_CONFLICT`, +`RECEIPT_SCHEMA_UNREADY` and `RECEIPT_HISTORY_CORRUPT` remain silent operator paths. +Cancellation propagates unchanged and emits no operator error. The buyer receives +the same safe code and message as before, with an actually empty exception +cause/context; transport formatting and deliberate caller-owned `context` echo +remain unchanged. That echo is not part of the error diagnostic. + +Use the static boundary, class and source coordinates to locate the failing SDK +or adopter seam. Repair connectivity or the installed schema, re-run readiness, +and restart using the rollout procedure below; retry the original immutable batch +and key after recovery. Do not enable raw exception/SQL logging to diagnose this +path or rewrite receipt history to make an error disappear. + ## Wire and replay contract Requests negotiate `adcp_version: "3.2-rc.3"`. This is the wire release spelling; diff --git a/docs/reporting-webhook-activity.md b/docs/reporting-webhook-activity.md index 04f8d084b..d2a0aa52b 100644 --- a/docs/reporting-webhook-activity.md +++ b/docs/reporting-webhook-activity.md @@ -53,6 +53,47 @@ Relationship notification support is independent and supplies no evidence for either flag. The memory implementation supports shared conformance vectors; it cannot justify a durable claim. +## Schema proof lifetime and readiness recovery + +`ReportingActivitySupport` remains a frozen public composition value. Its private +cache holds only a completed **positive schema proof**, scoped to that one support +instance, the exact B or B+C wiring identities, and the packaged required-object +manifest contracts. Separate support instances never share proof, even on the +same pool. Each pool must retain its deployment's database/search-path configuration; +do not change session search paths behind a serving support instance. + +Cold concurrent discovery calls share one catalog scan and one pool checkout. +Composite B+C activity validates both required contracts once against that capture. +After success, discovery needs no catalog checkout, including when the pool is +busy. Only primitive proof state survives completion: synchronous startup +validation using `asyncio.run` can be followed by discovery on another event loop. +An overlapping cold validator on a different loop fails closed; finish startup +before serving. Canceling a discovery waiter does not cancel other waiters' scan. +False results, failed/canceled scans and exceptions never become positive proof. +The next call retries after repair. + +Every request still recomputes its capability response and checks its claims, +exact component types, object identity, pool wiring, notification enablement, +B+C scheduling/union wiring, and account-listing/projector topology. A false +request claim does not warm the cache. Schema evidence grants no additional +materializer, status or higher-tier readiness and does not cache authorization. + +For deployment changes, stop admission, drain workers and requests, migrate with +the deployment connection and `await ledger.create_schema()`, construct a fresh +support/server, validate readiness, then resume serving. Manifest validation is +read-only; it does not install schema or private fairness-bootstrap objects. +On readiness failure, keep admission closed, repair the required migration/object, +and retry validation. Preserve existing immutable history and pending-effect +recovery rules throughout rollback or restart. + +For controlled BYO DDL/tests on an existing composition, drain callers, call +`support.invalidate_schema_validation()` **before** the DDL, complete the change, +then validate before admitting new work. Invalidation advances an epoch: a +previous in-flight scan cannot publish into the new epoch. A scan already in +progress may finish, but its caller cannot use that invalidated proof. Arbitrary +out-of-band DDL while serving is unsupported without this procedure or a fresh +support after drain/migrate/restart; discovery does not automatically detect it. + ## Canonical consumer and account visibility Call `resolve_reporting_consumer(auth_info=..., agent=...)` when registering a diff --git a/src/adcp/reporting/outbox/_schema.py b/src/adcp/reporting/outbox/_schema.py index a130b89db..c63fd2f3c 100644 --- a/src/adcp/reporting/outbox/_schema.py +++ b/src/adcp/reporting/outbox/_schema.py @@ -117,6 +117,13 @@ async def validate_schema(connection: Any, *, activity: bool = False) -> None: raise ReportingNotificationError( "notification_schema_unready:catalog_unavailable" ) from None + _validate_schema_objects(installed, activity=activity) + + +def _validate_schema_objects( + installed: dict[str, dict[str, Any]], *, activity: bool = False +) -> None: + """Apply the packaged B contract to an already captured catalog.""" if not REQUIRED_OBJECTS: raise ReportingNotificationError("notification_schema_unready:manifest_missing") for key, expected in REQUIRED_OBJECTS.items(): diff --git a/src/adcp/reporting/outbox/status_schema.py b/src/adcp/reporting/outbox/status_schema.py index c83a17dc6..a8c8ff010 100644 --- a/src/adcp/reporting/outbox/status_schema.py +++ b/src/adcp/reporting/outbox/status_schema.py @@ -25,6 +25,13 @@ async def validate_status_schema( # separately; this manifest's activity switch validates only C objects. await validate_schema(connection) installed = await schema_objects(connection) + _validate_status_objects(installed, activity=activity, status=status) + + +def _validate_status_objects( + installed: dict[str, dict[str, Any]], *, activity: bool = False, status: bool = True +) -> None: + """Apply only C's contract; the caller must also prove its B foundation.""" if not REQUIRED_STATUS_OBJECTS: raise ReportingNotificationError("status_schema_unready:manifest_missing") for key, expected in REQUIRED_STATUS_OBJECTS.items(): diff --git a/src/adcp/reporting/outbox/support.py b/src/adcp/reporting/outbox/support.py index a76f4fdd2..3c62043e5 100644 --- a/src/adcp/reporting/outbox/support.py +++ b/src/adcp/reporting/outbox/support.py @@ -2,7 +2,8 @@ from __future__ import annotations -from dataclasses import dataclass +import asyncio +from dataclasses import dataclass, field from typing import TYPE_CHECKING, Any from adcp.reporting.ledger.notification_models import ReportingNotificationError @@ -14,6 +15,30 @@ from adcp.reporting.outbox.worker import ReportingNotificationWorker +_ProofKey = tuple[str, tuple[int, ...], str, str] + + +@dataclass +class _SchemaFlight: + key: _ProofKey + epoch: int + task: asyncio.Task[bool] | None = None + + +@dataclass +class _SchemaValidation: + epoch: int = 0 + positive: _ProofKey | None = None + flight: _SchemaFlight | None = None + + +def _consume_schema_failure(task: asyncio.Task[bool]) -> None: + # A canceled waiter does not cancel other callers' proof. If every waiter + # leaves, retrieve the result without logging catalog/driver exceptions. + if not task.cancelled(): + task.exception() + + @dataclass(frozen=True) class ReportingActivitySupport: """The concrete writer/store/projector chain scheduled by the adopter. @@ -28,6 +53,101 @@ class ReportingActivitySupport: ledger: ReportingLedgerStore projector: ReportingActivityProjector | None = None status: ReportingStatusSupport | None = None + _schema_validation: _SchemaValidation = field( + default_factory=_SchemaValidation, init=False, repr=False, compare=False + ) + + def invalidate_schema_validation(self) -> None: + """Discard schema evidence before controlled DDL, after draining callers. + + An older in-flight scan cannot publish into the new epoch. A fresh + support/server instance after migration is the normal deployment path; + serving-time DDL is not automatically detected by discovery. + """ + state = self._schema_validation + state.epoch += 1 + state.positive = None + state.flight = None + + async def _schema_proven(self, pool: Any, *, composite: bool) -> bool: + from adcp.reporting.outbox import _schema + + assert self.projector is not None + wiring: tuple[object, ...] = ( + self.worker, + self.ledger, + self.projector, + self.projector.store, + self.worker.outbox, + self.worker.activity, + self.worker.cipher, + self.worker.subscriptions, + self.worker.signing, + pool, + ) + status_contract = "" + if composite: + from adcp.reporting.outbox import status_schema + + assert self.status is not None + wiring += ( + self.status, + self.status.store, + self.status.worker, + self.status.worker.outbox, + ) + status_contract = _schema._digest(status_schema.REQUIRED_STATUS_OBJECTS) + key: _ProofKey = ( + "B+C" if composite else "B", + tuple(id(component) for component in wiring), + _schema._digest(_schema.REQUIRED_OBJECTS), + status_contract, + ) + state = self._schema_validation + if state.positive == key: + return True + flight = state.flight + if flight is None: + flight = _SchemaFlight(key, state.epoch) + state.flight = flight + flight.task = asyncio.create_task(self._scan_schema(pool, flight, composite=composite)) + flight.task.add_done_callback(_consume_schema_failure) + task = flight.task + if task is None or flight.key != key or task.get_loop() is not asyncio.get_running_loop(): + # Overlapping, differently wired/looped startup is not evidence. + # Completed startup leaves no task or loop in this instance. + return False + if not await asyncio.shield(task): + return False + # Wiring may have changed while the catalog checkout was suspended. + # Rerun the same cheap checks before accepting the completed proof. + return await self.durable() + + async def _scan_schema(self, pool: Any, flight: _SchemaFlight, *, composite: bool) -> bool: + from adcp.reporting.outbox import _schema + + state = self._schema_validation + try: + async with pool.connection() as connection: + try: + installed = await _schema.schema_objects(connection) + except Exception: + raise ReportingNotificationError( + "notification_schema_unready:catalog_unavailable" + ) from None + # One physical catalog capture, each required contract once. + _schema._validate_schema_objects(installed, activity=True) + if composite: + from adcp.reporting.outbox import status_schema + + status_schema._validate_status_objects(installed, activity=True, status=False) + if state.epoch != flight.epoch or state.flight is not flight: + return False + state.positive = flight.key + return True + finally: + if state.flight is flight: + state.flight = None async def durable(self) -> bool: from adcp.reporting.ledger.delivery import InMemoryReportingReconciliationStore @@ -59,7 +179,6 @@ async def durable(self) -> bool: # Lazy imports retain base-install operation without the [pg] extra. from adcp.reporting.ledger.delivery_pg import PgReportingReconciliationStore from adcp.reporting.ledger.pg import PgReportingLedgerStore - from adcp.reporting.outbox._schema import validate_schema from adcp.reporting.outbox.pg import PgReportingOutbox if ( @@ -71,21 +190,22 @@ async def durable(self) -> bool: or not self.ledger._notifications_enabled ): raise ReportingNotificationError("activity_chain_unready") - async with self.ledger._pool.connection() as connection: - await validate_schema(connection, activity=True) - return True + return await self._schema_proven(self.ledger._pool, composite=False) async def _composite_durable(self) -> bool: """Only the closed B+C read/purge union may replace the original B reader.""" assert self.status is not None if not self.status.scheduled: return False - from adcp.reporting.outbox._schema import validate_schema + from adcp.reporting.ledger.delivery_pg import PgReportingReconciliationStore + from adcp.reporting.ledger.pg import PgReportingLedgerStore from adcp.reporting.outbox.pg import PgReportingOutbox from adcp.reporting.outbox.routing import ReportingEnvelopeCipher from adcp.reporting.outbox.status_activity_pg import PgReportingActivityUnionStore - from adcp.reporting.outbox.status_pg import PgStatusNotificationStore - from adcp.reporting.outbox.status_schema import validate_status_schema + from adcp.reporting.outbox.status_pg import ( + PgReportingStatusOutbox, + PgStatusNotificationStore, + ) from adcp.reporting.outbox.worker import ReportingNotificationWorker store = self.status.store @@ -94,6 +214,9 @@ async def _composite_durable(self) -> bool: reader = self.projector.store if self.projector is not None else None if ( type(self.worker) is not ReportingNotificationWorker + or type(store) is not PgStatusNotificationStore + or type(self.ledger) not in {PgReportingLedgerStore, PgReportingReconciliationStore} + or not isinstance(self.ledger, PgReportingLedgerStore) or type(self.worker.outbox) is not PgReportingOutbox or not isinstance(self.worker.outbox, PgReportingOutbox) or type(self.worker.cipher) is not ReportingEnvelopeCipher @@ -108,15 +231,16 @@ async def _composite_durable(self) -> bool: or self.status.worker.outbox is not store.outbox or self.ledger is not store.ledger or self.worker.outbox._pool is not store.ledger._pool + or type(store.outbox) is not PgReportingStatusOutbox + or store.outbox._pool is not store.ledger._pool + or self.worker.outbox._clock is not store.outbox._clock + or not store.ledger._notifications_enabled or self.worker.cipher is not self.status.worker.cipher or self.worker.subscriptions is not self.status.worker.subscriptions or self.worker.signing is not self.status.worker.signing ): raise ReportingNotificationError("activity_chain_unready") - async with store.ledger._pool.connection() as connection: - await validate_schema(connection, activity=True) - await validate_status_schema(connection, activity=True, status=False) - return True + return await self._schema_proven(store.ledger._pool, composite=True) async def capability_flags( self, *, account_activity: ReportingActivityProjector | None = None diff --git a/src/adcp/reporting/receipts/_diagnostics.py b/src/adcp/reporting/receipts/_diagnostics.py new file mode 100644 index 000000000..aa34f5805 --- /dev/null +++ b/src/adcp/reporting/receipts/_diagnostics.py @@ -0,0 +1,71 @@ +"""Closed, payload-free diagnostics for unexpected receipt storage failures.""" + +from __future__ import annotations + +import logging +import re +from typing import Literal + +_Boundary = Literal[ + "handler", + "store.create_schema", + "store.receipt_ingestion_ready", + "store.ingest_receipt_batch", + "store.read_receipt_boundaries", +] +_BOUNDARIES = frozenset( + { + "handler", + "store.create_schema", + "store.receipt_ingestion_ready", + "store.ingest_receipt_batch", + "store.read_receipt_boundaries", + } +) +_LOGGER = logging.getLogger("adcp.reporting.receipts") +_IDENTIFIER = re.compile(r"[A-Za-z_][A-Za-z_0-9]*(?:\.[A-Za-z_][A-Za-z_0-9]*)*", re.ASCII) + + +def _coordinate(value: object) -> str: + # Invalid/dynamic paths are discarded, not partially retained by escaping. + return ( + value + if type(value) is str and len(value) <= 160 and _IDENTIFIER.fullmatch(value) + else "unknown" + ) + + +def _storage_failure(error: Exception, *, boundary: _Boundary) -> None: + """Emit only class and source coordinates, never the exception or its text.""" + if not _LOGGER.isEnabledFor(logging.ERROR): + return + origin = error.__traceback__ + while origin is not None and origin.tb_next is not None: + origin = origin.tb_next + fields = { + "code": "RECEIPT_STORAGE_UNAVAILABLE", + "boundary": boundary if boundary in _BOUNDARIES else "handler", + "exception_type": _coordinate(type(error).__name__), + "origin_module": ( + _coordinate(origin.tb_frame.f_globals.get("__name__")) if origin else "unknown" + ), + "origin_function": _coordinate(origin.tb_frame.f_code.co_name) if origin else "unknown", + "origin_line": origin.tb_lineno if origin else 0, + } + # Construct a plain record so an ambient LogRecordFactory cannot attach + # request context. Paths, task names and thread/process names can themselves + # contain adopter data; none are part of this diagnostic's contract. + record = logging.LogRecord( + _LOGGER.name, logging.ERROR, "", 0, "Receipt storage is unavailable", (), None + ) + record.threadName = None + record.processName = None + if hasattr(record, "taskName"): + record.taskName = None + record.__dict__.update(fields) + try: + _LOGGER.handle(record) + except Exception: + # A failing operator sink cannot replace the existing safe buyer error. + # Never try a second logger, which could duplicate or expose the failure. + return diff --git a/src/adcp/reporting/receipts/handler.py b/src/adcp/reporting/receipts/handler.py index c564b4e88..106cb9fbb 100644 --- a/src/adcp/reporting/receipts/handler.py +++ b/src/adcp/reporting/receipts/handler.py @@ -11,6 +11,7 @@ from adcp.reporting.ledger.delivery_models import ReportingDeliveryPrincipal from adcp.reporting.ledger.notification_models import ReportingNotificationError from adcp.reporting.outbox.identity import canonical_consumer, resolve_reporting_consumer +from adcp.reporting.receipts._diagnostics import _storage_failure from adcp.reporting.receipts.errors import ReportingReceiptError from adcp.reporting.receipts.store import ReportingReceiptBatchStore from adcp.reporting.receipts.wire import TASK, validate_receipt_request @@ -244,9 +245,10 @@ async def sync_reporting_receipts( return await self.receipt_store.ingest_receipt_batch(request, caller=caller) except ReportingReceiptError as error: code, message = error.code, str(error) - except Exception: + except Exception as error: + _storage_failure(error, boundary="handler") unavailable = ReportingReceiptError("RECEIPT_STORAGE_UNAVAILABLE") code, message = unavailable.code, str(unavailable) # Leave the exception scope before translating. Credential/ACL adapters - # may raise provider errors; neither logs nor exception chains retain them. + # may raise provider errors; only safe origin coordinates were logged. raise ADCPTaskError(operation=TASK, errors=[Error(code=code, message=message)]) diff --git a/src/adcp/reporting/receipts/pg.py b/src/adcp/reporting/receipts/pg.py index 673deec8a..b219b3ac6 100644 --- a/src/adcp/reporting/receipts/pg.py +++ b/src/adcp/reporting/receipts/pg.py @@ -29,6 +29,7 @@ from adcp.reporting.ledger.store import LedgerConflictError from adcp.reporting.materializer.capture import private_snapshot from adcp.reporting.materializer.pg import PgReportingMaterializerStore, _now +from adcp.reporting.receipts._diagnostics import _Boundary, _storage_failure from adcp.reporting.receipts.capture import ReportingReceiptBoundary, decode_receipt_boundary from adcp.reporting.receipts.errors import ReportingReceiptError from adcp.reporting.receipts.records import ( @@ -51,20 +52,26 @@ def _storage_errors( - fn: Callable[_P, Coroutine[Any, Any, _R]], -) -> Callable[_P, Coroutine[Any, Any, _R]]: - @wraps(fn) - async def wrapped(*args: _P.args, **kwargs: _P.kwargs) -> _R: - try: - return await fn(*args, **kwargs) - except ReportingReceiptError: - raise - except Exception: - unavailable = ReportingReceiptError("RECEIPT_STORAGE_UNAVAILABLE") - # Outside the driver exception scope: no SQL/provider detail in __context__. - raise unavailable + boundary: _Boundary, +) -> Callable[[Callable[_P, Coroutine[Any, Any, _R]]], Callable[_P, Coroutine[Any, Any, _R]]]: + def decorate( + fn: Callable[_P, Coroutine[Any, Any, _R]], + ) -> Callable[_P, Coroutine[Any, Any, _R]]: + @wraps(fn) + async def wrapped(*args: _P.args, **kwargs: _P.kwargs) -> _R: + try: + return await fn(*args, **kwargs) + except ReportingReceiptError: + raise + except Exception as error: + _storage_failure(error, boundary=boundary) + unavailable = ReportingReceiptError("RECEIPT_STORAGE_UNAVAILABLE") + # Outside the driver exception scope: no SQL/provider detail in __context__. + raise unavailable - return wrapped + return wrapped + + return decorate class PgReportingReceiptStore(PgReportingMaterializerStore): @@ -75,7 +82,7 @@ class PgReportingReceiptStore(PgReportingMaterializerStore): public conformance clock was supplied to the inherited constructor. """ - @_storage_errors + @_storage_errors("store.create_schema") async def create_schema(self) -> None: async with self._connection() as connection, connection.transaction(): await self._create_schema_on(connection) @@ -83,7 +90,7 @@ async def create_schema(self) -> None: await connection.execute(root.joinpath("reporting_materializer.sql").read_text()) await connection.execute(root.joinpath("reporting_receipt_ingestion.sql").read_text()) - @_storage_errors + @_storage_errors("store.receipt_ingestion_ready") async def receipt_ingestion_ready(self) -> bool: async with self._connection() as connection: await validate_receipt_schema(connection, notifications=self._notifications_enabled) @@ -265,7 +272,7 @@ def _assemble_receipt_response( raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT") return cast(dict[str, Any], json.loads(_json(response))) - @_storage_errors + @_storage_errors("store.ingest_receipt_batch") async def ingest_receipt_batch( self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal ) -> dict[str, Any]: @@ -363,7 +370,7 @@ async def _capture_receipt_on(self, connection: Any, receipt: ReportingReceiptRe ), ) - @_storage_errors + @_storage_errors("store.read_receipt_boundaries") async def read_receipt_boundaries( self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100 ) -> tuple[ReportingReceiptBoundary, ...]: diff --git a/tests/conformance/reporting/_hardening_installed.py b/tests/conformance/reporting/_hardening_installed.py new file mode 100644 index 000000000..c88169cef --- /dev/null +++ b/tests/conformance/reporting/_hardening_installed.py @@ -0,0 +1,157 @@ +"""Run copied conformance fixtures against non-editable installed SDK bytes only.""" + +import contextlib +import hashlib +import importlib +import json +import os +import sys +import time +from importlib.resources import files +from pathlib import Path + +import pytest + + +class Results: + def __init__(self): + self.passed = self.failed = self.skipped = self.errors = self.deselected = 0 + self.failures = [] + + def pytest_runtest_logreport(self, report): + if report.skipped: + self.skipped += 1 + elif report.failed: + if report.when == "call": + self.failed += 1 + self.failures.append((report.nodeid, report.longreprtext)) + else: + self.errors += 1 + elif report.when == "call": + self.passed += 1 + + def pytest_collectreport(self, report): + if report.failed: + self.errors += 1 + + def pytest_deselected(self, items): + self.deselected += len(items) + + +def main(settings): + root = Path(settings["fixtures"]) + workspace = Path(settings["workspace"]) + assert sys.version_info[:2] == tuple(settings["python"]) + assert not any(Path(p).resolve().is_relative_to(workspace) for p in sys.path) + # This directory contains copied test fixtures, never src/adcp or an adcp alias. + assert not (root / "adcp").exists() and not (root / "src").exists() + sys.path.insert(0, str(root)) + origins = {} + for name, expected in settings["modules"].items(): + path = Path(importlib.import_module(name).__file__).resolve() + assert "site-packages" in str(path) and not path.is_relative_to(workspace) + assert hashlib.sha256(path.read_bytes()).hexdigest() == expected + origins[name] = str(path) + for name, expected in settings["assets"].items(): + raw = files("adcp.reporting").joinpath(name).read_bytes() + assert hashlib.sha256(raw).hexdigest() == expected + if settings["driver_absent"]: + assert importlib.util.find_spec("psycopg") is None + assert importlib.util.find_spec("psycopg_pool") is None + os.environ.pop("ADCP_PG_TEST_URL", None) + evidence = Path(settings["evidence"]) + evidence.mkdir(parents=True, exist_ok=True, mode=0o700) + phases = [("green", None)] + if settings["parent"]: + phases = [ + ( + "negative", + "startup_then_sequential or concurrent_cold or mounted_unexpected" + " or original_pg_execute", + ), + ("preservation", "expected_closed or cancellation or actual_domain_rejections"), + ] + outputs = [] + for name, selection in phases: + log = evidence / f"{settings['label']}-{name}.log" + recorder = Results() + command = [ + str(root / "tests/conformance/reporting/test_reporting_activity_schema_proof.py"), + str(root / "tests/conformance/reporting/test_reporting_receipt_diagnostics.py"), + "-v", + "-s", + "-ra", + "-o", + "asyncio_mode=auto", + "-p", + "no:cacheprovider", + "--basetemp", + str(root / f"temp-{name}"), + ] + if selection: + command += ["-k", selection] + started = time.monotonic() + with ( + log.open("w") as stream, + contextlib.redirect_stdout(stream), + contextlib.redirect_stderr(stream), + ): + code = int(pytest.main(command, plugins=[recorder])) + runtime = round(time.monotonic() - started, 3) + valid = code == 0 and recorder.failed == recorder.errors == 0 and recorder.passed > 0 + if name == "negative": + valid = ( + code == 1 + and recorder.failed == 30 + and recorder.errors == recorder.skipped == 0 + and all( + ( + "assert 0 == 1" in reason + if "receipt_diagnostics" in node + else any( + f"assert ({scans}, {checkouts}) == (1, 1)" in reason + for scans, checkouts in ((9, 9), (27, 9), (12, 12), (36, 12)) + ) + ) + for node, reason in recorder.failures + ) + ) + result = { + "phase": name, + "command": command, + "pytest_exit": code, + "valid": valid, + "passed": recorder.passed, + "failed": recorder.failed, + "errors": recorder.errors, + "skipped": recorder.skipped, + "deselected": recorder.deselected, + "runtime_seconds": runtime, + "log_path": str(log), + "log_sha256": hashlib.sha256(log.read_bytes()).hexdigest(), + "log_bytes": log.stat().st_size, + } + outputs.append(result) + (evidence / f"{settings['label']}-{name}.json").write_text( + json.dumps(result, indent=2) + "\n" + ) + # Inspect all SDK modules loaded by pytest, not only the initial shortlist. + for name, module in tuple(sys.modules.items()): + if name == "adcp" or name.startswith("adcp."): + path = getattr(module, "__file__", None) + if path is not None: + assert Path(path).resolve().is_relative_to(Path(sys.prefix)) + print( + json.dumps( + { + "python": sys.version, + "origins": origins, + "assets": settings["assets"], + "results": outputs, + } + ) + ) + + +if __name__ == "__main__": + main(json.load(sys.stdin)) diff --git a/tests/conformance/reporting/_hardening_packaging.py b/tests/conformance/reporting/_hardening_packaging.py new file mode 100644 index 000000000..50d2752ff --- /dev/null +++ b/tests/conformance/reporting/_hardening_packaging.py @@ -0,0 +1,115 @@ +"""Installed hardening probes reuse fixtures, never current SDK source modules.""" + +import hashlib +import json +import os +import shutil +import subprocess +import zipfile +from pathlib import Path + +from .test_reporting_notification_packaging import ROOT, run_step + +MODULES = ( + "adcp.reporting.outbox.support", + "adcp.reporting.outbox._schema", + "adcp.reporting.outbox.status_schema", + "adcp.reporting.receipts.handler", + "adcp.reporting.receipts.pg", + "adcp.reporting.receipts._diagnostics", +) +ASSETS = ( + "outbox/required_schema.json", + "outbox/required_status_schema.json", + "ledger/reporting_feed.sql", + "feed/required_schema.json", + "ledger/reporting_materializer.sql", + "materializer/required_schema.json", + "ledger/reporting_receipt_ingestion.sql", + "receipts/required_schema.json", +) +B23 = "50e35f0ae3540f19b40e8fc460f5870dfe018bf9" + + +def installed_hardening(root, python, wheel, *, label, parent=False, driver_absent=False): + fixture_root = root / f"hardening-{label}" + fixture_root.mkdir(mode=0o700) + shutil.copytree( + ROOT / "tests", fixture_root / "tests", ignore=shutil.ignore_patterns("__pycache__") + ) + script = fixture_root / "run_installed.py" + shutil.copy2(Path(__file__).with_name("_hardening_installed.py"), script) + installer = ( + [shutil.which("uv"), "pip", "install", "--python", str(python)] + if shutil.which("uv") + else [str(python), "-m", "pip", "install"] + ) + run_step( + [ + *installer, + "pytest==9.1.1", + "pytest-asyncio==1.4.0", + "respx==0.23.1", + "asgi-lifespan==2.1.0", + ], + label=f"{label}-conformance-dependencies", + cwd=root, + timeout=180, + ) + with zipfile.ZipFile(wheel) as archive: + modules = { + name: hashlib.sha256(archive.read(name.replace(".", "/") + ".py")).hexdigest() + for name in MODULES + if name.replace(".", "/") + ".py" in archive.namelist() + } + assets = { + name: hashlib.sha256(archive.read("adcp/reporting/" + name)).hexdigest() + for name in ASSETS + } + for path in [ + *(name.replace(".", "/") + ".py" for name in modules), + *("adcp/reporting/" + name for name in assets), + ]: + expected = ( + subprocess.check_output(["git", "show", f"{B23}:src/{path}"], cwd=ROOT) + if parent + else (ROOT / "src" / path).read_bytes() + ) + assert archive.read(path) == expected + settings = { + "workspace": str(ROOT), + "fixtures": str(fixture_root), + "label": label, + "modules": modules, + "assets": assets, + "parent": parent, + "driver_absent": driver_absent, + "python": [3, 10], + "source": ( + B23 + if parent + else subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=ROOT, text=True).strip() + ), + "evidence": os.environ.get("ADCP_HARDENING_EVIDENCE", str(root / "hardening-evidence")), + } + result = json.loads( + run_step( + [str(python), "-I", str(script)], + label=f"{label}-installed-hardening", + cwd=fixture_root, + value=settings, + timeout=420, + ) + ) + print( + json.dumps( + { + "installed_hardening": label, + "wheel_sha256": hashlib.sha256(wheel.read_bytes()).hexdigest(), + **result, + } + ), + flush=True, + ) + assert all(item["valid"] for item in result["results"]), result["results"] + return result diff --git a/tests/conformance/reporting/test_reporting_activity_migration.py b/tests/conformance/reporting/test_reporting_activity_migration.py index 1f4bf0ed0..3f086572a 100644 --- a/tests/conformance/reporting/test_reporting_activity_migration.py +++ b/tests/conformance/reporting/test_reporting_activity_migration.py @@ -253,6 +253,8 @@ async def test_required_unusable_index_blocks_capability_boot(flag): worker, reliable.store, ReportingActivityProjector(outbox) ) assert await support.durable() + # Controlled DDL after a completed startup proof must invalidate it. + support.invalidate_schema_validation() # The task-owned PG16 admin fixture can model the catalog state left # by a failed concurrent index build without timing a real crash. async with reliable.blobs.pool.connection() as conn: diff --git a/tests/conformance/reporting/test_reporting_activity_schema_proof.py b/tests/conformance/reporting/test_reporting_activity_schema_proof.py new file mode 100644 index 000000000..b57217f03 --- /dev/null +++ b/tests/conformance/reporting/test_reporting_activity_schema_proof.py @@ -0,0 +1,522 @@ +"""Real catalog accounting through the supported activity/capability composition.""" + +import asyncio +from concurrent.futures import ThreadPoolExecutor +from contextlib import AsyncExitStack, asynccontextmanager, contextmanager +from dataclasses import FrozenInstanceError, replace + +import pytest + +from adcp.decisioning import ( + create_adcp_server_from_platform, + validate_capabilities_response_shape, + validate_capabilities_response_shape_async, +) +from adcp.decisioning.capabilities import Account, MediaBuy, WebhookSigning +from adcp.reporting.ledger.notification_models import ReportingNotificationError +from adcp.reporting.outbox import ( + ReportingActivityProjector, + ReportingActivitySupport, + ReportingNotificationWorker, + ReportingStatusSupport, +) +from adcp.types import ReportingDeliveryCapabilities +from tests.test_decisioning_capabilities_projection import _SalesPlatform +from tests.test_reporting_ledger import _OFFERING + +from ._reliable_support import NotificationHarness, reliable_factory +from .test_reporting_webhook_activity import worker_for + + +@contextmanager +def mounted_activity(support): + reporting = ReportingDeliveryCapabilities.model_validate( + { + "supported": True, + "configuration_task": "sync_accounts", + "status_task": "get_reporting_status", + "offerings": [_OFFERING], + "automated_recovery_window_seconds": 3600, + "status_retention_days": 30, + } + ).model_copy(update={"readiness_notification": None, "status_notification": None}) + + class Listing: + def list(self, filter=None): + return [] + + class Platform(_SalesPlatform): + capabilities = replace( + _SalesPlatform.capabilities, specialisms=[], supported_protocols=["media_buy"] + ) + accounts = Listing() + claim = True + claim_account = False + calls = 0 + + def get_adcp_capabilities_for_request(self, params=None, context=None): + self.calls += 1 + return replace( + self.capabilities, + experimental_features=["media_buy.reporting_delivery"], + media_buy=MediaBuy( + supported_pricing_models=["cpm"], + reporting_delivery=reporting.model_copy( + update={"supports_webhook_activity": self.claim} + ), + ), + account=Account.model_validate( + { + "supported_billing": ["operator"], + "notifications": { + "supported": True, + "registration_task": "sync_accounts", + "read_task": "list_accounts", + "event_types": ["account.status_changed"], + "supports_webhook_activity": self.claim_account, + }, + } + ), + webhook_signing=WebhookSigning( + supported=True, + profile="adcp/webhook-signing/v1", + algorithms=["ed25519"], + delivery_retry_horizon_seconds=86400, + ), + webhook_signing_managed_externally=True, + ) + + handler, executor, _ = create_adcp_server_from_platform( + Platform(), + reporting_activity=support, + account_activity=support.projector, + auto_emit_task_webhooks=False, + validate_at_init=False, + ) + try: + yield handler + finally: + executor.shutdown(wait=True) + + +async def compose_activity(reliable, composite): + n = NotificationHarness(reliable) + outbox = n.outbox + worker = worker_for(n, outbox) + status = None + reader = outbox + if composite: + from adcp.reporting.outbox.status_activity_pg import PgReportingActivityUnionStore + from adcp.reporting.outbox.status_pg import PgStatusNotificationStore + + store = PgStatusNotificationStore(reliable.store) + await store.create_schema() + c_worker = ReportingNotificationWorker( + outbox=store.outbox, + activity=store.outbox, + subscriptions=n.subscriptions, + signing=n.signing, + cipher=n.cipher, + ) + status = ReportingStatusSupport(store, c_worker, scheduled=True) + reader = PgReportingActivityUnionStore(outbox, store.outbox) + return ReportingActivitySupport( + worker, reliable.store, ReportingActivityProjector(reader), status + ) + + +class CatalogAccounting: + """Instrument real psycopg execution and acquisition, never replace a validator.""" + + def __init__(self, monkeypatch, pool): + from psycopg import AsyncConnection + + self.checkouts = 0 + self.scans = 0 + self.catalog_queries = 0 + self.pause = False + self.entered = asyncio.Event() + self.release = asyncio.Event() + connection = pool.connection + execute = AsyncConnection.execute + + @asynccontextmanager + async def checkout(*args, **kwargs): + async with connection(*args, **kwargs) as conn: + self.checkouts += 1 + yield conn + + async def query(conn, command, *args, **kwargs): + if isinstance(command, str): + if command.startswith("SELECT c.oid, c.relname, c.relkind"): + self.scans += 1 + if any(name in command for name in ("pg_class", "pg_attribute", "pg_proc")): + self.catalog_queries += 1 + result = await execute(conn, command, *args, **kwargs) + # Pause after the final catalog result is captured. An invalidated + # old scan can still be positive, so epoch publication is exercised. + if self.pause and isinstance(command, str) and "FROM pg_proc p" in command: + self.pause = False + self.entered.set() + await self.release.wait() + return result + + monkeypatch.setattr(pool, "connection", checkout) + monkeypatch.setattr(AsyncConnection, "execute", query) + + +@pytest.fixture(params=[False, True], ids=["b", "b+c"]) +async def activity_proof(request, monkeypatch): + async with reliable_factory("postgres", notifications=True, autocommit=True) as reliable: + support = await compose_activity(reliable, request.param) + accounting = CatalogAccounting(monkeypatch, reliable.blobs.pool) + with mounted_activity(support) as handler: + yield support, handler, accounting, reliable.blobs.pool + + +def assert_primitive_cache(support): + def primitive(value): + if type(value) in (str, int, bool, type(None)): + return True + return type(value) is tuple and all(primitive(v) for v in value) + + assert all(primitive(value) for value in vars(support._schema_validation).values()) + + +async def test_startup_then_sequential_discovery_checks_catalog_and_pool_once(activity_proof): + support, handler, accounting, _ = activity_proof + await validate_capabilities_response_shape_async(handler) + queries = accounting.catalog_queries + for _ in range(8): + response = await handler.get_adcp_capabilities() + assert response["media_buy"]["reporting_delivery"]["supports_webhook_activity"] is True + assert handler._platform.calls == 9 + assert (accounting.scans, accounting.checkouts) == (1, 1) + assert accounting.catalog_queries == queries and queries > 10 + assert_primitive_cache(support) + + +async def test_concurrent_cold_discovery_single_flights_real_catalog(activity_proof): + support, handler, accounting, _ = activity_proof + responses = await asyncio.gather(*(handler.get_adcp_capabilities() for _ in range(12))) + assert len(responses) == handler._platform.calls == 12 + assert (accounting.scans, accounting.checkouts) == (1, 1) + assert_primitive_cache(support) + + +async def test_request_false_claim_does_not_warm_schema_proof(activity_proof): + support, handler, accounting, _ = activity_proof + handler._platform.claim = False + for _ in range(3): + await handler.get_adcp_capabilities() + assert (accounting.scans, accounting.checkouts) == (0, 0) + handler._platform.claim = True + await handler.get_adcp_capabilities() + handler._platform.claim = False + await handler.get_adcp_capabilities() + handler._platform.claim = True + await handler.get_adcp_capabilities() + assert (accounting.scans, accounting.checkouts) == (1, 1) + assert_primitive_cache(support) + + +async def test_warm_core_discovery_does_not_queue_behind_saturated_pool(activity_proof): + _, handler, accounting, pool = activity_proof + await validate_capabilities_response_shape_async(handler) + async with AsyncExitStack() as stack: + for _ in range(pool.max_size): + await stack.enter_async_context(pool.connection()) + acquired = accounting.checkouts + await asyncio.wait_for(handler.get_adcp_capabilities(), timeout=1) + assert accounting.checkouts == acquired + assert accounting.scans == 1 + + +async def test_failed_schema_is_retried_after_repair(activity_proof): + support, handler, accounting, pool = activity_proof + async with pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_webhook_attempts DISABLE TRIGGER reporting_webhook_attempt_guard" + ) + for _ in range(2): + with pytest.raises(ReportingNotificationError, match="notification_schema_unready"): + await handler.get_adcp_capabilities() + assert_primitive_cache(support) + assert accounting.scans == 2 + async with pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_webhook_attempts ENABLE TRIGGER reporting_webhook_attempt_guard" + ) + await handler.get_adcp_capabilities() + await handler.get_adcp_capabilities() + assert accounting.scans == 3 + + +async def test_explicit_invalidation_rechecks_broken_required_object(activity_proof): + support, handler, accounting, pool = activity_proof + await handler.get_adcp_capabilities() + support.invalidate_schema_validation() + async with pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_webhook_attempts DISABLE TRIGGER reporting_webhook_attempt_guard" + ) + with pytest.raises(ReportingNotificationError, match="notification_schema_unready"): + await handler.get_adcp_capabilities() + assert accounting.scans == 2 + assert_primitive_cache(support) + + +async def test_old_positive_scan_cannot_publish_after_epoch_invalidation(activity_proof): + support, handler, accounting, pool = activity_proof + accounting.pause = True + old = asyncio.create_task(handler.get_adcp_capabilities()) + try: + await asyncio.wait_for(accounting.entered.wait(), timeout=10) + support.invalidate_schema_validation() + async with pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_webhook_attempts DISABLE TRIGGER" + " reporting_webhook_attempt_guard" + ) + accounting.release.set() + with pytest.raises(ReportingNotificationError): + await old + assert_primitive_cache(support) + with pytest.raises(ReportingNotificationError, match="notification_schema_unready"): + await handler.get_adcp_capabilities() + assert accounting.scans == 2 + finally: + accounting.release.set() + await asyncio.gather(old, return_exceptions=True) + + +@pytest.mark.parametrize("mutation", ["recorder", "reader", "pool", "notifications", "listing"]) +async def test_warm_proof_does_not_bypass_dynamic_topology(activity_proof, mutation): + support, handler, accounting, _ = activity_proof + handler._platform.claim_account = True + await handler.get_adcp_capabilities() + if mutation == "recorder": + support.worker.activity = None + elif mutation == "reader": + handler._account_activity = ReportingActivityProjector(support.projector.store) + elif mutation == "pool": + support.worker.outbox._pool = object() + elif mutation == "notifications": + support.ledger._notifications_enabled = False + else: + handler._platform.accounts.list = None + with pytest.raises(ReportingNotificationError): + await handler.get_adcp_capabilities() + assert accounting.scans == 1 + + +async def test_equal_frozen_support_instances_have_independent_proof(activity_proof): + support, handler, accounting, _ = activity_proof + equal = replace(support) + assert equal == support + with pytest.raises(FrozenInstanceError): + support.projector = None + await handler.get_adcp_capabilities() + with mounted_activity(equal) as other: + await other.get_adcp_capabilities() + await other.get_adcp_capabilities() + assert equal == support + assert (accounting.scans, accounting.checkouts) == (2, 2) + assert equal._schema_validation is not support._schema_validation + + +async def test_synchronous_startup_then_runtime_loop_has_no_loop_bound_cache(activity_proof): + support, handler, accounting, _ = activity_proof + # The supported synchronous validator really calls asyncio.run in a + # separate thread. The runtime loop continues to own the live PG pool. + with ThreadPoolExecutor(max_workers=1) as startup: + await asyncio.wrap_future(startup.submit(validate_capabilities_response_shape, handler)) + assert_primitive_cache(support) + await handler.get_adcp_capabilities() + assert (accounting.scans, accounting.checkouts) == (1, 1) + + +async def test_memory_only_claims_and_frozen_constructor_remain_unchanged(): + async with reliable_factory("memory", notifications=True) as reliable: + support = await compose_activity(reliable, False) + assert not await support.durable() + assert not await support.durable() + with mounted_activity(support) as handler: + handler._platform.claim = False + await handler.get_adcp_capabilities() + handler._platform.claim = True + with pytest.raises(ReportingNotificationError, match="requires_durable_reporting"): + await handler.get_adcp_capabilities() + + +async def test_cancelled_waiter_leaves_other_cold_callers_single_flight(activity_proof): + support, handler, accounting, _ = activity_proof + accounting.pause = True + first = asyncio.create_task(handler.get_adcp_capabilities()) + second = asyncio.create_task(handler.get_adcp_capabilities()) + try: + await asyncio.wait_for(accounting.entered.wait(), 10) + first.cancel() + with pytest.raises(asyncio.CancelledError): + await first + assert support._schema_validation.positive is None + accounting.release.set() + await second + assert (accounting.scans, accounting.checkouts) == (1, 1) + assert_primitive_cache(support) + finally: + accounting.release.set() + await asyncio.gather(first, second, return_exceptions=True) + + +async def test_cancelled_scan_and_driver_failure_never_become_positive(activity_proof, monkeypatch): + from psycopg import AsyncConnection, OperationalError + + support, handler, accounting, _ = activity_proof + accounting.pause = True + call = asyncio.create_task(handler.get_adcp_capabilities()) + try: + await asyncio.wait_for(accounting.entered.wait(), 10) + support._schema_validation.flight.task.cancel() + with pytest.raises(asyncio.CancelledError): + await call + finally: + accounting.release.set() + await asyncio.gather(call, return_exceptions=True) + assert support._schema_validation.positive is None + assert_primitive_cache(support) + original = AsyncConnection.execute + fired = [] + + async def fail_after_catalog_query(connection, command, *args, **kwargs): + result = await original(connection, command, *args, **kwargs) + if isinstance(command, str) and command.startswith("SELECT c.oid, c.relname, c.relkind"): + fired.append(True) + raise OperationalError("not retained by schema validation") + return result + + with monkeypatch.context() as patch: + patch.setattr(AsyncConnection, "execute", fail_after_catalog_query) + with pytest.raises(ReportingNotificationError, match="catalog_unavailable"): + await handler.get_adcp_capabilities() + assert fired == [True] and support._schema_validation.positive is None + assert_primitive_cache(support) + await handler.get_adcp_capabilities() + await handler.get_adcp_capabilities() + assert (accounting.scans, accounting.checkouts) == (3, 3) + + +@pytest.mark.parametrize("mutation", ["status-pool", "status-clock", "status-writer", "scheduling"]) +async def test_composite_warm_proof_rechecks_closed_union_wiring(monkeypatch, mutation): + async with reliable_factory("postgres", notifications=True, autocommit=True) as reliable: + support = await compose_activity(reliable, True) + accounting = CatalogAccounting(monkeypatch, reliable.blobs.pool) + with mounted_activity(support) as handler: + await handler.get_adcp_capabilities() + if mutation == "status-pool": + support.status.store.outbox._pool = object() + elif mutation == "status-clock": + support.status.store.outbox._clock = object() + elif mutation == "status-writer": + support.status.worker.activity = None + else: + # The public dataclass remains frozen; adversarial private + # mutation must still not turn an old proof into scheduling. + object.__setattr__(support.status, "scheduled", False) + with pytest.raises(ReportingNotificationError): + await handler.get_adcp_capabilities() + assert accounting.scans == 1 + + +async def test_separate_search_paths_and_instances_do_not_share_positive_proof(monkeypatch): + async with ( + reliable_factory("postgres", notifications=True, autocommit=True) as first, + reliable_factory("postgres", notifications=True, autocommit=True) as second, + ): + one = await compose_activity(first, False) + two = await compose_activity(second, False) + assert first.blobs.pool.conninfo == second.blobs.pool.conninfo + assert first.blobs.pool.kwargs["options"] != second.blobs.pool.kwargs["options"] + accounting = CatalogAccounting(monkeypatch, first.blobs.pool) + with mounted_activity(one) as a, mounted_activity(two) as b: + await a.get_adcp_capabilities() + async with second.blobs.pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_webhook_attempts DISABLE TRIGGER" + " reporting_webhook_attempt_guard" + ) + with pytest.raises(ReportingNotificationError, match="notification_schema_unready"): + await b.get_adcp_capabilities() + async with second.blobs.pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_webhook_attempts ENABLE TRIGGER" + " reporting_webhook_attempt_guard" + ) + await b.get_adcp_capabilities() + await a.get_adcp_capabilities() + await b.get_adcp_capabilities() + assert accounting.scans == 3 and accounting.checkouts == 1 + + +async def test_b_and_composite_prove_each_packaged_contract_once(monkeypatch): + from adcp.reporting.outbox import _schema, status_schema + + async with reliable_factory("postgres", notifications=True, autocommit=True) as reliable: + composite = await compose_activity(reliable, True) + b_only = ReportingActivitySupport( + composite.worker, + composite.ledger, + ReportingActivityProjector(composite.worker.outbox), + ) + accounting = CatalogAccounting(monkeypatch, reliable.blobs.pool) + calls = [] + b_validator, c_validator = ( + _schema._validate_schema_objects, + status_schema._validate_status_objects, + ) + + def b_contract(installed, **options): + calls.append(("B", options)) + b_validator(installed, **options) + + def c_contract(installed, **options): + calls.append(("C", options)) + c_validator(installed, **options) + + monkeypatch.setattr(_schema, "_validate_schema_objects", b_contract) + monkeypatch.setattr(status_schema, "_validate_status_objects", c_contract) + for support in (b_only, composite): + assert await support.durable() + assert await support.durable() + assert calls == [ + ("B", {"activity": True}), + ("B", {"activity": True}), + ("C", {"activity": True, "status": False}), + ] + assert (accounting.scans, accounting.checkouts) == (2, 2) + assert b_only._schema_validation.positive != composite._schema_validation.positive + async with reliable.blobs.pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_status_webhook_attempts DISABLE TRIGGER" + " reporting_status_webhook_attempt_guard" + ) + assert await b_only.durable() + # A new server instance is the normal migration/restart boundary. + restarted = replace(composite) + with pytest.raises(ReportingNotificationError, match="status_schema_unready"): + await restarted.durable() + assert accounting.scans == 3 + + +async def test_positive_proof_is_bound_to_packaged_manifest_contract(activity_proof, monkeypatch): + from adcp.reporting.outbox import _schema + + support, handler, accounting, _ = activity_proof + await handler.get_adcp_capabilities() + key = "table:reporting_webhook_attempts" + with monkeypatch.context() as patch: + patch.setitem(_schema.REQUIRED_OBJECTS, key, {"fingerprint": "0" * 64, "enabled": True}) + with pytest.raises(ReportingNotificationError, match="notification_schema_unready:changed"): + await handler.get_adcp_capabilities() + assert accounting.scans == 2 + assert_primitive_cache(support) diff --git a/tests/conformance/reporting/test_reporting_feed_hardening_installed.py b/tests/conformance/reporting/test_reporting_feed_hardening_installed.py new file mode 100644 index 000000000..6010c15cf --- /dev/null +++ b/tests/conformance/reporting/test_reporting_feed_hardening_installed.py @@ -0,0 +1,191 @@ +"""Actual approved B2.3 and current floor wheels at both changed boundaries.""" + +import asyncio +import hashlib +import json +import os +import shutil +import subprocess +import zipfile +from pathlib import Path + +import pytest + +from ._feed_support import feed_harness, feed_request, mixed_case, walk, without_feed +from ._hardening_packaging import ASSETS, B23, MODULES, installed_hardening +from .test_reporting_feed_installed_pg import ( + b1_wheels, + built_distribution, + feed_wheels, + installed_feed, +) +from .test_reporting_feed_packaging import feed_modules +from .test_reporting_feed_process import feed_process +from .test_reporting_materializer_rolling import build_frozen +from .test_reporting_notification_packaging import ROOT, run_step + +__all__ = ["b1_wheels", "built_distribution", "feed_wheels", "installed_feed"] + + +async def test_installed_python310_proof_and_receipt_diagnostics(installed_feed, feed_wheels): + root, python, _, _, installed = installed_feed + _, wheels, _ = feed_wheels + wheel = wheels[installed["distribution"]] + await asyncio.to_thread( + installed_hardening, root, python, wheel, label=installed["distribution"] + "-pg" + ) + + +@pytest.fixture(scope="module") +def approved_b23(tmp_path_factory, request): + # Preserve all nine original frozen inputs; this is an additional artifact. + root, _, _, identity = build_frozen("b23-hardening-control", tmp_path_factory, request, sha=B23) + interpreter = os.environ.get("ADCP_PYTHON310") + if interpreter is None: + pytest.skip("ADCP_PYTHON310 supplies the installed floor cell") + environment = root / "python310" + run_step( + [interpreter, "-m", "venv", str(environment)], label="b23-python310-environment", cwd=root + ) + python = environment / "bin/python" + wheel = next((root / "dist").glob("*.whl")) + installer = ( + [shutil.which("uv"), "pip", "install", "--python", str(python)] + if shutil.which("uv") + else [str(python), "-m", "pip", "install"] + ) + run_step( + [*installer, f"{wheel}[pg]", "asgi-lifespan==2.1.0"], + label="b23-python310-wheel-install", + cwd=root, + timeout=180, + ) + with zipfile.ZipFile(wheel) as archive: + modules = {} + for name in set(feed_modules()) | (set(MODULES) - {"adcp.reporting.receipts._diagnostics"}): + path = name.replace(".", "/") + ".py" + if path not in archive.namelist(): + path = name.replace(".", "/") + "/__init__.py" + raw = archive.read(path) + assert raw == subprocess.check_output(["git", "show", f"{B23}:src/{path}"], cwd=ROOT) + modules[name] = hashlib.sha256(raw).hexdigest() + assets = {} + for name in ASSETS: + raw = archive.read("adcp/reporting/" + name) + assert raw == subprocess.check_output( + ["git", "show", f"{B23}:src/adcp/reporting/{name}"], cwd=ROOT + ) + assets[name] = hashlib.sha256(raw).hexdigest() + script, helper = root / "feed_process.py", root / "receipt_transport.py" + shutil.copy2(Path(__file__).with_name("_feed_process.py"), script) + shutil.copy2(Path(__file__).with_name("_receipt_transport.py"), helper) + installed = { + **identity, + "modules": modules, + "assets": assets, + "python": [3, 10], + "tree": "16c55b24340aeea3524c781f19bfa64d1491ac34", + } + return root, python, script, helper, installed, wheel + + +def test_actual_approved_b23_installed_negative_and_preservation_controls(approved_b23): + root, python, _, _, identity, wheel = approved_b23 + result = installed_hardening(root, python, wheel, label="b23-parent", parent=True) + print( + json.dumps({"approved_b23_installed_control": identity, "results": result["results"]}), + flush=True, + ) + + +@pytest.mark.parametrize("notifications", [False, True]) +async def test_actual_b23_to_child_installed_restart_preserves_pages_and_receipt_replay( + approved_b23, installed_feed, notifications +): + _, parent_python, parent_script, parent_helper, parent, _ = approved_b23 + root, python, script, helper, current = installed_feed + old = { + "python": parent_python, + "script": parent_script, + "helper": parent_helper, + "installed": parent, + } + new = {"python": python, "script": script, "helper": helper, "installed": current} + async with feed_harness("postgres", notifications=notifications) as h: + s, receipt_request, receipt_response = await mixed_case(h) + # The approved binary really mounts the ingress and returns its durable + # replay before it writes page one; fixtures only supply populated data. + async with feed_process(h, s, receipt_request, action="receipt", **old) as child: + admitted = await child.event("done") + assert await asyncio.wait_for(child.process.wait(), 5) == 0 + assert admitted["result"] == receipt_response + async with feed_process(h, s, feed_request(s), pause="committed", **old) as child: + first = (await child.event("committed"))["result"] + await child.kill() + original = await h.store.read_reporting_feed_snapshot( + first["ledger_snapshot_id"], caller=s.binding.principal + ) + expected = await walk(h.store, feed_request(s), s.binding.principal, first=first) + await h.store.set_revision_readable( + account_id=s.obligation.account_id, + reporting_revision_id=s.revision.reporting_revision_id, + readable=False, + ) + # Normal stop/migrate/restart uses the current installed migration path. + ready = json.loads( + await asyncio.to_thread( + run_step, + [str(python), "-I", str(script)], + label="b23-to-child-installed-restart", + cwd=root, + value={ + "conninfo": h.pool.conninfo, + "kwargs": h.pool.kwargs, + "notifications": notifications, + "action": "install", + "installed": current, + }, + timeout=90, + ) + ) + assert ready["result"]["feed_objects"] == 33 + before = without_feed(await h.image()) + continuation = feed_request( + s, pagination={"cursor": first["pagination"]["cursor"], "max_results": 1} + ) + for v1 in (False, True): + async with feed_process( + h, s, continuation, action="walk", transport="a2a", v1=v1, **new + ) as child: + continued = await child.event("done") + assert await asyncio.wait_for(child.process.wait(), 5) == 0 + assert continued["result"]["pages"] == expected[0][1:] + assert continued["result"]["binding"] == original.binding + assert continued["result"]["version"] == original.representation_version + assert continued["result"]["ownership_mode"] == original.ownership_mode == "absent" + async with feed_process(h, s, receipt_request, action="receipt", **new) as child: + replayed = await child.event("done") + assert await asyncio.wait_for(child.process.wait(), 5) == 0 + assert replayed["result"] == receipt_response + assert without_feed(await h.image()) == before + assert ( + await h.store.read_reporting_feed_snapshot( + first["ledger_snapshot_id"], caller=s.binding.principal + ) + == original + ) + print( + json.dumps( + { + "b23_to_child_restart": current["distribution"], + "notifications": notifications, + "parent": parent, + "current": current, + "parent_origins": admitted["origins"], + "current_origins": replayed["origins"], + "page_count": len(expected[0]), + "checkpoint": first["changes_checkpoint"], + } + ), + flush=True, + ) diff --git a/tests/conformance/reporting/test_reporting_feed_packaging.py b/tests/conformance/reporting/test_reporting_feed_packaging.py index 332c04525..36fef48c3 100644 --- a/tests/conformance/reporting/test_reporting_feed_packaging.py +++ b/tests/conformance/reporting/test_reporting_feed_packaging.py @@ -140,3 +140,6 @@ def test_python310_feed_without_pg_exports_sql_and_strict_adopter(request, kind) ), flush=True, ) + from ._hardening_packaging import installed_hardening + + installed_hardening(root, python, wheels[kind], label=kind + "-base", driver_absent=True) diff --git a/tests/conformance/reporting/test_reporting_receipt_diagnostics.py b/tests/conformance/reporting/test_reporting_receipt_diagnostics.py new file mode 100644 index 000000000..d6c20711b --- /dev/null +++ b/tests/conformance/reporting/test_reporting_receipt_diagnostics.py @@ -0,0 +1,452 @@ +"""Unexpected original failures are useful to operators without retaining payloads.""" + +import asyncio +import json +import logging +from types import SimpleNamespace + +import pytest + +from adcp.exceptions import ADCPTaskError +from adcp.reporting.receipts import ReportingReceiptError, ReportingReceiptHandler +from adcp.server.base import ToolContext + +from ._receipt_support import receipt_case, receipt_harness, request_for +from ._receipt_transport import MountedReceipts, error_code + +SECRETS = ( + "unexpected-request-secret-64ee282a", + "private-account-64ee282a", + "private-consumer-64ee282a", + "private-idempotency-64ee282a", + "private-continuation-64ee282a", + "private-receipt-64ee282a", + "postgresql://private-auth-token-64ee282a@provider.invalid/financial", + "SELECT private_financial_value_64ee282a FROM provider_payload", +) +SAFE_MESSAGE = "receipt storage is unavailable; retry the same batch and key" +DIAGNOSTIC_FIELDS = { + "code", + "boundary", + "exception_type", + "origin_module", + "origin_function", + "origin_line", +} + + +def operator_errors(caplog): + return [record for record in caplog.records if record.levelno >= logging.ERROR] + + +def assert_diagnostic(caplog, boundary, exception_type, *, origin_function=None): + records = operator_errors(caplog) + assert len(records) == 1 + record = records[0] + assert record.name == "adcp.reporting.receipts" + assert record.code == "RECEIPT_STORAGE_UNAVAILABLE" + assert record.boundary == boundary + assert record.exception_type == exception_type + assert record.origin_module == __name__ + if origin_function is not None: + assert record.origin_function == origin_function + assert type(record.origin_line) is int and record.origin_line > 0 + assert record.args == () and record.exc_info is None and record.exc_text is None + assert record.stack_info is None + standard = set(logging.makeLogRecord({}).__dict__) | {"message", "asctime"} + assert set(record.__dict__) - standard == DIAGNOSTIC_FIELDS + # No repr fallback: every retained value must itself be JSON serializable. + serialized = json.dumps(record.__dict__, sort_keys=True) + for secret in SECRETS: + assert secret not in serialized and secret not in caplog.text + return record + + +async def sensitive_case(h): + s = await receipt_case(h, account_id=SECRETS[1], consumer_id=SECRETS[2]) + request = request_for(s, key=SECRETS[3]) + request["receipts"][0]["reporting_receipt_id"] = SECRETS[5] + request["context"] = {"request_secret": SECRETS[0], "continuation": SECRETS[4]} + return s, request + + +async def mounted_call(mount, transport, request): + async with mount.client() as client: + if transport == "mcp": + return await mount.mcp(client, request) + return await mount.a2a(client, request, v1=transport == "a2a-1.0") + + +def assert_safe_wire(status, payload, request): + assert status == 200 + assert error_code(payload) == "RECEIPT_STORAGE_UNAVAILABLE" + error = payload["adcp_error"] if "adcp_error" in payload else payload["errors"][0] + expected = ( + "sync_reporting_receipts failed: " + SAFE_MESSAGE + if "adcp_error" in payload + else SAFE_MESSAGE + ) + assert error["message"] == expected + # The existing transport echoes caller context. It must remain unchanged; + # only the safe error, not the already supplied context, is diagnostic text. + assert payload.get("context") == request.get("context") + assert set(payload) <= {"context", "adcp_error", "errors"} + serialized = json.dumps(error, sort_keys=True) + for secret in SECRETS: + assert secret not in serialized + + +@pytest.mark.parametrize("transport", ["mcp", "a2a-0.3", "a2a-1.0"]) +@pytest.mark.parametrize("boundary", ["resolver", "custom-store"]) +async def test_mounted_unexpected_original_failure_logs_once_and_keeps_safe_wire( + transport, boundary, caplog +): + caplog.set_level(logging.ERROR) + async with receipt_harness("memory") as h: + s, request = await sensitive_case(h) + + class CustomStore: + async def ingest_receipt_batch(self, request, *, caller): + raise RuntimeError(" | ".join(SECRETS)) + + mount = MountedReceipts( + h if boundary == "resolver" else SimpleNamespace(store=CustomStore()) + ) + mount.authorize(s) + + async def resolver(reference, context, consumer): + # The private chain must not survive on the buyer error or record. + try: + raise ValueError(SECRETS[6]) + except ValueError as cause: + raise RuntimeError(" | ".join(SECRETS)) from cause + + if boundary == "resolver": + mount.handler._receipt_account_resolver = resolver + status, payload = await mounted_call(mount, transport, request) + assert_safe_wire(status, payload, request) + assert_diagnostic( + caplog, + "handler", + "RuntimeError", + origin_function="resolver" if boundary == "resolver" else "ingest_receipt_batch", + ) + + +def inject_driver_failure(monkeypatch, point, failure): + from psycopg import AsyncConnection + + fired = [] + original_execute = AsyncConnection.execute + original_command = AsyncConnection._exec_command + + async def execute(connection, query, *args, **kwargs): + result = await original_execute(connection, query, *args, **kwargs) + if ( + not fired + and isinstance(query, str) + and query.startswith("INSERT INTO reporting_receipt_ingestion_results") + ): + fired.append(point) + raise failure + return result + + def commit(connection, command, *args, **kwargs): + if not fired and command == b"COMMIT": + fired.append(point) + raise failure + return (yield from original_command(connection, command, *args, **kwargs)) + + if point == "execute": + monkeypatch.setattr(AsyncConnection, "execute", execute) + else: + monkeypatch.setattr(AsyncConnection, "_exec_command", commit) + return fired + + +@pytest.mark.parametrize("notifications", [False, True]) +@pytest.mark.parametrize("point", ["execute", "commit"]) +@pytest.mark.parametrize("transport", ["store", "handler", "mcp", "a2a-0.3", "a2a-1.0"]) +async def test_original_pg_execute_and_commit_failures_log_once_across_translation( + notifications, point, transport, monkeypatch, caplog +): + caplog.set_level(logging.ERROR) + async with receipt_harness("postgres", notifications=notifications) as h: + from psycopg import OperationalError + + s, request = await sensitive_case(h) + before = await h.image() + mount = MountedReceipts(h) + mount.authorize(s) + with monkeypatch.context() as patch: + fired = inject_driver_failure(patch, point, OperationalError(" | ".join(SECRETS))) + if transport in {"store", "handler"}: + error_type = ReportingReceiptError if transport == "store" else ADCPTaskError + with pytest.raises(error_type) as caught: + if transport == "store": + await h.store.ingest_receipt_batch(request, caller=s.binding.principal) + else: + await mount.handler.sync_reporting_receipts( + request, ToolContext(caller_identity=s.binding.consumer_id) + ) + assert caught.value.__cause__ is None and caught.value.__context__ is None + if transport == "store": + assert caught.value.code == "RECEIPT_STORAGE_UNAVAILABLE" + assert str(caught.value) == SAFE_MESSAGE + else: + assert caught.value.errors[0].code == "RECEIPT_STORAGE_UNAVAILABLE" + assert caught.value.errors[0].message == SAFE_MESSAGE + else: + status, payload = await mounted_call(mount, transport, request) + assert_safe_wire(status, payload, request) + assert fired == [point] + assert_diagnostic( + caplog, "store.ingest_receipt_batch", "OperationalError", origin_function=point + ) + assert await h.image() == before + assert (await h.store.ingest_receipt_batch(request, caller=s.binding.principal))["results"][ + 0 + ]["result"] == "recorded" + + +@pytest.mark.parametrize( + "code", + [ + "INVALID_REQUEST", + "UNAUTHORIZED", + "IDEMPOTENCY_CONFLICT", + "RECEIPT_SCHEMA_UNREADY", + "RECEIPT_HISTORY_CORRUPT", + ], +) +@pytest.mark.parametrize("transport", ["mcp", "a2a-0.3", "a2a-1.0"]) +async def test_expected_closed_receipt_errors_are_silent_on_mounted_transports( + code, transport, caplog +): + caplog.set_level(logging.ERROR) + async with receipt_harness("memory") as h: + s, request = await sensitive_case(h) + + class ExpectedFailureStore: + async def ingest_receipt_batch(self, request, *, caller): + raise ReportingReceiptError(code) + + mount = MountedReceipts(SimpleNamespace(store=ExpectedFailureStore())) + mount.authorize(s) + status, payload = await mounted_call(mount, transport, request) + assert status == 200 and error_code(payload) == code + assert operator_errors(caplog) == [] + + +@pytest.mark.parametrize("boundary", ["resolver", "custom-store", "pg"]) +async def test_cancellation_propagates_without_translation_or_diagnostic( + boundary, caplog, monkeypatch +): + caplog.set_level(logging.ERROR) + async with receipt_harness("postgres" if boundary == "pg" else "memory") as h: + s, request = await sensitive_case(h) + mount = MountedReceipts(h) + mount.authorize(s) + + async def cancel(*args, **kwargs): + raise asyncio.CancelledError(SECRETS[0]) + + with monkeypatch.context() as patch: + if boundary == "resolver": + mount.handler._receipt_account_resolver = cancel + elif boundary == "custom-store": + patch.setattr(h.store, "ingest_receipt_batch", cancel) + else: + inject_driver_failure(patch, "execute", asyncio.CancelledError(SECRETS[0])) + with pytest.raises(asyncio.CancelledError): + await mount.handler.sync_reporting_receipts( + request, ToolContext(caller_identity=s.binding.consumer_id) + ) + assert operator_errors(caplog) == [] + + +async def test_unexpected_exception_is_never_formatted_and_translation_context_is_empty(caplog): + caplog.set_level(logging.ERROR) + + class UnformattableError(RuntimeError): + def __str__(self): + raise AssertionError("unexpected exception was stringified") + + def __repr__(self): + raise AssertionError("unexpected exception was represented") + + async with receipt_harness("memory") as h: + _, request = await sensitive_case(h) + + async def resolver(reference, context, consumer): + raise UnformattableError(*SECRETS) + + handler = ReportingReceiptHandler(h.store, resolve_account=resolver) + with pytest.raises(ADCPTaskError) as caught: + await handler.sync_reporting_receipts(request, ToolContext(caller_identity=SECRETS[2])) + assert caught.value.__context__ is None and caught.value.__cause__ is None + assert_diagnostic(caplog, "handler", "UnformattableError", origin_function="resolver") + + +@pytest.mark.parametrize("transport", ["mcp", "a2a-0.3", "a2a-1.0"]) +@pytest.mark.parametrize("notifications", [False, True]) +@pytest.mark.parametrize( + "code", + [ + "INVALID_REQUEST", + "UNAUTHORIZED", + "IDEMPOTENCY_CONFLICT", + "RECEIPT_SCHEMA_UNREADY", + "RECEIPT_HISTORY_CORRUPT", + ], +) +async def test_actual_domain_rejections_do_not_emit_operator_errors( + transport, notifications, code, caplog +): + caplog.set_level(logging.ERROR) + async with receipt_harness("postgres", notifications=notifications) as h: + s, request = await sensitive_case(h) + mount = MountedReceipts(h) + mount.authorize(s) + if code != "INVALID_REQUEST": + await h.store.ingest_receipt_batch(request, caller=s.binding.principal) + if code == "INVALID_REQUEST": + request["receipts"][0]["received_at"] = "2026-09-18T00:00:00Z" + elif code == "UNAUTHORIZED": + mount.grants.clear() + elif code == "IDEMPOTENCY_CONFLICT": + request["context"]["changed"] = True + elif code == "RECEIPT_SCHEMA_UNREADY": + async with h.pool.connection() as connection: + await connection.execute( + "ALTER TABLE reporting_receipt_ingestion_results" + " DISABLE TRIGGER reporting_receipt_ingestion_result" + ) + else: + # Privileged corruption fixture, with all schema guards restored + # before invoking the production decoder. This is a domain error. + from psycopg.types.json import Jsonb + + from adcp.reporting.canonical_json import canonical_json_sha256_v1 + + async with h.pool.connection() as connection, connection.transaction(): + value = ( + await ( + await connection.execute( + "SELECT final_response FROM reporting_receipt_ingestion_batches" + ) + ).fetchone() + )[0] + value["results"][0]["receipt"]["reporting_receipt_id"] = "corrupt-other-receipt" + await connection.execute("SET LOCAL session_replication_role = replica") + await connection.execute( + "UPDATE reporting_receipt_ingestion_batches" + " SET final_response=%s,final_sha256=%s", + (Jsonb(value), canonical_json_sha256_v1(value)), + ) + before = await h.image() + status, payload = await mounted_call(mount, transport, request) + assert status == 200 and error_code(payload) == code + assert operator_errors(caplog) == [] + assert await h.image() == before + + +async def test_diagnostic_does_not_capture_task_names_or_ambient_record_context(caplog): + caplog.set_level(logging.ERROR) + original_factory = logging.getLogRecordFactory() + task = asyncio.current_task() + original_name = task.get_name() + + def request_factory(*args, **kwargs): + record = original_factory(*args, **kwargs) + record.request_context = SECRETS + return record + + async with receipt_harness("memory") as h: + _, request = await sensitive_case(h) + + async def resolver(reference, context, consumer): + raise RuntimeError(*SECRETS) + + handler = ReportingReceiptHandler(h.store, resolve_account=resolver) + try: + logging.setLogRecordFactory(request_factory) + task.set_name(SECRETS[1]) + with pytest.raises(ADCPTaskError): + await handler.sync_reporting_receipts( + request, ToolContext(caller_identity=SECRETS[2]) + ) + finally: + logging.setLogRecordFactory(original_factory) + task.set_name(original_name) + record = assert_diagnostic(caplog, "handler", "RuntimeError", origin_function="resolver") + assert record.threadName is None and record.processName is None + assert getattr(record, "taskName", None) is None + assert not hasattr(record, "request_context") + + +async def test_origin_sanitizer_discards_paths_and_invalid_module_names(caplog): + caplog.set_level(logging.ERROR) + namespace = {"__name__": SECRETS[6], "RuntimeError": RuntimeError} + source = "async def resolver(*args):\n raise RuntimeError('private provider failure')\n" + exec(compile(source, "/private/provider/" + SECRETS[3] + ".py", "exec"), namespace) + async with receipt_harness("memory") as h: + _, request = await sensitive_case(h) + handler = ReportingReceiptHandler(h.store, resolve_account=namespace["resolver"]) + with pytest.raises(ADCPTaskError) as caught: + await handler.sync_reporting_receipts(request, ToolContext(caller_identity=SECRETS[2])) + assert caught.value.__cause__ is None and caught.value.__context__ is None + records = operator_errors(caplog) + assert len(records) == 1 + record = records[0] + assert (record.origin_module, record.origin_function, record.origin_line) == ( + "unknown", + "resolver", + 2, + ) + assert record.pathname == "" and record.exc_info is None and record.exc_text is None + serialized = json.dumps(record.__dict__) + assert "private provider failure" not in serialized + assert all(secret not in serialized for secret in SECRETS) + + +@pytest.mark.parametrize("boundary", ["resolver", "pg"]) +async def test_failing_operator_sink_cannot_replace_safe_translation(boundary, monkeypatch): + logger = logging.getLogger("adcp.reporting.receipts") + emitted = [] + + class BrokenSink(logging.Handler): + def emit(self, record): + emitted.append(record) + raise RuntimeError(SECRETS[0]) + + sink = BrokenSink() + async with receipt_harness("postgres" if boundary == "pg" else "memory") as h: + s, request = await sensitive_case(h) + mount = MountedReceipts(h) + mount.authorize(s) + + async def resolver(*args): + raise RuntimeError(*SECRETS) + + with monkeypatch.context() as patch: + if boundary == "resolver": + mount.handler._receipt_account_resolver = resolver + else: + fired = inject_driver_failure(patch, "execute", RuntimeError(*SECRETS)) + logger.addHandler(sink) + try: + with pytest.raises(ADCPTaskError) as caught: + await mount.handler.sync_reporting_receipts( + request, ToolContext(caller_identity=s.binding.consumer_id) + ) + finally: + logger.removeHandler(sink) + if boundary == "pg": + assert fired == ["execute"] + assert caught.value.errors[0].code == "RECEIPT_STORAGE_UNAVAILABLE" + assert caught.value.errors[0].message == SAFE_MESSAGE + assert caught.value.__context__ is None and caught.value.__cause__ is None + assert len(emitted) == 1 + assert emitted[0].exc_info is None and emitted[0].exc_text is None + assert all(secret not in json.dumps(emitted[0].__dict__) for secret in SECRETS) From a09878f67ab397a4b51b3314e3e8a5e87cf96da5 Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Fri, 18 Sep 2026 08:56:48 +0000 Subject: [PATCH 2/3] test(reporting): cover every raw receipt storage diagnostic boundary MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Independent review found that only `handler` and `store.ingest_receipt_batch` had committed boundary-label assertions, leaving `store.create_schema`, `store.receipt_ingestion_ready` and `store.read_receipt_boundaries` — all documented parts of the diagnostic contract — with no regression protection. A mislabelled decoration or a label missing from the runtime allowlist would silently reattribute an operator diagnostic with nothing failing. Exercise all four decorated raw boundaries on real PostgreSQL with genuine driver failures, asserting the injection fired, the exact safe code/message, empty cause/context, each boundary's own literal label, allowlist membership, exactly one record and a payload-free serialized LogRecord. Add an exhaustiveness guard so a new allowlisted store boundary cannot ship without its own executed regression, two mutation negative controls proving the label assertions are load-bearing, and a pin that a driver failure inside `validate_receipt_schema` stays the deliberately silent RECEIPT_SCHEMA_UNREADY path. Tests only; no production, documentation or workflow behaviour changes. The installed negative and preservation selections still resolve to exactly 30 and 48 cases, and the PostgreSQL-only additions skip correctly when the driver is absent. Co-Authored-By: Claude Opus 5 (1M context) --- .../test_reporting_receipt_diagnostics.py | 158 ++++++++++++++++++ 1 file changed, 158 insertions(+) diff --git a/tests/conformance/reporting/test_reporting_receipt_diagnostics.py b/tests/conformance/reporting/test_reporting_receipt_diagnostics.py index d6c20711b..486983cfd 100644 --- a/tests/conformance/reporting/test_reporting_receipt_diagnostics.py +++ b/tests/conformance/reporting/test_reporting_receipt_diagnostics.py @@ -450,3 +450,161 @@ async def resolver(*args): assert len(emitted) == 1 assert emitted[0].exc_info is None and emitted[0].exc_text is None assert all(secret not in json.dumps(emitted[0].__dict__) for secret in SECRETS) + + +STORE_BOUNDARY_CASES = ( + ("create_schema", "reporting_receipt_ingestion"), + ("ingest_receipt_batch", "INSERT INTO reporting_receipt_ingestion_results"), + ("read_receipt_boundaries", "reporting_receipt_ingestion_boundaries"), + # Readiness is decorated outside its validator only; a driver failure inside + # that validator is the deliberately silent RECEIPT_SCHEMA_UNREADY path. + ("receipt_ingestion_ready", None), +) +BOUNDARY_IDS = [case[0] for case in STORE_BOUNDARY_CASES] + + +def test_store_boundary_cases_cover_every_allowlisted_store_boundary(): + """A new decorated boundary cannot ship without its own executed regression.""" + from adcp.reporting.receipts._diagnostics import _BOUNDARIES + + covered = {f"store.{method}" for method, _ in STORE_BOUNDARY_CASES} + assert covered == _BOUNDARIES - {"handler"} + + +def _boundary_call(h, s, request, method): + return { + "create_schema": lambda: h.store.create_schema(), + "ingest_receipt_batch": lambda: h.store.ingest_receipt_batch( + request, caller=s.binding.principal + ), + "read_receipt_boundaries": lambda: h.store.read_receipt_boundaries( + caller=s.binding.principal + ), + "receipt_ingestion_ready": lambda: h.store.receipt_ingestion_ready(), + }[method] + + +@pytest.mark.parametrize("method,marker", STORE_BOUNDARY_CASES, ids=BOUNDARY_IDS) +async def test_each_raw_store_boundary_logs_its_own_label_once(caplog, monkeypatch, method, marker): + """Every decorated raw boundary owns its literal label and stays payload free.""" + caplog.set_level(logging.ERROR) + async with receipt_harness("postgres") as h: + from psycopg import AsyncConnection, OperationalError + from psycopg_pool import PoolTimeout + + from adcp.reporting.receipts._diagnostics import _BOUNDARIES + + s, request = await sensitive_case(h) + fired = [] + original = AsyncConnection.execute + + async def execute(connection, query, *args, **kwargs): + if not fired and isinstance(query, str) and marker in query: + fired.append(method) + raise OperationalError(SECRETS[7]) + return await original(connection, query, *args, **kwargs) + + def refuse(*args, **kwargs): + fired.append(method) + raise PoolTimeout(SECRETS[6]) + + caplog.clear() + with monkeypatch.context() as patch: + if marker is None: + patch.setattr(h.store._pool, "connection", refuse) + else: + patch.setattr(AsyncConnection, "execute", execute) + with pytest.raises(ReportingReceiptError) as caught: + await _boundary_call(h, s, request, method)() + assert fired == [method] + assert caught.value.code == "RECEIPT_STORAGE_UNAVAILABLE" + assert str(caught.value) == SAFE_MESSAGE + assert caught.value.__cause__ is None and caught.value.__context__ is None + record = assert_diagnostic( + caplog, + f"store.{method}", + "PoolTimeout" if marker is None else "OperationalError", + origin_function="refuse" if marker is None else "execute", + ) + assert record.boundary in _BOUNDARIES + + +@pytest.mark.parametrize("mutation", ["allowlist-omits-label", "mislabelled-decoration"]) +async def test_raw_store_boundary_label_assertion_is_load_bearing(caplog, monkeypatch, mutation): + """Negative control: a wrong label really is observable, so the label assertion bites. + + Production is correct, so the red side is produced by mutating the label + contract itself rather than by a setup or marker failure. The injection must + still fire on an otherwise valid call, and redaction must survive the mutation. + """ + caplog.set_level(logging.ERROR) + async with receipt_harness("postgres") as h: + from psycopg import AsyncConnection, OperationalError + + from adcp.reporting.receipts import _diagnostics + from adcp.reporting.receipts import pg as receipt_pg + + s, request = await sensitive_case(h) + fired = [] + original = AsyncConnection.execute + + async def execute(connection, query, *args, **kwargs): + if ( + not fired + and isinstance(query, str) + and "reporting_receipt_ingestion_boundaries" in query + ): + fired.append("read_receipt_boundaries") + raise OperationalError(SECRETS[7]) + return await original(connection, query, *args, **kwargs) + + caplog.clear() + with monkeypatch.context() as patch: + patch.setattr(AsyncConnection, "execute", execute) + if mutation == "allowlist-omits-label": + # The production fallback silently reattributes an unlisted label. + patch.setattr(_diagnostics, "_BOUNDARIES", frozenset({"handler"})) + expected = "handler" + else: + undecorated = type(h.store).read_receipt_boundaries.__wrapped__ + patch.setattr( + type(h.store), + "read_receipt_boundaries", + receipt_pg._storage_errors("store.create_schema")(undecorated), + ) + expected = "store.create_schema" + with pytest.raises(ReportingReceiptError) as caught: + await h.store.read_receipt_boundaries(caller=s.binding.principal) + assert fired == ["read_receipt_boundaries"] + assert caught.value.code == "RECEIPT_STORAGE_UNAVAILABLE" + assert caught.value.__cause__ is None and caught.value.__context__ is None + # The mutated label is emitted, so the positive assertion above would fail. + assert expected != "store.read_receipt_boundaries" + assert_diagnostic(caplog, expected, "OperationalError", origin_function="execute") + + +async def test_readiness_validator_failure_stays_silent_schema_unready(caplog, monkeypatch): + """The documented silent RECEIPT_SCHEMA_UNREADY path must not start logging.""" + caplog.set_level(logging.ERROR) + async with receipt_harness("postgres") as h: + from psycopg import AsyncConnection, OperationalError + + await sensitive_case(h) + fired = [] + original = AsyncConnection.execute + + async def execute(connection, query, *args, **kwargs): + if not fired and isinstance(query, str) and "pg_class" in query: + fired.append(True) + raise OperationalError(SECRETS[7]) + return await original(connection, query, *args, **kwargs) + + caplog.clear() + with monkeypatch.context() as patch: + patch.setattr(AsyncConnection, "execute", execute) + with pytest.raises(ReportingReceiptError) as caught: + await h.store.receipt_ingestion_ready() + assert fired == [True] + assert caught.value.code == "RECEIPT_SCHEMA_UNREADY" + assert operator_errors(caplog) == [] + assert all(secret not in caplog.text for secret in SECRETS) From ab12e1511058f0a9e75a8228f650d9a302423640 Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Thu, 24 Sep 2026 22:36:19 +0000 Subject: [PATCH 3/3] fix(reporting): qualify hardening against integrated feed baseline --- .github/workflows/ci.yml | 7 ++++--- docs/reporting-receipt-ingress.md | 7 +++++++ .../conformance/reporting/_hardening_packaging.py | 2 +- .../test_reporting_activity_schema_proof.py | 15 ++++++++++----- .../test_reporting_feed_hardening_installed.py | 11 +++++++---- .../test_reporting_receipt_diagnostics.py | 7 ++++--- 6 files changed, 33 insertions(+), 16 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index aa619fd06..7b3581c4f 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -524,7 +524,7 @@ jobs: feed_tests=() for test_file in tests/conformance/reporting/test_reporting_feed_*.py; do case "$test_file" in - tests/conformance/reporting/test_reporting_feed_rolling.py|tests/conformance/reporting/test_reporting_feed_packaging.py|tests/conformance/reporting/test_reporting_feed_installed_pg.py) ;; + tests/conformance/reporting/test_reporting_feed_rolling.py|tests/conformance/reporting/test_reporting_feed_packaging.py|tests/conformance/reporting/test_reporting_feed_installed_pg.py|tests/conformance/reporting/test_reporting_feed_hardening_installed.py) ;; *) feed_tests+=("$test_file") ;; esac done @@ -624,8 +624,9 @@ jobs: --health-retries 10 steps: - uses: actions/checkout@v6 - - name: Fetch the exact approved B2.3 hardening comparison artifact - run: git fetch --no-tags origin 50e35f0ae3540f19b40e8fc460f5870dfe018bf9 + - name: Fetch the exact integrated B2.3 hardening comparison artifact + timeout-minutes: 1 + run: git fetch --no-tags origin 2d777ace7b4bf8be519ce0abd4fd0a25ed4f1da7 - uses: actions/setup-python@v6 id: feed-python310 with: diff --git a/docs/reporting-receipt-ingress.md b/docs/reporting-receipt-ingress.md index 4b70af65c..ce74c3966 100644 --- a/docs/reporting-receipt-ingress.md +++ b/docs/reporting-receipt-ingress.md @@ -248,7 +248,14 @@ integrated artifacts remove those dependencies; a C-only cluster would hide portability regressions, so the gates require URL and drivers without a locale pin. The pre-`17ee407a` A, pre-`0f34c666` B, pre-`967b6e28` C, pre-`5487f2bd` B1 and pre-`3fd62121` B2.1 rolling exclusions must accompany release notes. + A's notification-readiness closure after C is compared on both sides and does not excuse new regressions. Optional notifications may stay explicitly disabled; complete polling still depends on the later B2.3/B2.4 components. Full buyer adjustment automation and `client.reporting` remain named downstream #1172 work. + +The schema-proof and receipt-diagnostic hardening comparison uses integrated +B2.3 commit `2d777ace`. It checks the unchanged safe buyer error and frozen-feed +restart boundary against that binary, while separately demonstrating its repeated +catalog work and absent operator diagnostics. It does not qualify earlier B2.3 +snapshots; this pre-`2d777ace` comparison limit must also accompany release notes. diff --git a/tests/conformance/reporting/_hardening_packaging.py b/tests/conformance/reporting/_hardening_packaging.py index 50d2752ff..80855d237 100644 --- a/tests/conformance/reporting/_hardening_packaging.py +++ b/tests/conformance/reporting/_hardening_packaging.py @@ -28,7 +28,7 @@ "ledger/reporting_receipt_ingestion.sql", "receipts/required_schema.json", ) -B23 = "50e35f0ae3540f19b40e8fc460f5870dfe018bf9" +B23 = "2d777ace7b4bf8be519ce0abd4fd0a25ed4f1da7" def installed_hardening(root, python, wheel, *, label, parent=False, driver_absent=False): diff --git a/tests/conformance/reporting/test_reporting_activity_schema_proof.py b/tests/conformance/reporting/test_reporting_activity_schema_proof.py index b57217f03..16f5a29d7 100644 --- a/tests/conformance/reporting/test_reporting_activity_schema_proof.py +++ b/tests/conformance/reporting/test_reporting_activity_schema_proof.py @@ -339,8 +339,10 @@ async def test_synchronous_startup_then_runtime_loop_has_no_loop_bound_cache(act async def test_memory_only_claims_and_frozen_constructor_remain_unchanged(): async with reliable_factory("memory", notifications=True) as reliable: support = await compose_activity(reliable, False) - assert not await support.durable() - assert not await support.durable() + hardening_operation_1 = await support.durable() + assert not hardening_operation_1 + hardening_operation_2 = await support.durable() + assert not hardening_operation_2 with mounted_activity(support) as handler: handler._platform.claim = False await handler.get_adcp_capabilities() @@ -486,8 +488,10 @@ def c_contract(installed, **options): monkeypatch.setattr(_schema, "_validate_schema_objects", b_contract) monkeypatch.setattr(status_schema, "_validate_status_objects", c_contract) for support in (b_only, composite): - assert await support.durable() - assert await support.durable() + hardening_operation_4 = await support.durable() + assert hardening_operation_4 + hardening_operation_5 = await support.durable() + assert hardening_operation_5 assert calls == [ ("B", {"activity": True}), ("B", {"activity": True}), @@ -500,7 +504,8 @@ def c_contract(installed, **options): "ALTER TABLE reporting_status_webhook_attempts DISABLE TRIGGER" " reporting_status_webhook_attempt_guard" ) - assert await b_only.durable() + hardening_operation_3 = await b_only.durable() + assert hardening_operation_3 # A new server instance is the normal migration/restart boundary. restarted = replace(composite) with pytest.raises(ReportingNotificationError, match="status_schema_unready"): diff --git a/tests/conformance/reporting/test_reporting_feed_hardening_installed.py b/tests/conformance/reporting/test_reporting_feed_hardening_installed.py index 6010c15cf..6aa91893d 100644 --- a/tests/conformance/reporting/test_reporting_feed_hardening_installed.py +++ b/tests/conformance/reporting/test_reporting_feed_hardening_installed.py @@ -84,7 +84,7 @@ def approved_b23(tmp_path_factory, request): "modules": modules, "assets": assets, "python": [3, 10], - "tree": "16c55b24340aeea3524c781f19bfa64d1491ac34", + "tree": "2f71a273c0218e7ffc490fb4df243d02711cff7b", } return root, python, script, helper, installed, wheel @@ -117,7 +117,8 @@ async def test_actual_b23_to_child_installed_restart_preserves_pages_and_receipt # replay before it writes page one; fixtures only supply populated data. async with feed_process(h, s, receipt_request, action="receipt", **old) as child: admitted = await child.event("done") - assert await asyncio.wait_for(child.process.wait(), 5) == 0 + hardening_operation_1 = await asyncio.wait_for(child.process.wait(), 5) + assert hardening_operation_1 == 0 assert admitted["result"] == receipt_response async with feed_process(h, s, feed_request(s), pause="committed", **old) as child: first = (await child.event("committed"))["result"] @@ -158,14 +159,16 @@ async def test_actual_b23_to_child_installed_restart_preserves_pages_and_receipt h, s, continuation, action="walk", transport="a2a", v1=v1, **new ) as child: continued = await child.event("done") - assert await asyncio.wait_for(child.process.wait(), 5) == 0 + hardening_operation_3 = await asyncio.wait_for(child.process.wait(), 5) + assert hardening_operation_3 == 0 assert continued["result"]["pages"] == expected[0][1:] assert continued["result"]["binding"] == original.binding assert continued["result"]["version"] == original.representation_version assert continued["result"]["ownership_mode"] == original.ownership_mode == "absent" async with feed_process(h, s, receipt_request, action="receipt", **new) as child: replayed = await child.event("done") - assert await asyncio.wait_for(child.process.wait(), 5) == 0 + hardening_operation_2 = await asyncio.wait_for(child.process.wait(), 5) + assert hardening_operation_2 == 0 assert replayed["result"] == receipt_response assert without_feed(await h.image()) == before assert ( diff --git a/tests/conformance/reporting/test_reporting_receipt_diagnostics.py b/tests/conformance/reporting/test_reporting_receipt_diagnostics.py index 486983cfd..392f661fd 100644 --- a/tests/conformance/reporting/test_reporting_receipt_diagnostics.py +++ b/tests/conformance/reporting/test_reporting_receipt_diagnostics.py @@ -204,9 +204,10 @@ async def test_original_pg_execute_and_commit_failures_log_once_across_translati caplog, "store.ingest_receipt_batch", "OperationalError", origin_function=point ) assert await h.image() == before - assert (await h.store.ingest_receipt_batch(request, caller=s.binding.principal))["results"][ - 0 - ]["result"] == "recorded" + hardening_operation_1 = await h.store.ingest_receipt_batch( + request, caller=s.binding.principal + ) + assert (hardening_operation_1)["results"][0]["result"] == "recorded" @pytest.mark.parametrize(