diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 86415d58a..7b3581c4f 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. @@ -523,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 @@ -605,7 +606,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 permissions: contents: read services: @@ -623,6 +624,9 @@ jobs: --health-retries 10 steps: - uses: actions/checkout@v6 + - 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: @@ -636,21 +640,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 578e48690..ce74c3966 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.6"`. This is the wire release spelling; @@ -205,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/docs/reporting-webhook-activity.md b/docs/reporting-webhook-activity.md index 41953996c..22aea2e34 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 09791d331..050017734 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 @@ -241,9 +242,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..80855d237 --- /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 = "2d777ace7b4bf8be519ce0abd4fd0a25ed4f1da7" + + +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 963d4186c..52a2bafe4 100644 --- a/tests/conformance/reporting/test_reporting_activity_migration.py +++ b/tests/conformance/reporting/test_reporting_activity_migration.py @@ -254,6 +254,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..16f5a29d7 --- /dev/null +++ b/tests/conformance/reporting/test_reporting_activity_schema_proof.py @@ -0,0 +1,527 @@ +"""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) + 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() + 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): + 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}), + ("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" + ) + 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"): + 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..6aa91893d --- /dev/null +++ b/tests/conformance/reporting/test_reporting_feed_hardening_installed.py @@ -0,0 +1,194 @@ +"""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": "2f71a273c0218e7ffc490fb4df243d02711cff7b", + } + 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") + 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"] + 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") + 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") + 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 ( + 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..392f661fd --- /dev/null +++ b/tests/conformance/reporting/test_reporting_receipt_diagnostics.py @@ -0,0 +1,611 @@ +"""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 + 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( + "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) + + +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)