From 5e1cb8c014fa2bdacf1e59ae77e21d2148c25f89 Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Mon, 28 Sep 2026 08:36:58 +0000 Subject: [PATCH 1/3] test(reporting): require full installed lifecycle interop matrix --- .github/workflows/ci.yml | 27 +- .../ci/reporting_interop/full_scenario.json | 81 ++++ .../reporting_interop/python_full_server.py | 283 ++++++++++++++ .../run_full_installed_artifact_matrix.py | 355 ++++++++++++++++++ .../ci/reporting_interop/ts_full_receipt.cjs | 134 +++++++ ...eporting_full_installed_artifact_matrix.py | 99 +++++ 6 files changed, 976 insertions(+), 3 deletions(-) create mode 100644 scripts/ci/reporting_interop/full_scenario.json create mode 100644 scripts/ci/reporting_interop/python_full_server.py create mode 100644 scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py create mode 100644 scripts/ci/reporting_interop/ts_full_receipt.cjs create mode 100644 tests/test_reporting_full_installed_artifact_matrix.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 053b556a2..d9640c1e8 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -1348,9 +1348,9 @@ jobs: python -m venv "${INPUT_ROOT}/python/stable" python -m venv "${INPUT_ROOT}/python/candidate" "${INPUT_ROOT}/python/stable/bin/pip" install --retries 5 --timeout 30 \ - "${INPUT_ROOT}/python/adcp-8.0.0b16-py3-none-any.whl[pg]" + "${INPUT_ROOT}/python/adcp-8.0.0b16-py3-none-any.whl[pg]" pytest==9.0.2 "${INPUT_ROOT}/python/candidate/bin/pip" install --retries 5 --timeout 30 \ - "${INPUT_ROOT}/python/adcp-8.0.0b18-py3-none-any.whl[pg]" + "${INPUT_ROOT}/python/adcp-8.0.0b18-py3-none-any.whl[pg]" pytest==9.0.2 cp scripts/ci/reporting_interop/npm/rc47/package.json \ scripts/ci/reporting_interop/npm/rc47/package-lock.json "${INPUT_ROOT}/ts-stable/" @@ -1390,11 +1390,32 @@ jobs: --pg-admin-url postgresql://postgres@127.0.0.1:5432/postgres \ --output "${RUNNER_TEMP}/reporting-installed-artifact-matrix" + - name: Run installed full reporting lifecycle 2x2 + shell: bash + run: | + set -euo pipefail + INPUT_ROOT="${RUNNER_TEMP}/reporting-1199-inputs" + "${INPUT_ROOT}/python/candidate/bin/python" -I \ + scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py \ + --python-stable-runtime "${INPUT_ROOT}/python/stable/bin/python" \ + --python-stable-artifact "${INPUT_ROOT}/python/adcp-8.0.0b16-py3-none-any.whl" \ + --python-candidate-runtime "${INPUT_ROOT}/python/candidate/bin/python" \ + --python-candidate-artifact "${INPUT_ROOT}/python/adcp-8.0.0b18-py3-none-any.whl" \ + --typescript-stable-install "${INPUT_ROOT}/ts-stable" \ + --typescript-stable-archive "${INPUT_ROOT}/ts-stable/sdk-14.0.0-rc.47.tgz" \ + --typescript-candidate-install "${INPUT_ROOT}/ts-candidate" \ + --typescript-candidate-archive "${INPUT_ROOT}/ts-candidate/sdk-14.0.0-rc.48.tgz" \ + --node-runtime "$(command -v node)" \ + --pg-admin-url postgresql://postgres@127.0.0.1:5432/postgres \ + --output "${RUNNER_TEMP}/reporting-full-installed-artifact-matrix" + - if: always() uses: actions/upload-artifact@v7 with: name: reporting-installed-artifact-matrix-${{ github.run_attempt }} - path: ${{ runner.temp }}/reporting-installed-artifact-matrix + path: | + ${{ runner.temp }}/reporting-installed-artifact-matrix + ${{ runner.temp }}/reporting-full-installed-artifact-matrix if-no-files-found: warn reporting-installed-artifact-required-gate: diff --git a/scripts/ci/reporting_interop/full_scenario.json b/scripts/ci/reporting_interop/full_scenario.json new file mode 100644 index 000000000..3749e2331 --- /dev/null +++ b/scripts/ci/reporting_interop/full_scenario.json @@ -0,0 +1,81 @@ +{ + "expected_outcome": { + "definitive": true, + "inspection_calls": 0, + "obligation_count": 1, + "reasons_absent": [ + "consumer_status_pending", + "receipt_required" + ] + }, + "expected_periods": [ + { + "automatedRecoveryWindowSeconds": 21600, + "canonicalization": { + "id": "reference-jcs-rows-v1", + "primaryKeys": [ + "row_id" + ], + "sha256": "5566f62a29487d54d755e4e9c35c97661bb243d158d4bf7a712fb1c7eb3bcb80", + "uri": "https://contracts.example.test/reference-canonicalization.json" + }, + "coverage": { + "covered_package_ids": [], + "fully_covered_media_buy_ids": [ + "mb_acct_a" + ], + "media_buy_ids": [ + "mb_acct_a" + ], + "package_ids": [], + "partially_covered_media_buy_ids": [], + "status": "full", + "unknown_media_buy_ids": [], + "unknown_package_ids": [], + "unsupported_media_buy_ids": [], + "unsupported_package_ids": [] + }, + "coverageRequirement": "full", + "deliveryConfigId": "daily", + "deliveryConfigVersion": 1, + "deliveryMethod": "warehouse_materialization", + "deliverySlaSeconds": 3600, + "destinationRef": "destination", + "feedPurpose": "billing", + "mediaBuyIds": [ + "mb_acct_a" + ], + "officialFinality": { + "basis": "source_final", + "policyId": "fixture-official-v1" + }, + "periodEnd": "2026-09-01T01:00:00Z", + "periodSourceTimezone": "UTC", + "periodStart": "2026-09-01T00:00:00Z", + "reconciliationMode": "consumer_receipt", + "reportDefinitionId": "reference-report-v1", + "reportDefinitionSha256": "2c988a2c008b84097be962248c433dc3acd9663372f89a4b1c7c1235ee635785", + "reportDefinitionUri": "https://contracts.example.test/reference-definition.json", + "reportingProfile": "paid_media_delivery", + "requiredFinality": "official", + "schemaDialect": "https://json-schema.org/draft/2020-12/schema", + "schemaRefPolicy": "local_fragment_only", + "schemaSha256": "aa4927d97b5e5889c4d79bd5461da81657cafea75a61d266ab452bfa0fc5d799", + "schemaUri": "https://contracts.example.test/reference-row-schema.json", + "schemaVersion": "1.0.0", + "verificationProfile": "canonical_digest" + } + ], + "now": "2026-09-28T07:31:39.502842Z", + "request": { + "account": { + "account_id": "acct_a" + }, + "period": { + "end": "2026-09-01T01:00:00Z", + "start": "2026-09-01T00:00:00Z" + }, + "view": "periods" + }, + "schema_version": 1 +} diff --git a/scripts/ci/reporting_interop/python_full_server.py b/scripts/ci/reporting_interop/python_full_server.py new file mode 100644 index 000000000..34339c180 --- /dev/null +++ b/scripts/ci/reporting_interop/python_full_server.py @@ -0,0 +1,283 @@ +#!/usr/bin/env python3 +"""Installed-wheel PostgreSQL seller for the full #1199 reporting lane. + +Only the fixture driver comes from this checkout. The SDK is imported from the +selected installed wheel, and the runner gives each invocation a fresh database. +""" + +from __future__ import annotations + +import argparse +import json +import sys +from contextlib import AsyncExitStack +from dataclasses import replace +from pathlib import Path +from types import SimpleNamespace +from typing import Any + +import pytest +from psycopg_pool import AsyncConnectionPool + +import adcp +from adcp.server import ADCPHandler, serve +from adcp.server.auth import BearerTokenAuth, Principal, auth_context_factory + +ROOT = Path(__file__).resolve().parents[3] +SDK_SOURCE = ROOT / "src" / "adcp" +if Path(adcp.__file__).resolve().is_relative_to(SDK_SOURCE): + raise RuntimeError("full matrix seller must use the installed Python wheel") + +sys.path.insert(0, str(ROOT)) +from tests.conformance.reporting._generation_support import ( # noqa: E402 + configuration, + obligation_for, + revision_for, +) +from tests.conformance.reporting._production_support import production_harness # noqa: E402 +from tests.conformance.reporting._projection_support import drain # noqa: E402 +from tests.conformance.reporting._reliable_support import ( # noqa: E402 + DeterministicReceiverStore, + ScriptedNotificationReceiver, + SimulatedCrash, + _BytesStore, + notification_subscription, + notification_verification_keys, +) + +EVENTS = ( + "reporting.ledger_changed", + "reporting.status_changed", + "reporting.delivery_ready", +) + + +def _arguments() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--port", type=int, required=True) + parser.add_argument("--database-url", required=True) + parser.add_argument("--destination", type=Path, required=True) + parser.add_argument("--evidence", type=Path, required=True) + return parser.parse_args() + + +async def _verify_notification_replay(harness: Any) -> dict[str, Any]: + """Exercise signed replay and account activity in this cell's installed SDK.""" + from adcp.reporting.production.notifications import production_notification_workers + from adcp.signing.jwks import StaticJwksResolver + from adcp.signing.webhook_verifier import WebhookVerifyOptions, verify_webhook_signature + + h = harness + h.subscriptions.put( + notification_subscription( + subscriber="buyer", + principal=h.item.binding.consumer_id, + events=EVENTS, + url="https://receiver.example.test/reporting", + ) + ) + patcher = pytest.MonkeyPatch() + blobs = _BytesStore(h.pool) + await blobs.create_schema() + receiver_store = DeterministicReceiverStore(blobs, h.notification_failures) + receiver = ScriptedNotificationReceiver( + SimpleNamespace(clock=h.clock, failures=h.notification_failures, receiver=receiver_store) + ) + receiver.install(patcher) + try: + # The original revision was committed before the subscriber existed. + # A distinct ordinary revision supplies a real post-registration event. + core = replace(configuration(), delivery_config_id="interop-webhook") + await h.store.put_configuration(core) + obligation = await h.store.commit_obligation( + replace(obligation_for(core), reporting_obligation_id="interop-webhook-obligation") + ) + revision, rows = revision_for(obligation, suffix="interop-webhook") + await h.store.commit_revision(revision, rows) + worker = h.production.notification_workers[0] + for _ in range(100): + if not await worker.expand_one(account_id="acct_a"): + break + else: + raise RuntimeError("webhook expansion did not reach idle") + if not await worker.outbox.list_deliveries(account_id="acct_a"): + raise RuntimeError("webhook delivery was not queued") + h.notification_failures.at("http.accepted", SimulatedCrash()) + try: + await worker.deliver_one(account_id="acct_a") + except SimulatedCrash: + pass + else: + raise RuntimeError("webhook crash boundary did not fire") + first = receiver.received[-1] + if await receiver_store.read("acct_a", first.idempotency_key) != first.body: + raise RuntimeError("receiver did not retain accepted webhook bytes") + async with worker.outbox._connection() as connection: + await connection.execute( + "UPDATE reporting_notification_deliveries" + " SET lease_expires_at=clock_timestamp()-interval '1 second'" + " WHERE account_id=%s AND state='leased'", + ("acct_a",), + ) + h.signing.generation = 2 + restarted = production_notification_workers( + h.store, + h.projection, + subscriptions=h.subscriptions, + cipher=worker.cipher, + signing=worker.signing, + )[0] + for _ in range(100): + if not await restarted.deliver_one(account_id="acct_a"): + break + else: + raise RuntimeError("webhook retry did not reach idle") + replay = [r for r in receiver.received if r.idempotency_key == first.idempotency_key] + if len(replay) != 2 or replay[0].body != replay[1].body: + raise RuntimeError("webhook retry changed the signed body or idempotency key") + for generation, attempt in enumerate(replay, start=1): + signature_input = next( + ( + value + for key, value in attempt.headers.items() + if key.lower() == "signature-input" + ), + "", + ) + if f"#key-{generation}" not in signature_input: + raise RuntimeError("webhook retry did not use the expected signing generation") + options = WebhookVerifyOptions( + jwks_resolver=StaticJwksResolver({"keys": notification_verification_keys()}), + clock=lambda: h.clock().timestamp(), + ) + for attempt in replay: + verify_webhook_signature( + method="POST", + url="https://receiver.example.test" + attempt.target, + headers=attempt.headers, + body=attempt.body, + options=options, + ) + activity = await restarted.outbox.list_activity( + account_id="acct_a", consumer_id=h.item.binding.consumer_id, limit=20 + ) + if not activity or not any( + record.outcome and record.outcome.status == "success" for record in activity + ): + raise RuntimeError("account activity omitted the successful webhook replay") + return { + "attempt_count": len(replay), + "body_unchanged": True, + "idempotency_key_unchanged": True, + "signature_generations": [1, 2], + "verified_signature_count": len(replay), + } + finally: + patcher.undo() + + +def main() -> None: + args = _arguments() + holder: dict[str, Any] = {} + stack = AsyncExitStack() + pool = AsyncConnectionPool( + args.database_url, + min_size=1, + max_size=4, + open=False, + kwargs={"autocommit": True}, + ) + + async def startup() -> None: + await pool.open() + await pool.wait(timeout=15) + harness = await stack.enter_async_context( + production_harness( + "postgres", + args.destination, + existing_pool=pool, + reconciled=True, + notifications=True, + notification_delivery=True, + feedback=True, + count=2, + ) + ) + await harness.production.activate(account_id=harness.item.config.account_id) + result = await harness.production.materializer.run_once() + if result.state not in ("verified", "idle"): + raise RuntimeError(f"materialization did not verify: {result}") + await drain(harness.projection, harness.item.config.account_id) + holder["harness"] = harness + + async def shutdown() -> None: + try: + harness = holder.get("harness") + if harness is not None: + await harness.production.aclose() + result = await _verify_notification_replay(harness) + args.evidence.write_text( + json.dumps(result, sort_keys=True) + "\n", encoding="utf-8" + ) + finally: + await stack.aclose() + await pool.close() + + class Handler(ADCPHandler): + advertised_tools = { + "get_adcp_capabilities", + "get_reporting_status", + "sync_reporting_receipts", + "sync_reporting_status", + "get_media_buy_delivery", + "sync_accounts", + } + + async def get_adcp_capabilities(self, params: Any, context: Any = None) -> Any: + return await holder["harness"].production.handler.get_adcp_capabilities(params, context) + + async def get_reporting_status(self, params: Any, context: Any = None) -> Any: + return await holder["harness"].production.handler.get_reporting_status(params, context) + + async def sync_reporting_receipts(self, params: Any, context: Any = None) -> Any: + return await holder["harness"].production.handler.sync_reporting_receipts( + params, context + ) + + async def sync_reporting_status(self, params: Any, context: Any = None) -> Any: + return await holder["harness"].production.handler.sync_reporting_status(params, context) + + async def get_media_buy_delivery(self, params: Any, context: Any = None) -> Any: + return await holder["harness"].production.handler.get_media_buy_delivery( + params, context + ) + + async def sync_accounts(self, params: Any, context: Any = None) -> Any: + return await holder["harness"].production.handler.sync_accounts(params, context) + + def validate_token(token: str) -> Principal | None: + if token != "full-matrix-token": + return None + return Principal( + caller_identity="https://buyer.example.test/agent", + tenant_id="matrix-tenant", + metadata={"account_id": "acct_a"}, + ) + + serve( + Handler(), + name="reporting-full-installed-seller", + host="127.0.0.1", + port=args.port, + transport="both", + auth=BearerTokenAuth(validate_token=validate_token), + context_factory=auth_context_factory, + stateless_http=True, + allowed_hosts=("127.0.0.1", "localhost"), + on_startup=(startup,), + on_shutdown=(shutdown,), + ) + + +if __name__ == "__main__": + main() diff --git a/scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py b/scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py new file mode 100644 index 000000000..41cbff4f5 --- /dev/null +++ b/scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py @@ -0,0 +1,355 @@ +#!/usr/bin/env python3 +"""Run #1199's full installed Python/TypeScript reporting lifecycle matrix. + +Each cell owns a fresh PostgreSQL database, installed Python seller process, +installed TypeScript client, and retained output. An error in any cell fails the +single aggregate result; historical and mutable artifacts receive no credit. +""" + +from __future__ import annotations + +import importlib.util +import json +import secrets +import subprocess +import sys +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +HERE = Path(__file__).resolve().parent +ROOT = HERE.parents[2] +SELLER = HERE / "python_full_server.py" +RECEIPT = HERE / "ts_full_receipt.cjs" +BUYER = HERE / "ts_mcp_buyer.cjs" +SCENARIO = HERE / "full_scenario.json" +PINS = HERE / "pins.json" +TOKEN = "full-matrix-token" + + +def _load_core() -> Any: + spec = importlib.util.spec_from_file_location( + "reporting_installed_core_matrix", HERE / "run_installed_artifact_matrix.py" + ) + if spec is None or spec.loader is None: + raise RuntimeError("cannot load installed Core matrix") + module = importlib.util.module_from_spec(spec) + sys.modules[spec.name] = module + spec.loader.exec_module(module) + return module + + +core = _load_core() + + +def _client_command( + script: Path, + *, + node: Path, + install: Any, + version: str, + endpoint: str, +) -> list[str]: + pin = core.foundation._typescript_pin(install) + return [ + str(node), + str(script), + "--package-lock", + str(install.lock), + "--expected-version", + pin["version"], + "--expected-integrity", + pin["integrity"], + "--adcp-version", + version, + "--endpoint", + endpoint, + "--auth-env", + "ADCP_INTEROP_BUYER_TOKEN", + ] + + +def _full_result( + *, + receipt: dict[str, Any] | None, + buyer: dict[str, Any] | None, + webhook: dict[str, Any] | None, +) -> bool: + if not ( + receipt + and receipt.get("status") == "passed" + and receipt.get("exact_revision_read") is True + and receipt.get("accepted_receipt_count") == 1 + and receipt.get("reconciliation_status") == "accepted" + and buyer + and buyer.get("status") == "passed" + and buyer.get("reconciliation", {}).get("definitive") is True + and buyer.get("inspection_calls") == 0 + and webhook + and webhook.get("attempt_count", 0) >= 2 + and webhook.get("body_unchanged") is True + and webhook.get("idempotency_key_unchanged") is True + and webhook.get("verified_signature_count") == 2 + and webhook.get("signature_generations") == [1, 2] + ): + return False + obligations = buyer["reconciliation"].get("obligations", []) + return ( + len(obligations) == 1 + and obligations[0].get("reportingObligationId") == "rpo_acct_a" + and receipt["reporting_revision_id"] == "production-revision" + ) + + +def _blocking_acceptance(rows: list[dict[str, Any]]) -> bool: + return ( + len(rows) == 4 + and not core._cross_cell_isolation_errors(rows) + and all(row["positive_full_lifecycle"] for row in rows) + ) + + +def _run_cell( + cell: dict[str, Any], + *, + contract: dict[str, Any], + python_input: tuple[Any, Any], + typescript_input: tuple[Any, Any], + node: Path, + admin_url: str, + database: str, + output: Path, +) -> dict[str, Any]: + foundation = core.foundation + runtime, _artifact = python_input + install, _archive = typescript_input + database_identity = foundation._create_database(admin_url, database) + database_url = foundation._node_database_url(admin_url, database) + version = core._seller_wire_version(contract, cell["python_role"]) + port = foundation._free_port() + endpoint = f"http://127.0.0.1:{port}/mcp" + scenario = json.loads(SCENARIO.read_text(encoding="utf-8")) + scenario["now"] = datetime.now(timezone.utc).isoformat() + scenario_path = output / "scenario.json" + scenario_path.write_text(json.dumps(scenario, sort_keys=True) + "\n", encoding="utf-8") + seller_out = output / "seller.stdout.log" + seller_err = output / "seller.stderr.log" + receipt_out = output / "receipt.stdout.json" + receipt_err = output / "receipt.stderr.log" + buyer_out = output / "buyer.stdout.json" + buyer_err = output / "buyer.stderr.log" + webhook_path = output / "webhook.json" + receipt: dict[str, Any] | None = None + buyer: dict[str, Any] | None = None + webhook: dict[str, Any] | None = None + issues: list[str] = [] + process_identity: dict[str, Any] | None = None + with seller_out.open("wb") as stdout, seller_err.open("wb") as stderr: + process = subprocess.Popen( + [ + str(runtime.executable), + "-I", + str(SELLER), + "--port", + str(port), + "--database-url", + database_url, + "--destination", + str(output / "destination.sqlite"), + "--evidence", + str(webhook_path), + ], + cwd=output, + stdin=subprocess.DEVNULL, + stdout=stdout, + stderr=stderr, + start_new_session=True, + ) + try: + process_identity = { + "pid": process.pid, + "start_token": foundation._linux_process_start_token(process.pid), + } + foundation._wait_port(process, port) + receipt_exit, receipt = core._client_result( + _client_command( + RECEIPT, node=node, install=install, version=version, endpoint=endpoint + ), + cwd=output, + token=TOKEN, + stdout_path=receipt_out, + stderr_path=receipt_err, + ) + if receipt_exit != 0 or receipt is None or receipt.get("status") != "passed": + issues.append("receipt_or_exact_read_failed") + else: + buyer_exit, buyer = core._client_result( + [ + *_client_command( + BUYER, node=node, install=install, version=version, endpoint=endpoint + ), + "--input", + str(scenario_path), + ], + cwd=output, + token=TOKEN, + stdout_path=buyer_out, + stderr_path=buyer_err, + ) + if buyer_exit != 0 or buyer is None or buyer.get("status") != "passed": + issues.append("semantic_reconciliation_failed") + except ( + core.MatrixError, + foundation.HarnessError, + OSError, + subprocess.TimeoutExpired, + ) as error: + issues.append(f"cell_process_failed:{type(error).__name__}:{error}") + finally: + foundation._stop(process) + if webhook_path.is_file(): + try: + webhook = json.loads(webhook_path.read_text(encoding="utf-8")) + except json.JSONDecodeError: + issues.append("webhook_evidence_invalid") + else: + issues.append("webhook_replay_not_completed") + positive = not issues and _full_result(receipt=receipt, buyer=buyer, webhook=webhook) + if not positive and not issues: + issues.append("full_lifecycle_assertion_failed") + retained = [scenario_path, seller_out, seller_err] + retained.extend( + path + for path in (receipt_out, receipt_err, buyer_out, buyer_err, webhook_path) + if path.is_file() + ) + return { + "id": cell["id"], + "python_role": cell["python_role"], + "typescript_role": cell["typescript_role"], + "database_identity": database_identity, + "seller_process_identity": process_identity, + "output_directory": output.name, + "positive_full_lifecycle": positive, + "issues": issues, + "receipt": receipt, + "buyer": buyer, + "webhook": webhook, + "evidence": [core._retained(path, output) for path in retained], + } + + +def main() -> None: + core.foundation = core._load_foundation() + foundation = core.foundation + args = core._arguments() + contract = core._load_contract(args.contract) + protocol_source = json.loads(PINS.read_text(encoding="utf-8"))["protocol"]["required_contract"] + if protocol_source["version"] != contract["artifacts"]["python"]["candidate"]["protocol"]: + raise core.MatrixError("installed candidate protocol differs from the signed source pin") + args.output = args.output.absolute() + args.output.mkdir(parents=True, exist_ok=False) + python_inputs = { + "stable": ( + foundation.PythonRuntime("previous", args.python_stable_runtime), + foundation.PythonArtifact("previous", args.python_stable_artifact), + ), + "candidate": ( + foundation.PythonRuntime("wheel", args.python_candidate_runtime), + foundation.PythonArtifact("wheel", args.python_candidate_artifact), + ), + } + typescript_inputs = { + "stable": ( + foundation.TypeScriptInstall("previous", args.typescript_stable_install), + foundation.TypeScriptArchive("previous", args.typescript_stable_archive), + ), + "candidate": ( + foundation.TypeScriptInstall("candidate", args.typescript_candidate_install), + foundation.TypeScriptArchive("candidate", args.typescript_candidate_archive), + ), + } + identities = core._validate_artifact_contract( + contract, + python_inputs=python_inputs, + typescript_inputs=typescript_inputs, + node=args.node_runtime, + ) + run_id = secrets.token_hex(6) + created: list[str] = [] + rows: list[dict[str, Any]] = [] + try: + for index, cell in enumerate(contract["cells"]): + output = args.output / cell["id"] + output.mkdir() + database = f"adcp_1199_full_{run_id}_{index}" + created.append(database) + rows.append( + _run_cell( + cell, + contract=contract, + python_input=python_inputs[cell["python_role"]], + typescript_input=typescript_inputs[cell["typescript_role"]], + node=args.node_runtime, + admin_url=args.pg_admin_url, + database=database, + output=output, + ) + ) + finally: + if not args.keep_databases: + foundation._drop_databases(args.pg_admin_url, created) + isolation_errors = core._cross_cell_isolation_errors(rows) + accepted = _blocking_acceptance(rows) + report = { + "schema_version": 1, + "issue": 1199, + "run_id": run_id, + "status": "passed" if accepted else "failed", + "blocking_acceptance": accepted, + "declared_artifacts": contract["artifacts"], + "protocol_source": protocol_source, + "artifact_identities": identities, + "cross_cell_isolation_errors": isolation_errors, + "cells": rows, + "harness": { + **foundation._git_identity(), + "entrypoints": { + name: core._retained(path, ROOT) + for name, path in { + "seller": SELLER, + "receipt": RECEIPT, + "buyer": BUYER, + "scenario": SCENARIO, + }.items() + }, + "contract": core._retained(args.contract, ROOT), + "protocol_pins": core._retained(PINS, ROOT), + }, + } + (args.output / "results.json").write_text( + json.dumps(report, indent=2, sort_keys=True) + "\n", encoding="utf-8" + ) + print( + json.dumps( + { + "status": report["status"], + "run_id": run_id, + "cells": [ + { + "id": row["id"], + "positive_full_lifecycle": row["positive_full_lifecycle"], + "issues": row["issues"], + } + for row in rows + ], + }, + sort_keys=True, + ) + ) + if not accepted: + raise SystemExit(1) + + +if __name__ == "__main__": + main() diff --git a/scripts/ci/reporting_interop/ts_full_receipt.cjs b/scripts/ci/reporting_interop/ts_full_receipt.cjs new file mode 100644 index 000000000..52b092fd1 --- /dev/null +++ b/scripts/ci/reporting_interop/ts_full_receipt.cjs @@ -0,0 +1,134 @@ +#!/usr/bin/env node +'use strict'; + +/* Public installed TypeScript client for the receipt and exact-read stage. */ + +const fs = require('node:fs'); +const path = require('node:path'); +const { createRequire } = require('node:module'); +const { callWithProbeAbort, createProbeAbort } = require('./ts_probe_abort.cjs'); + +function argsFrom(argv) { + const result = {}; + for (let i = 0; i < argv.length; i += 2) { + const key = argv[i]; + if (!key?.startsWith('--') || !argv[i + 1] || argv[i + 1].startsWith('--')) { + throw new Error('arguments must be --key value pairs'); + } + if (Object.hasOwn(result, key.slice(2))) throw new Error(`duplicate ${key}`); + result[key.slice(2)] = argv[i + 1]; + } + return result; +} + +function required(args, key) { + if (!args[key]) throw new Error(`missing --${key}`); + return args[key]; +} + +async function main() { + const args = argsFrom(process.argv.slice(2)); + const lock = path.resolve(required(args, 'package-lock')); + const packageRequire = createRequire(path.join(path.dirname(lock), 'package.json')); + const packageName = '@adcp/sdk'; + const installed = packageRequire(`${packageName}/package.json`); + const pinned = JSON.parse(fs.readFileSync(lock, 'utf8')).packages?.[`node_modules/${packageName}`]; + if (installed.version !== required(args, 'expected-version') || + pinned?.version !== installed.version || + pinned?.integrity !== required(args, 'expected-integrity')) { + throw new Error('installed TypeScript artifact does not match its immutable pin'); + } + const sdk = packageRequire(packageName); + if (typeof sdk.ProtocolClient?.callTool !== 'function' || + typeof sdk.unwrapProtocolResponse !== 'function' || + typeof sdk.closeMCPConnections !== 'function') { + throw new Error('installed TypeScript client lacks the required public calls'); + } + const endpoint = new URL(required(args, 'endpoint')); + if (!['http:', 'https:'].includes(endpoint.protocol) || endpoint.username || + endpoint.password || endpoint.search || endpoint.hash) { + throw new Error('invalid seller endpoint'); + } + const token = process.env[required(args, 'auth-env')]; + if (!token) throw new Error('seller authorization token is missing'); + const version = required(args, 'adcp-version'); + const agent = { + agent_uri: endpoint.toString(), auth_token: token, + id: 'full-interop-seller', name: 'full-interop-seller', protocol: 'mcp', + }; + const abort = createProbeAbort('installed full-lifecycle receipt'); + const call = async (name, params) => { + const raw = await callWithProbeAbort(abort, signal => sdk.ProtocolClient.callTool( + agent, name, params, { + adcpVersion: version, wireAdcpVersion: version, signal, + transport: { requestTimeoutMs: 30_000, maxResponseBytes: 16 * 1024 * 1024 }, + }, + )); + return sdk.unwrapProtocolResponse(raw, name, 'mcp', { responseAdcpVersion: version }); + }; + try { + const account = { account_id: 'acct_a' }; + const revision = await call('get_reporting_status', { + account, view: 'revision', reporting_revision_id: 'production-revision', + }); + const materializations = revision.materializations || []; + if (materializations.length !== 1 || + materializations[0].reporting_revision_id !== 'production-revision' || + materializations[0].status !== 'delivered') { + throw new Error('exact revision read did not return the delivered winner'); + } + const materialization = materializations[0]; + const proof = materialization.verification; + if (!proof || proof.row_count !== 2 || + proof.verification_profile !== 'canonical_digest' || + !proof.canonical_content_digest?.value) { + throw new Error('exact revision read lacks canonical managed evidence'); + } + const submission = { + account, idempotency_key: '10000000-0000-4000-8000-000000000001', + receipts: [{ + reporting_receipt_id: 'interop-receipt-0001', + reporting_obligation_id: materialization.reporting_obligation_id, + reporting_revision_id: materialization.reporting_revision_id, + reporting_materialization_id: materialization.reporting_materialization_id, + status: 'accepted', verification_profile: proof.verification_profile, + observed_row_count: proof.row_count, + observed_control_totals: proof.control_totals, + observed_canonical_content_digest: proof.canonical_content_digest, + observed_at: new Date().toISOString(), + }], + }; + const accepted = await call('sync_reporting_receipts', submission); + if (accepted.status !== 'completed' || accepted.results?.length !== 1 || + accepted.results[0].receipt?.reporting_revision_id !== 'production-revision' || + accepted.results[0].receipt?.status !== 'accepted') { + throw new Error('receipt was not accepted for the exact delivered revision'); + } + const period = await call('get_reporting_status', { account, view: 'periods' }); + if (period.periods?.length !== 1 || period.periods[0].reconciliation_status !== 'accepted' || + period.periods[0].receipt_count !== 1) { + throw new Error('accepted receipt is absent from Reconciled Billing status'); + } + return { + status: 'passed', package_version: installed.version, + reporting_revision_id: materialization.reporting_revision_id, + reporting_materialization_id: materialization.reporting_materialization_id, + canonical_content_digest: proof.canonical_content_digest.value, + accepted_receipt_count: period.periods[0].receipt_count, + reconciliation_status: period.periods[0].reconciliation_status, + exact_revision_read: true, + }; + } finally { + abort.abort('receipt stage settled'); + await sdk.closeMCPConnections(); + abort.dispose(); + } +} + +main().then( + report => process.stdout.write(`${JSON.stringify(report)}\n`), + error => { + process.stderr.write(`${String(error.message || error)}\n`); + process.exitCode = 1; + }, +); diff --git a/tests/test_reporting_full_installed_artifact_matrix.py b/tests/test_reporting_full_installed_artifact_matrix.py new file mode 100644 index 000000000..38ddbea33 --- /dev/null +++ b/tests/test_reporting_full_installed_artifact_matrix.py @@ -0,0 +1,99 @@ +"""Fail-closed controls for the installed full reporting lifecycle gate.""" + +from __future__ import annotations + +import copy +import importlib.util +import sys +from pathlib import Path + +import pytest + +ROOT = Path(__file__).resolve().parents[1] +SCRIPTS = ROOT / "scripts" / "ci" / "reporting_interop" +sys.path.insert(0, str(SCRIPTS)) +SPEC = importlib.util.spec_from_file_location( + "reporting_full_installed_matrix", SCRIPTS / "run_full_installed_artifact_matrix.py" +) +assert SPEC is not None and SPEC.loader is not None +matrix = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(matrix) + + +def _cell(index: int) -> dict: + return { + "positive_full_lifecycle": True, + "database_identity": { + "name": f"adcp_1199_full_cell_{index}", + "cluster_identity": { + "system_identifier": "cluster-1", + "postmaster_started_at": "2026-09-28T00:00:00Z", + }, + }, + "seller_process_identity": {"pid": 1000 + index, "start_token": str(2000 + index)}, + "output_directory": f"cell-{index}", + } + + +def test_full_gate_requires_four_positive_isolated_cells() -> None: + cells = [_cell(index) for index in range(4)] + assert matrix._blocking_acceptance(cells) + + omitted = copy.deepcopy(cells) + omitted.pop() + assert not matrix._blocking_acceptance(omitted) + + failed = copy.deepcopy(cells) + failed[2]["positive_full_lifecycle"] = False + assert not matrix._blocking_acceptance(failed) + + leaked = copy.deepcopy(cells) + leaked[3]["database_identity"] = leaked[0]["database_identity"] + assert not matrix._blocking_acceptance(leaked) + + +def _full_evidence() -> dict: + return { + "receipt": { + "status": "passed", + "exact_revision_read": True, + "accepted_receipt_count": 1, + "reconciliation_status": "accepted", + "reporting_revision_id": "production-revision", + }, + "buyer": { + "status": "passed", + "inspection_calls": 0, + "reconciliation": { + "definitive": True, + "obligations": [{"reportingObligationId": "rpo_acct_a"}], + }, + }, + "webhook": { + "attempt_count": 2, + "body_unchanged": True, + "idempotency_key_unchanged": True, + "verified_signature_count": 2, + "signature_generations": [1, 2], + }, + } + + +@pytest.mark.parametrize( + ("stage", "field", "value"), + [ + ("receipt", "exact_revision_read", False), + ("receipt", "reconciliation_status", "pending"), + ("buyer", "inspection_calls", 1), + ("webhook", "body_unchanged", False), + ("webhook", "verified_signature_count", 1), + ("webhook", "signature_generations", [1, 1]), + ], +) +def test_full_gate_refuses_missing_lifecycle_evidence( + stage: str, field: str, value: object +) -> None: + evidence = _full_evidence() + assert matrix._full_result(**evidence) + evidence[stage][field] = value + assert not matrix._full_result(**evidence) From 4ceac7baaa1a7ea3001c6b653467dc4bdf2f1413 Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Mon, 28 Sep 2026 09:05:54 +0000 Subject: [PATCH 2/3] test(reporting): guard required full lifecycle CI stage --- ...st_reporting_full_installed_artifact_matrix.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/tests/test_reporting_full_installed_artifact_matrix.py b/tests/test_reporting_full_installed_artifact_matrix.py index 38ddbea33..a83c4520f 100644 --- a/tests/test_reporting_full_installed_artifact_matrix.py +++ b/tests/test_reporting_full_installed_artifact_matrix.py @@ -52,6 +52,21 @@ def test_full_gate_requires_four_positive_isolated_cells() -> None: assert not matrix._blocking_acceptance(leaked) +def test_required_installed_job_runs_both_core_and_full_lifecycle_stages() -> None: + workflow = (ROOT / ".github/workflows/ci.yml").read_text(encoding="utf-8") + job_start = workflow.index(" reporting-installed-artifact-matrix:") + job_end = workflow.index("\n reporting-installed-artifact-required-gate:", job_start) + job = workflow[job_start:job_end] + gate_end = workflow.index("\n v3-reference-seller-tests:", job_end) + gate = workflow[job_end:gate_end] + + assert "run_installed_artifact_matrix.py" in job + assert "run_full_installed_artifact_matrix.py" in job + assert "pytest==9.0.2" in job + assert "needs: reporting-installed-artifact-matrix" in gate + assert 'needs.reporting-installed-artifact-matrix.result }}" != "success"' in gate + + def _full_evidence() -> dict: return { "receipt": { From 39ae07c51951a3e40a6590f8c191e2ed7feb7bbb Mon Sep 17 00:00:00 2001 From: Brian O'Kelley Date: Mon, 28 Sep 2026 09:23:29 +0000 Subject: [PATCH 3/3] test(reporting): allow durable webhook replay to finish during shutdown --- scripts/ci/reporting_interop/run_foundation_matrix.py | 4 ++-- .../reporting_interop/run_full_installed_artifact_matrix.py | 4 +++- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/scripts/ci/reporting_interop/run_foundation_matrix.py b/scripts/ci/reporting_interop/run_foundation_matrix.py index 025cae4a8..07a7844a0 100644 --- a/scripts/ci/reporting_interop/run_foundation_matrix.py +++ b/scripts/ci/reporting_interop/run_foundation_matrix.py @@ -1462,12 +1462,12 @@ def _run_resource_location_comparison( } -def _stop(process: subprocess.Popen[bytes]) -> None: +def _stop(process: subprocess.Popen[bytes], *, grace_seconds: float = 10) -> None: if process.poll() is not None: return try: os.killpg(process.pid, signal.SIGTERM) - process.wait(timeout=10) + process.wait(timeout=grace_seconds) except subprocess.TimeoutExpired: os.killpg(process.pid, signal.SIGKILL) process.wait(timeout=5) diff --git a/scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py b/scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py index 41cbff4f5..752479ab3 100644 --- a/scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py +++ b/scripts/ci/reporting_interop/run_full_installed_artifact_matrix.py @@ -206,7 +206,9 @@ def _run_cell( ) as error: issues.append(f"cell_process_failed:{type(error).__name__}:{error}") finally: - foundation._stop(process) + # Shutdown verifies the durable notification retry against PostgreSQL. + # Give that bounded work time to write evidence before SIGKILL. + foundation._stop(process, grace_seconds=90) if webhook_path.is_file(): try: webhook = json.loads(webhook_path.read_text(encoding="utf-8"))