Module adcp.reporting
Reliable Reporting: buyer-side reconciliation and seller-side production.
Two halves of the same ledger live under this package.
Buyer side (:mod:adcp.reporting._reconcile, re-exported here)
:func:reconcile_reporting_core() and friends read a seller's
get_reporting_status ledger and decide whether reporting for a scope is
definitive.
These names were previously importable from adcp.reporting
when it was a single module; that import path is unchanged.
Seller side (:mod:adcp.reporting.source, :mod:adcp.reporting.conformance)
The transport-independent producer contract: an executor fetches one frozen
slice from a reporting source and returns one immutable
:class:~adcp.reporting.source.SourceBatchManifestV1 whose bytes are
canonical_json_utf8_v1.
The conformance validators are the executable
definition of "conforming"; run them in your own test suite.
:mod:<code><a title="adcp.reporting.inline_source" href="inline_source.html">adcp.reporting.inline\_source</a></code> is the on-ramp: wrap the delivery fetch
you already have and it produces conforming publications for you.
:mod:<code><a title="adcp.reporting.ledger" href="ledger/index.html">adcp.reporting.ledger</a></code> is the other half -- obligations, immutable
revisions, the <code>reporting.core</code> health projection, and a
framework-agnostic <code>get\_reporting\_status</code> handler.
Submodules are imported lazily (:pep:562) so import adcp.reporting stays
cheap for buyers who never touch the producer contract.
Sub-modules
adcp.reporting.adjustment_evidence-
Buyer adjustment evidence and deterministic receipt construction …
adcp.reporting.canonical_json-
canonical_json_utf8_v1— the byte-exact reporting evidence encoding … adcp.reporting.conformance-
Executable conformance for reporting source executors …
adcp.reporting.currency-
Currency checks shared by trusted obligation writers and source adapters …
adcp.reporting.evidence-
Small immutable public evidence values, independent of transports and credentials.
adcp.reporting.feed-
Frozen authorized reporting feeds; PostgreSQL remains an optional lazy import.
adcp.reporting.fixtures-
Redacted, deterministic fixtures for the reporting source contract …
adcp.reporting.inline_source-
Wrap an existing delivery fetch into a conforming reporting source executor …
adcp.reporting.inline_storage-
Borrowed-pool PostgreSQL storage for inline objects and committed replay seals …
adcp.reporting.ledger-
Reliable Reporting
reporting.core, seller side … adcp.reporting.materializer-
Destination contracts, verification and optional durable materialization …
adcp.reporting.migration-
Explicit, stopped-worker upgrade of reporting state lacking caller ownership …
adcp.reporting.outbox-
Optional transactional reporting notifications, activity and status lifecycle.
adcp.reporting.ownership-
Additive, page-local revision ownership without changing protocol schemas …
adcp.reporting.production-
Production reporting composition with explicit, frozen provider contracts …
adcp.reporting.projection-
Versioned private reporting projection, independent of optional delivery workers.
adcp.reporting.receipts-
Durable seller receipt ingress. Optional PG driver is imported lazily.
adcp.reporting.revision_selection-
Strict, pure publication selection shared by seller and buyer projections …
adcp.reporting.service-
Adapter-first orchestration for AdCP Reliable Reporting …
adcp.reporting.service_lifecycle-
Owned lifecycle and admission for the adapter-first reporting service …
adcp.reporting.service_production-
Domain inputs for the adapter-first managed reporting factory …
adcp.reporting.source-
The seller-side reporting source contract: one frozen slice in, one immutable manifest out …
adcp.reporting.submissions-
Buyer-only durable receipt submission, separate from reconciliation planning …
adcp.reporting.testing-
Deterministic adapter, failure, and service lifecycle test helpers …
Functions
async def assert_reporting_lifecycle_replay(expected: ReportingLifecycleSnapshot,
store: ReportingLedgerStore,
*,
page_size: int = 500,
max_pages: int = 1000) ‑> None-
Expand source code
async def assert_reporting_lifecycle_replay( expected: ReportingLifecycleSnapshot, store: ReportingLedgerStore, *, page_size: int = 500, max_pages: int = 1000, ) -> None: """Assert no lost, changed, or extra publication after restart and retry.""" actual = await capture_reporting_lifecycle( store, account_id=expected.obligation.account_id, reporting_obligation_id=expected.obligation.reporting_obligation_id, page_size=page_size, max_pages=max_pages, ) if actual != expected: raise AssertionError("reporting obligation, revisions, or row bytes changed after replay")Assert no lost, changed, or extra publication after restart and retry.
def build_reporting_receipt(context: ReportingInspectionContext,
observation: ReportingObservation,
*,
reporting_receipt_id: str | None = None,
observed_at: datetime | None = None) ‑> ReportingReceipt-
Expand source code
def build_reporting_receipt( context: ReportingInspectionContext, observation: ReportingObservation, *, reporting_receipt_id: str | None = None, observed_at: datetime | None = None, ) -> ReportingReceipt: materialization = context.materialization revision = context.revision if not materialization.verification or not materialization.resource: raise ReportingReconciliationError( "MATERIALIZATION_NOT_READY", "cannot receipt an unverified materialization" ) failures: list[str] = [] if observation.row_count != revision.row_count: failures.append("ROW_COUNT_MISMATCH") if _totals(observation.control_totals) != _totals(revision.control_totals): failures.append("CONTROL_TOTAL_MISMATCH") profile = _enum(materialization.verification.verification_profile) if profile == "canonical_digest" and ( not revision.canonical_content_digest or _json(observation.canonical_content_digest) != _json(revision.canonical_content_digest) ): failures.append("CANONICAL_DIGEST_MISMATCH") if ( profile == "manifest_checksums" and observation.manifest_sha256 != materialization.resource.manifest_sha256 ): failures.append("MANIFEST_DIGEST_MISMATCH") if ( profile == "native_commit" and observation.native_version_ref != materialization.resource.native_version_ref ): failures.append("NATIVE_VERSION_MISMATCH") payload: dict[str, object] = { "reporting_receipt_id": reporting_receipt_id or f"reporting-receipt:{uuid4()}", "reporting_obligation_id": context.obligation.reporting_obligation_id, "reporting_revision_id": revision.reporting_revision_id, "reporting_materialization_id": materialization.reporting_materialization_id, "status": "rejected" if failures else "accepted", "verification_profile": profile, "observed_row_count": observation.row_count, "observed_control_totals": observation.control_totals, "observed_at": observed_at or datetime.now(timezone.utc), } if observation.canonical_content_digest: payload["observed_canonical_content_digest"] = observation.canonical_content_digest if observation.manifest_sha256: payload["observed_manifest_sha256"] = observation.manifest_sha256 if observation.native_version_ref: payload["observed_native_version_ref"] = observation.native_version_ref if observation.consumer_commit_ref: payload["consumer_commit_ref"] = observation.consumer_commit_ref if failures: payload["rejection_codes"] = failures return ReportingReceipt.model_validate(payload) async def capture_reporting_lifecycle(store: ReportingLedgerStore,
*,
account_id: str,
reporting_obligation_id: str,
page_size: int = 500,
max_pages: int = 1000) ‑> adcp.reporting._testing_lifecycle.ReportingLifecycleSnapshot-
Expand source code
async def capture_reporting_lifecycle( store: ReportingLedgerStore, *, account_id: str, reporting_obligation_id: str, page_size: int = 500, max_pages: int = 1000, ) -> ReportingLifecycleSnapshot: """Read and verify one published obligation through public store methods. Works with in-memory, PostgreSQL, or adopter stores. All pages must repeat the exact revision identity/count and exhaust within ``max_pages``. The complete walk is hashed against the committed Core binding, including the typed control totals when the publisher supplied managed evidence. Keep workers quiescent during capture/replay comparison. For memory tests, retain the original stores across service instances; that models service recreation only. PostgreSQL tests should create fresh store instances over the retained database. No private state is copied or repaired by the kit. """ if type(page_size) is not int or not 1 <= page_size <= 500: raise ValueError("page_size must be an integer between 1 and 500") if type(max_pages) is not int or max_pages <= 0: raise ValueError("max_pages must be a positive integer") obligation = await store.get_obligation( account_id=account_id, reporting_obligation_id=reporting_obligation_id ) if obligation is None: raise AssertionError("reporting obligation did not survive lifecycle recovery") if (obligation.account_id, obligation.reporting_obligation_id) != ( account_id, reporting_obligation_id, ): raise AssertionError("obligation read returned a different account or identity") revisions = tuple( sorted( await store.list_revisions( account_id=account_id, reporting_obligation_id=reporting_obligation_id ), key=lambda item: item.reporting_revision_id, ) ) if not revisions: raise AssertionError("lifecycle conformance requires at least one published revision") if len({item.reporting_revision_id for item in revisions}) != len(revisions): raise AssertionError("revision listing repeated an identity") contents: list[bytes] = [] for revision in revisions: revision_id = revision.reporting_revision_id if (revision.account_id, revision.reporting_obligation_id) != ( account_id, reporting_obligation_id, ): raise AssertionError("revision read returned a different account or obligation") exact = await store.get_revision(account_id=account_id, reporting_revision_id=revision_id) if exact != revision: raise AssertionError("revision listing disagrees with the exact revision read") rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() for _ in range(max_pages): page = await store.read_revision_rows( account_id=account_id, reporting_revision_id=revision_id, cursor=cursor, limit=page_size, ) if ( page.reporting_revision_id != revision_id or type(page.total_count) is not int or page.total_count != revision.row_count or type(page.has_more) is not bool or len(page.rows) > page_size ): raise AssertionError("revision page changed its identity, count, or page bound") rows.extend(page.rows) if len(rows) > revision.row_count: raise AssertionError("revision walk returned more rows than its binding") if not page.has_more: if page.cursor is not None: raise AssertionError("final revision page retained a continuation cursor") break if not page.rows or not isinstance(page.cursor, str) or not page.cursor: raise AssertionError("revision continuation made no progress") if page.cursor in seen: raise AssertionError("revision continuation repeated a cursor") seen.add(page.cursor) cursor = page.cursor else: raise AssertionError("revision walk exceeded max_pages") if len(rows) != revision.row_count: raise AssertionError("revision walk ended before all bound rows were read") digest = revision_content_sha256( reporting_revision_id=revision_id, row_count=revision.row_count, control_totals=revision.control_totals, control_total_evidence=revision.managed_control_totals, reporting_rows=rows, ) if digest != revision.revision_content_sha256: raise AssertionError("revision rows do not match the committed content binding") contents.append(canonical_json_utf8_v1(rows)) return ReportingLifecycleSnapshot(obligation, revisions, tuple(contents))Read and verify one published obligation through public store methods.
Works with in-memory, PostgreSQL, or adopter stores. All pages must repeat the exact revision identity/count and exhaust within
max_pages. The complete walk is hashed against the committed Core binding, including the typed control totals when the publisher supplied managed evidence.Keep workers quiescent during capture/replay comparison. For memory tests, retain the original stores across service instances; that models service recreation only. PostgreSQL tests should create fresh store instances over the retained database. No private state is copied or repaired by the kit.
def classify_content_mismatch(*,
obligation: ReportingObligation,
revision: ReportingRevision,
reading: ReportingContentReading,
definition: ReportingPinnedDefinition | None = None) ‑> Literal['scope_media_buy_missing', 'coverage_short', 'metric_missing', 'schema_nonconformant', 'currency_mismatch', 'period_mismatch'] | None-
Expand source code
def classify_content_mismatch( *, obligation: ReportingObligation, revision: ReportingRevision, reading: ReportingContentReading, definition: ReportingPinnedDefinition | None = None, ) -> MismatchCode | None: """Which agreed fact this revision contradicts, or ``None`` if it honours them. Checked in the spec's precedence order. Two orderings matter and are not arbitrary: * ``metric_missing`` is checked *before* ``schema_nonconformant``. The spec is explicit: a metric that is simply absent uses ``metric_missing`` even when the pinned schema declares it required. Checking structure first would collapse every missing metric into a generic schema failure and lose the one detail that tells the seller what to fix. * ``scope_media_buy_missing`` is checked before ``coverage_short`` because a missing buy is the coarser failure: reporting the package shortfall of a revision that is missing a whole buy describes a symptom, not the cause. Returns ``None`` rather than raising when nothing is wrong, so a caller can use it as the condition for filing ``received``. """ frozen_buys = _ids(obligation.media_buy_ids) present_buys = set(reading.media_buy_ids) if any(buy not in present_buys for buy in frozen_buys): return "scope_media_buy_missing" covered = _ids(getattr(obligation.coverage, "covered_package_ids", None)) if covered and not set(covered) <= set(reading.package_ids): return "coverage_short" if definition is not None: present_metrics = set(reading.metric_names) if any(name not in present_metrics for name in definition.metric_names): return "metric_missing" if not reading.schema_conformant: return "schema_nonconformant" if definition is not None: for name, unit in definition.metric_units.items(): observed = reading.metric_units.get(name) if observed is not None and observed != unit: return "currency_mismatch" for name, unit in definition.control_total_units.items(): observed = reading.control_total_units.get(name) if observed is not None and observed != unit: return "currency_mismatch" if reading.time_dimension_values: start = _utc(obligation.period.start) end = _utc(obligation.period.end) # Half-open: [start, end). A value exactly at the end belongs to the # next period, and a revision that carries it is reporting outside the # window the obligation froze. if any(not (start <= _utc(value) < end) for value in reading.time_dimension_values): return "period_mismatch" del revision # Identity is asserted by the caller's recomputed digest. return NoneWhich agreed fact this revision contradicts, or
Noneif it honours them.Checked in the spec's precedence order. Two orderings matter and are not arbitrary:
metric_missingis checked beforeschema_nonconformant. The spec is explicit: a metric that is simply absent usesmetric_missingeven when the pinned schema declares it required. Checking structure first would collapse every missing metric into a generic schema failure and lose the one detail that tells the seller what to fix.scope_media_buy_missingis checked beforecoverage_shortbecause a missing buy is the coarser failure: reporting the package shortfall of a revision that is missing a whole buy describes a symptom, not the cause.
Returns
Nonerather than raising when nothing is wrong, so a caller can use it as the condition for filingreceived. def consumer_status_chain_key(intent: ConsumerStatusIntent,
*,
account_id: str,
consumer_id: str | None = None) ‑> str-
Expand source code
def consumer_status_chain_key( intent: ConsumerStatusIntent, *, account_id: str, consumer_id: str | None = None, ) -> str: """The logical chain an intent belongs to, as a checkpoint key. Keyed exactly as the seller keys it: authenticated consumer, account, configuration generation, report definition, and the half-open period. ``account_id`` is required rather than optional because the spec's chain identity is account-scoped, and a checkpoint store shared across sellers -- or across accounts on one seller -- would otherwise let two chains with the same ``delivery_config_id`` overwrite each other's leaves. The buyer would then supersede the wrong statement, or none. ``consumer_id`` is accepted for a buyer that posts as more than one authenticated consumer. A single-identity buyer can leave it out; the key stays stable either way as long as the same value is used for reads and writes. Deliberately *not* keyed on the obligation id: a chain that began as ``obligation_missing`` attaches to the repaired obligation later and must not fork. """ return _chain_key_of( account_id, consumer_id, ( intent.delivery_config_id, intent.delivery_config_version, intent.report_definition_id, _iso(intent.period_start), _iso(intent.period_end), ), )The logical chain an intent belongs to, as a checkpoint key.
Keyed exactly as the seller keys it: authenticated consumer, account, configuration generation, report definition, and the half-open period.
account_idis required rather than optional because the spec's chain identity is account-scoped, and a checkpoint store shared across sellers – or across accounts on one seller – would otherwise let two chains with the samedelivery_config_idoverwrite each other's leaves. The buyer would then supersede the wrong statement, or none.consumer_idis accepted for a buyer that posts as more than one authenticated consumer. A single-identity buyer can leave it out; the key stays stable either way as long as the same value is used for reads and writes.Deliberately not keyed on the obligation id: a chain that began as
obligation_missingattaches to the repaired obligation later and must not fork. def evaluate_reporting_ledger(ledger: ReportingLedger,
*,
expected_periods: list[ExpectedReportingPeriod] | None = None,
now: datetime | None = None) ‑> adcp.reporting._reconcile.ReportingReconciliationResult-
Expand source code
def evaluate_reporting_ledger( ledger: ReportingLedger, *, expected_periods: list[ExpectedReportingPeriod] | None = None, now: datetime | None = None, ) -> ReportingReconciliationResult: """Evaluate exact retained evidence without performing reads or writes. Definitive and missing-period claims require an unchanged full snapshot from :func:`load_reporting_ledger`. Manually assembled ledgers are diagnostic only; this API does not certify a caller's incremental merge or recompute raw adjustment digests from normalized models. """ complete_read = ledger._read_fingerprint is not None if complete_read and ledger._read_fingerprint != _ledger_fingerprint(ledger): raise ReportingReconciliationError( "LEDGER_CHANGED", "the completed ledger was changed after loading" ) leaves, partition_complete, index = _validate_read_ledger(ledger) complete_read = complete_read and partition_complete now = now or datetime.now(timezone.utc) outcomes: list[ObligationReconciliation] = [] unique_revisions: dict[str, ReportingRevision] = {} for obligation in ledger.obligations: revision, materialization, selected_reasons = index.selections[ obligation.reporting_obligation_id ] reasons = list(selected_reasons) if not complete_read: reasons.append("UNVERIFIED_LEDGER_SNAPSHOT") if _enum(obligation.health) != "complete": reasons.append(f"OBLIGATION_{_enum(obligation.health).upper()}") if ( materialization and materialization.resource and materialization.resource.expires_at <= now ): reasons.append("RESOURCE_EXPIRED") if revision: unique_revisions[revision.reporting_revision_id] = revision if _enum(obligation.reconciliation_mode) == "consumer_receipt" and revision: if _enum(revision.finality) != "official" and "FINALITY_NOT_MET" not in reasons: reasons.append("FINALITY_NOT_MET") receipt = leaves.get( ("revision", obligation.reporting_obligation_id, revision.reporting_revision_id) ) # Validation binds this leaf to the materialization it names, not # the newest artifact selected independently for current readability. if receipt is None or _enum(receipt.status) != "accepted": reasons.append("MISSING_MATCHING_CONSUMER_RECEIPT") if revision: # Applicable corrections need accepted current evidence in either mode. applicable = index.adjustments_by_revision.get(revision.reporting_revision_id, []) if any( ( leaf := leaves.get( ("adjustment", a.reporting_adjustment_id, revision.reporting_revision_id) ) ) is None or _enum(leaf.status) != "accepted" for a in applicable ): reasons.append("MISSING_MATCHING_ADJUSTMENT_RECEIPT") outcomes.append( ObligationReconciliation( obligation.reporting_obligation_id, not reasons, revision.reporting_revision_id if revision else None, materialization.reporting_materialization_id if materialization else None, tuple(reasons), ) ) actual = { ( item.delivery_config_id, item.delivery_config_version, item.report_definition_id, _enum(item.feed_purpose), item.reporting_profile, _identifiers(item.media_buy_ids), item.period.start.isoformat(), item.period.end.isoformat(), item.period.source_timezone, ) for item in ledger.obligations } generations = { (g.delivery_config_id, g.delivery_config_version, _enum(g.feed_purpose)) for g in getattr(ledger.scope, "delivery_config_generations", []) } media_buy_ids = set(_identifiers(getattr(ledger.scope, "media_buy_ids", []))) expectations = [ ( item, _expected_in_scope(item, ledger.scope, generations, media_buy_ids), ( item.delivery_config_id, item.delivery_config_version, item.report_definition_id, item.feed_purpose, item.reporting_profile, tuple(sorted(item.media_buy_ids)), _iso(item.period_start), _iso(item.period_end), item.source_timezone, ) in actual, ) for item in expected_periods or [] ] # The producer describes configuration requirements here, not current # revision finality: [official] can be a complete unfiltered denominator. # ExpectedReportingPeriod has no trusted finality fact, however, so for an # *absent* obligation we cannot prove membership in a proper subset. This # intentionally withholds some valid absence claims; it is not evidence of # seller filtering. Positive returned obligations retain their exact proof. absence_finality_proven = {_enum(value) for value in getattr(ledger.scope, "finality", [])} == { "snapshot", "official", } missing = [ item for item, in_scope, present in expectations if complete_read and absence_finality_proven and in_scope and not present ] definitive = bool( complete_read and expected_periods is not None and all(in_scope and present for _, in_scope, present in expectations) and bool(getattr(ledger.scope, "scope_closed", False)) and bool(getattr(ledger.scope, "coverage_complete", False)) and all(item.definitive for item in outcomes) ) return ReportingReconciliationResult( definitive, ledger, outcomes, missing, totals_by_revision=[ (item.reporting_revision_id, item.row_count, item.control_totals) for item in unique_revisions.values() ], )Evaluate exact retained evidence without performing reads or writes.
Definitive and missing-period claims require an unchanged full snapshot from :func:
load_reporting_ledger(). Manually assembled ledgers are diagnostic only; this API does not certify a caller's incremental merge or recompute raw adjustment digests from normalized models. async def load_consumer_loop_view(client: Any,
request: GetReportingStatusRequest,
*,
capabilities: Mapping[str, Any] | None = None) ‑> adcp.reporting._consumer.ConsumerLoopView-
Expand source code
async def load_consumer_loop_view( client: Any, request: GetReportingStatusRequest, *, capabilities: Mapping[str, Any] | None = None, ) -> ConsumerLoopView: """Read the summary view for the buyer-side fields of the rc.3 loop. ``consumer_status_pending`` is ``None`` when the seller did not emit it, which means it does not advertise ``consumer_status_task``. That is distinct from ``0``: zero is "you owe nothing", ``None`` is "this seller has no loop to owe anything to", and conflating them would make a buyer think it was caught up on a seller that never asked. """ payload = request.model_dump(mode="json", exclude_none=True) payload["view"] = "summary" payload.pop("pagination", None) result = await client.get_reporting_status(GetReportingStatusRequest.model_validate(payload)) response = result.data if not result.success or response is None: from adcp.reporting._reconcile import ReportingReconciliationError raise ReportingReconciliationError( "STATUS_READ_FAILED", "get_reporting_status did not return a completed summary view" ) counts = getattr(response, "obligation_counts", None) pending = getattr(counts, "consumer_status_pending", None) if counts is not None else None return ConsumerLoopView( consumer_status_pending=pending, issues=tuple(getattr(response, "issues", None) or ()), operations_contact=ReportingOperationsContactView.from_capability(capabilities), )Read the summary view for the buyer-side fields of the rc.3 loop.
consumer_status_pendingisNonewhen the seller did not emit it, which means it does not advertiseconsumer_status_task. That is distinct from0: zero is "you owe nothing",Noneis "this seller has no loop to owe anything to", and conflating them would make a buyer think it was caught up on a seller that never asked. async def load_reporting_ledger(client: ReportingStatusClient,
request: GetReportingStatusRequest,
*,
max_snapshot_restarts: int = 2,
max_pages: int = 2048,
max_records: int = 200000,
max_bytes: int = 67108864) ‑> adcp.reporting._reconcile.ReportingLedger-
Expand source code
async def load_reporting_ledger( client: ReportingStatusClient, request: GetReportingStatusRequest, *, max_snapshot_restarts: int = 2, max_pages: int = 2048, max_records: int = 200_000, max_bytes: int = 64 * 1024 * 1024, ) -> ReportingLedger: """Load a complete authenticated periods snapshot, never an incremental delta. Restart at the first page, retaining scope and the requested page size. Incremental, health, finality and exact-revision result selectors cannot prove completeness and are removed. Page, received-row (including replays), and serialized typed-page byte limits bound the walk. Transport adapters must also bound raw responses. """ if ( type(max_snapshot_restarts) is not int or max_snapshot_restarts < 0 or type(max_pages) is not int or max_pages < 1 or type(max_records) is not int or max_records < 1 or type(max_bytes) is not int or max_bytes < 1 ): raise ValueError("reporting walk bounds must be positive (restarts may be zero)") prepared = None try: base = request.model_dump(mode="json", exclude_none=True, warnings="error") base["view"] = "periods" base.pop("pagination", None) for selector in ("changes_after", "reporting_revision_id", "health", "finality"): base.pop(selector, None) page_size = request.pagination.max_results if request.pagination else None requested_account = request.account.model_dump(mode="json", warnings="error").get( "account_id" ) prepared = (base, page_size, requested_account) except Exception: prepared = None if prepared is None: raise ReportingReconciliationError( "INVALID_STATUS_REQUEST", "reporting status request could not be constructed" ) base, page_size, requested_account = prepared for restart in range(max_snapshot_restarts + 1): try: obligations: dict[str, ReportingObligation] = {} revisions: dict[str, ReportingRevision] = {} materializations: dict[str, ReportingMaterialization] = {} receipts: dict[str, ReportingReceipt] = {} consumer_statuses: dict[str, Any] = {} adjustments: dict[str, ReportingAdjustment] = {} adjustment_receipts: dict[str, ReportingAdjustmentReceipt] = {} ownership: dict[str, str] = {} ownership_mode: bool | None = None frozen_metadata: str | None = None cursor: str | None = None seen_cursors: set[str] = set() snapshot_id: str | None = None ledger_as_of: datetime | None = None account_id: str | None = None scope: BaseModel | None = None total_count: int | None = None received_records = 0 received_bytes = 0 for _page_number in range(max_pages): payload = dict(base) pagination_request: dict[str, object] = {} if page_size is not None: pagination_request["max_results"] = page_size if cursor: pagination_request["cursor"] = cursor if pagination_request: payload["pagination"] = pagination_request page_request = None try: page_request = GetReportingStatusRequest.model_validate(payload) except Exception: page_request = None if page_request is None: # SDK-side construction and provider failures are distinct, # but neither may retain validation inputs or private context. raise ReportingReconciliationError( "INVALID_STATUS_REQUEST", "reporting status request could not be constructed", ) result = None try: result = await client.get_reporting_status(page_request) except Exception: result = None if result is None: # Leave the exception scope: do not retain a private provider # exception or Pydantic input in __context__ either. raise ReportingReconciliationError( "STATUS_READ_FAILED", "get_reporting_status could not be read" ) response = result.data if ( not result.success or response is None or _enum(response.view) != "periods" or _enum(response.status) != "completed" ): raise ReportingReconciliationError( "STATUS_READ_FAILED", "get_reporting_status did not return a completed periods view", ) received_bytes += len(response.model_dump_json(exclude_unset=True).encode("utf-8")) received_records += sum( len(records or []) for records in ( response.periods, response.revisions, response.materializations, response.receipts, response.consumer_statuses, response.adjustments, response.adjustment_receipts, ) ) if received_bytes > max_bytes or received_records > max_records: raise ReportingReconciliationError( "LEDGER_LIMIT_EXCEEDED", "ledger read budget exceeded" ) response = response.model_copy(deep=True) raw_page = response.model_dump(mode="json", exclude_none=True) if "ext" in response.model_fields_set and response.ext is None: raw_page["ext"] = None try: local = page_revision_ownership(raw_page) except ReportingOwnershipError: raise ReportingReconciliationError( "INVALID_REVISION_OWNERSHIP", "invalid page-local revision ownership" ) from None mode = local is not None if ownership_mode is not None and mode != ownership_mode: raise ReportingReconciliationError( "INVALID_REVISION_OWNERSHIP", "mixed ownership modes within one snapshot" ) ownership_mode = mode for revision_id, owner_id in (local or {}).items(): if revision_id in ownership and ownership[revision_id] != owner_id: raise ReportingReconciliationError( "INVALID_REVISION_OWNERSHIP", "ownership changed within one snapshot" ) ownership[revision_id] = owner_id metadata = _json( { k: raw_page.get(k) for k in ( "changes_checkpoint", "next_expected_at", "health", "issues", "obligation_counts", "coverage", "data_through", ) } ) if frozen_metadata is not None and metadata != frozen_metadata: raise ReportingReconciliationError( "SNAPSHOT_CHANGED", "frozen projection changed" ) frozen_metadata = metadata pagination = response.pagination if ( not response.ledger_snapshot_id or not response.ledger_as_of or not response.account_id or not response.scope or pagination is None or pagination.total_count is None or type(pagination.has_more) is not bool ): raise ReportingReconciliationError( "INCOMPLETE_LEDGER_PAGE", "get_reporting_status omitted ledger metadata" ) if requested_account is not None and response.account_id != requested_account: raise ReportingReconciliationError( "LEDGER_SCOPE_MISMATCH", "ledger does not match the requested account" ) if snapshot_id and snapshot_id != response.ledger_snapshot_id: raise ReportingReconciliationError("SNAPSHOT_CHANGED", "snapshot changed") if ledger_as_of and ledger_as_of != response.ledger_as_of: raise ReportingReconciliationError( "SNAPSHOT_CHANGED", "ledger boundary changed" ) if account_id and account_id != response.account_id: raise ReportingReconciliationError("SNAPSHOT_CHANGED", "account changed") if scope and _json(scope) != _json(response.scope): raise ReportingReconciliationError("SNAPSHOT_CHANGED", "denominator changed") if total_count is not None and total_count != pagination.total_count: raise ReportingReconciliationError("SNAPSHOT_CHANGED", "record total changed") snapshot_id = response.ledger_snapshot_id ledger_as_of = response.ledger_as_of account_id = response.account_id scope = response.scope total_count = pagination.total_count if total_count > max_records: raise ReportingReconciliationError( "LEDGER_LIMIT_EXCEEDED", "ledger record limit exceeded" ) for obligation in response.periods or []: _add_immutable( obligations, obligation.reporting_obligation_id, obligation, "obligation", ) for revision in response.revisions or []: _add_immutable(revisions, revision.reporting_revision_id, revision, "revision") for materialization in response.materializations or []: _add_immutable( materializations, materialization.reporting_materialization_id, materialization, "materialization", ) for receipt in response.receipts or []: _add_immutable(receipts, receipt.reporting_receipt_id, receipt, "receipt") for adjustment in response.adjustments or []: _add_immutable( adjustments, adjustment.reporting_adjustment_id, adjustment, "adjustment" ) for adjustment_receipt in response.adjustment_receipts or []: _add_immutable( adjustment_receipts, adjustment_receipt.reporting_receipt_id, adjustment_receipt, "adjustment receipt", ) for status in getattr(response, "consumer_statuses", None) or []: _add_immutable( consumer_statuses, status.reporting_status_id, status, "consumer status", ) if ( sum( len(records) for records in ( obligations, revisions, materializations, receipts, consumer_statuses, adjustments, adjustment_receipts, ) ) > max_records ): raise ReportingReconciliationError( "LEDGER_LIMIT_EXCEEDED", "ledger record limit exceeded" ) if not pagination.has_more: break cursor = pagination.cursor if not cursor or len(cursor) > 2048 or cursor in seen_cursors: raise ReportingReconciliationError( "CURSOR_LOOP", "ledger pagination did not advance" ) seen_cursors.add(cursor) else: raise ReportingReconciliationError( "LEDGER_LIMIT_EXCEEDED", "ledger page limit exceeded" ) count = ( len(obligations) + len(revisions) + len(materializations) + len(receipts) + len(consumer_statuses) + len(adjustments) + len(adjustment_receipts) ) if total_count != count: raise ReportingReconciliationError( "LEDGER_COUNT_MISMATCH", f"ledger declared {total_count} records but returned {count}", ) if not snapshot_id or not ledger_as_of or not account_id or not scope: raise ReportingReconciliationError( "EMPTY_LEDGER_RESPONSE", "get_reporting_status returned no ledger page" ) ledger = ReportingLedger( snapshot_id, ledger_as_of, account_id, scope, list(obligations.values()), list(revisions.values()), list(materializations.values()), list(receipts.values()), list(consumer_statuses.values()), ownership if ownership_mode else None, list(adjustments.values()), list(adjustment_receipts.values()), ) _, partition_complete, _ = _validate_read_ledger(ledger) if partition_complete: ledger._read_fingerprint = _ledger_fingerprint(ledger) return ledger except ReportingReconciliationError as error: if error.code != "SNAPSHOT_CHANGED" or restart == max_snapshot_restarts: raise raise ReportingReconciliationError("SNAPSHOT_CHANGED", "ledger never stabilized")Load a complete authenticated periods snapshot, never an incremental delta.
Restart at the first page, retaining scope and the requested page size. Incremental, health, finality and exact-revision result selectors cannot prove completeness and are removed. Page, received-row (including replays), and serialized typed-page byte limits bound the walk. Transport adapters must also bound raw responses.
def plan_consumer_statuses(obligations: Sequence[ReportingObligation],
*,
now: datetime,
automated_recovery_window: timedelta,
missing_expected_periods: Sequence[Any] = (),
readings: Mapping[str, ReportingContentReading] | None = None,
definition: ReportingPinnedDefinition | None = None,
revisions: Mapping[str, ReportingRevision] | None = None,
obligation_revisions: Mapping[str, Sequence[ReportingRevision]] | None = None,
current_statuses: Sequence[Any] = (),
checkpointed_leaves: Mapping[str, ConsumerStatusCheckpoint] | None = None,
account_id: str | None = None,
consumer_id: str | None = None,
status_id_prefix: str = 'rpcs') ‑> list[adcp.reporting._consumer.ConsumerStatusIntent]-
Expand source code
def plan_consumer_statuses( obligations: Sequence[ReportingObligation], *, now: datetime, automated_recovery_window: timedelta, missing_expected_periods: Sequence[Any] = (), readings: Mapping[str, ReportingContentReading] | None = None, definition: ReportingPinnedDefinition | None = None, revisions: Mapping[str, ReportingRevision] | None = None, obligation_revisions: Mapping[str, Sequence[ReportingRevision]] | None = None, current_statuses: Sequence[Any] = (), checkpointed_leaves: Mapping[str, ConsumerStatusCheckpoint] | None = None, account_id: str | None = None, consumer_id: str | None = None, status_id_prefix: str = "rpcs", ) -> list[ConsumerStatusIntent]: """Decide what the buyer owes the seller right now. One intent per *expected* period whose deadline -- ``expected_at`` plus the seller's advertised ``automated_recovery_window_seconds`` -- has passed. Periods that are not due yet are skipped: filing early is not wrong, but filing ``revision_missing`` before the seller is late would be a false accusation this loop is supposed to prevent. "Expected" deliberately includes periods the seller's ledger omitted. ``missing_expected_periods`` (what ``evaluate_reporting_ledger`` computed by diffing the buyer's own denominator against the ledger) each become ``obligation_missing`` -- keyed without a seller obligation id, because requiring one would make the first missing report invisible again, which is the exact failure this loop exists to surface. Status selection, in the spec's terms: * no obligation at all -> ``obligation_missing``; * obligation with no required revision -> ``revision_missing``; * revision advertised but unreadable -> ``unreadable`` with the reading's closed ``failure_code``; * revision read -> ``received``, or ``content_mismatch`` when :func:`classify_content_mismatch` finds a contradicted contract fact. The third and fourth cases need the caller's own read result, so an obligation that *has* a required revision and *no* reading raises :class:`ConsumerStatusPlanError` rather than defaulting. Defaulting either way makes a false claim: ``revision_missing`` accuses the seller of publishing nothing it demonstrably published, and ``received`` asserts bytes the buyer never looked at. ``obligation_revisions`` must carry each due obligation's *complete* retained history -- the same partition its ``revision_count`` declares. Finality is chosen by whole-history selection (:func:`~adcp.reporting.revision_selection.select_reporting_revision`), not by a count, so a partition that is absent, short, long or structurally damaged raises :class:`ConsumerStatusPlanError` instead of selecting from it. An obligation with ``revision_count == 0`` needs no entry. ``current_statuses`` is this caller's own status history from the same ledger read. An intent whose content matches the current leaf is skipped entirely -- re-filing an unchanged claim under a fresh id churns the chain for nothing, and a derived id that collided with the leaf would make the statement supersede *itself*. """ readings = readings or {} revisions = revisions or {} obligation_revisions = obligation_revisions or {} leaves = _index_current_leaves(current_statuses) checkpointed = checkpointed_leaves or {} if checkpointed and account_id is None: # Without the account the derived key cannot match the one the poster # wrote, so every lookup would miss silently and every statement would # go out with no supersedes. Refuse rather than quietly degrade. raise ConsumerStatusPlanError( "checkpointed_leaves requires account_id: the chain key is account-scoped, and " "a mismatched key misses silently rather than failing" ) boundary = _utc(now) plan: list[ConsumerStatusIntent] = [] def add(intent: ConsumerStatusIntent | None) -> None: """Append unless the current leaf already says exactly this.""" if intent is not None: plan.append(intent) for period in missing_expected_periods: expected_at = _expected_at_of(period) if expected_at is None: # Without the buyer's own expected_at there is no defensible # deadline and no way to satisfy "obligation_missing is valid only # at or after expected_at". Skip rather than invent one. continue due_at = expected_at + automated_recovery_window if boundary < due_at: continue add( _intend( status="obligation_missing", delivery_config_id=period.delivery_config_id, delivery_config_version=period.delivery_config_version, report_definition_id=period.report_definition_id, period_start=_parse(period.period_start), period_end=_parse(period.period_end), source_timezone=getattr(period, "source_timezone", None) or "UTC", status_as_of=max(boundary, expected_at), due_at=due_at, leaves=leaves, checkpointed=checkpointed, account_id=account_id, consumer_id=consumer_id, prefix=status_id_prefix, ) ) for obligation in obligations: obligation_id = obligation.reporting_obligation_id due_at = _utc(obligation.expected_at) + automated_recovery_window if boundary < due_at: continue reading = readings.get(obligation_id) history = obligation_revisions.get(obligation_id) if history is not None: _has_required_revision(obligation, history) # Reject damage even when a reading exists. status: ConsumerStatusValue mismatch: MismatchCode | None = None failure: ReportingFailureCode | None = None revision_id: str | None = None digest: str | None = None status_as_of = boundary if reading is not None and reading.failure_code is not None: status = "unreadable" failure = reading.failure_code revision_id = reading.reporting_revision_id elif reading is not None: revision = revisions.get(reading.reporting_revision_id) if revision is None: # The buyer claims to have read a revision this ledger snapshot # does not contain. Filing `received` would bind the statement # to a revision the seller cannot resolve, and the seller must # reject it -- so say so here, where the mistake is fixable. raise ConsumerStatusPlanError( f"reading names revision {reading.reporting_revision_id!r}, which is not " f"in this ledger snapshot for obligation {obligation_id!r}; re-read the " "snapshot rather than filing a status against a revision the seller " "cannot resolve" ) mismatch = classify_content_mismatch( obligation=obligation, revision=revision, reading=reading, definition=definition, ) status = "content_mismatch" if mismatch else "received" revision_id = reading.reporting_revision_id digest = reading.observed_revision_content_sha256 if status == "received": # Buyer-attributed arrival evidence, not planning time. status_as_of = _utc(reading.first_consumable_at or boundary) elif _has_required_revision(obligation, obligation_revisions.get(obligation_id)): raise ConsumerStatusPlanError( f"obligation {obligation_id!r} has a required revision but no reading was " "supplied; pass a ReportingContentReading (with failure_code when the read " "failed) rather than letting this default to revision_missing, which would " "accuse the seller of publishing nothing" ) else: status = "revision_missing" add( _intend( status=status, delivery_config_id=obligation.delivery_config_id, delivery_config_version=obligation.delivery_config_version, report_definition_id=obligation.report_definition_id, period_start=_utc(obligation.period.start), period_end=_utc(obligation.period.end), source_timezone=_source_timezone(obligation), status_as_of=status_as_of, due_at=due_at, leaves=leaves, checkpointed=checkpointed, account_id=account_id, consumer_id=consumer_id, prefix=status_id_prefix, obligation_id=obligation_id, revision_id=revision_id, digest=digest, mismatch=mismatch, failure=failure, ) ) return planDecide what the buyer owes the seller right now.
One intent per expected period whose deadline –
expected_atplus the seller's advertisedautomated_recovery_window_seconds– has passed. Periods that are not due yet are skipped: filing early is not wrong, but filingrevision_missingbefore the seller is late would be a false accusation this loop is supposed to prevent."Expected" deliberately includes periods the seller's ledger omitted.
missing_expected_periods(whatevaluate_reporting_ledger()computed by diffing the buyer's own denominator against the ledger) each becomeobligation_missing– keyed without a seller obligation id, because requiring one would make the first missing report invisible again, which is the exact failure this loop exists to surface.Status selection, in the spec's terms:
- no obligation at all ->
obligation_missing; - obligation with no required revision ->
revision_missing; - revision advertised but unreadable ->
unreadablewith the reading's closedfailure_code; - revision read ->
received, orcontent_mismatchwhen :func:classify_content_mismatch()finds a contradicted contract fact.
The third and fourth cases need the caller's own read result, so an obligation that has a required revision and no reading raises :class:
ConsumerStatusPlanErrorrather than defaulting. Defaulting either way makes a false claim:revision_missingaccuses the seller of publishing nothing it demonstrably published, andreceivedasserts bytes the buyer never looked at.obligation_revisionsmust carry each due obligation's complete retained history – the same partition itsrevision_countdeclares. Finality is chosen by whole-history selection (:func:~adcp.reporting.revision_selection.select_reporting_revision), not by a count, so a partition that is absent, short, long or structurally damaged raises :class:ConsumerStatusPlanErrorinstead of selecting from it. An obligation withrevision_count == 0needs no entry.current_statusesis this caller's own status history from the same ledger read. An intent whose content matches the current leaf is skipped entirely – re-filing an unchanged claim under a fresh id churns the chain for nothing, and a derived id that collided with the leaf would make the statement supersede itself. - no obligation at all ->
async def post_consumer_statuses(client: Any,
plan: Sequence[ConsumerStatusIntent],
*,
account_id: str,
consumer_id: str | None = None,
checkpoints: ConsumerStatusCheckpointStore | None = None,
idempotency_key_prefix: str = 'rpcs_batch') ‑> adcp.reporting._consumer.ConsumerStatusPostResult-
Expand source code
async def post_consumer_statuses( client: Any, plan: Sequence[ConsumerStatusIntent], *, account_id: str, consumer_id: str | None = None, checkpoints: ConsumerStatusCheckpointStore | None = None, idempotency_key_prefix: str = "rpcs_batch", ) -> ConsumerStatusPostResult: """Post a plan through ``sync_reporting_status``, batching as the schema allows. Each batch carries at most one statement per logical chain, because the spec rejects every duplicate-chain entry in a batch *without evaluating their supersession order* -- so two statements for one chain in one request lose both. Chains are therefore spread across batches rather than packed. ``recorded`` and ``unchanged`` are both successes: an exact retry replaying as ``unchanged`` is the idempotency contract working, not a problem to report. Only ``failed`` entries are surfaced as failures, and they are returned rather than raised so one bad statement does not discard the rest. The idempotency key is derived from the batch's contents, so a retry after a transport failure reuses it and the seller can replay rather than re-evaluate. """ if not plan: return ConsumerStatusPostResult() recorded: list[ConsumerStatusIntent] = [] unchanged: list[ConsumerStatusIntent] = [] failed: list[tuple[ConsumerStatusIntent, str, str]] = [] for batch in _batches(plan): by_id = {intent.reporting_status_id: intent for intent in batch} request = { "account": {"account_id": account_id}, "idempotency_key": _batch_idempotency_key(idempotency_key_prefix, batch), "statuses": [intent.to_wire() for intent in batch], } results = await _submit(client, request) for entry in results: intent = _match_result(entry, by_id) if intent is None: continue outcome = str(entry.get("result")) if outcome == "recorded": recorded.append(intent) elif outcome == "unchanged": unchanged.append(intent) else: errors = entry.get("errors") or [{}] failed.append( ( intent, str(errors[0].get("code", "UNKNOWN")), str(errors[0].get("message", "")), ) ) continue if checkpoints is not None: await checkpoints.put( consumer_status_chain_key( intent, account_id=account_id, consumer_id=consumer_id ), ConsumerStatusCheckpoint( reporting_status_id=intent.reporting_status_id, supersedes_reporting_status_id=intent.supersedes_reporting_status_id, ), ) return ConsumerStatusPostResult( recorded=tuple(recorded), unchanged=tuple(unchanged), failed=tuple(failed) )Post a plan through
sync_reporting_status, batching as the schema allows.Each batch carries at most one statement per logical chain, because the spec rejects every duplicate-chain entry in a batch without evaluating their supersession order – so two statements for one chain in one request lose both. Chains are therefore spread across batches rather than packed.
recordedandunchangedare both successes: an exact retry replaying asunchangedis the idempotency contract working, not a problem to report. Onlyfailedentries are surfaced as failures, and they are returned rather than raised so one bad statement does not discard the rest.The idempotency key is derived from the batch's contents, so a retry after a transport failure reuses it and the seller can replay rather than re-evaluate.
async def reconcile_reporting(client: ReportingReconciliationClient,
request: GetReportingStatusRequest,
inspect: Callable[[ReportingInspectionContext], Awaitable[ReportingObservation]] | None = None,
*,
expected_periods: list[ExpectedReportingPeriod],
resource_reader: ReportingResourceReader | None = None,
checkpoint_store: ReportingCheckpointStore | None = None,
max_snapshot_restarts: int = 2,
max_inspection_attempts: int = 3,
inspection_timeout_seconds: float = 30.0,
inspection_retry_backoff_seconds: float = 1.0,
now: datetime | None = None,
reporting_capabilities: ReportingDeliveryCapabilities | None = None) ‑> ReportingReconciliationResult-
Expand source code
async def reconcile_reporting( client: ReportingReconciliationClient, request: GetReportingStatusRequest, inspect: Callable[[ReportingInspectionContext], Awaitable[ReportingObservation]] | None = None, *, expected_periods: list[ExpectedReportingPeriod], resource_reader: ReportingResourceReader | None = None, checkpoint_store: ReportingCheckpointStore | None = None, max_snapshot_restarts: int = 2, max_inspection_attempts: int = 3, inspection_timeout_seconds: float = 30.0, inspection_retry_backoff_seconds: float = 1.0, now: datetime | None = None, reporting_capabilities: ReportingDeliveryCapabilities | None = None, ) -> ReportingReconciliationResult: """Reconcile a closed ledger, persist observations, and submit receipts. Pass ``resource_reader`` for the built-in manifest/file inspector, or ``inspect`` as an advanced adapter for warehouses and native shares. Each inspection is time-bounded. Typed transient failures retry with exponential backoff; permanent integrity failures stop immediately. """ if inspect is not None and resource_reader is not None: raise ValueError("pass inspect or resource_reader, not both") if resource_reader is not None: from adcp.reporting_inspection import ManifestReportingInspector inspect = ManifestReportingInspector(resource_reader) if ( not isinstance(max_inspection_attempts, int) or isinstance(max_inspection_attempts, bool) or max_inspection_attempts < 1 ): raise ValueError("max_inspection_attempts must be at least 1") if not isfinite(inspection_timeout_seconds) or inspection_timeout_seconds <= 0: raise ValueError("inspection_timeout_seconds must be finite and greater than 0") if not isfinite(inspection_retry_backoff_seconds) or inspection_retry_backoff_seconds < 0: raise ValueError("inspection_retry_backoff_seconds must be finite and not negative") ledger = await load_reporting_ledger( client, request, max_snapshot_restarts=max_snapshot_restarts ) if reporting_capabilities is not None: tiers = reporting_tiers(reporting_capabilities) if ReportingTier.MANAGED_DELIVERY not in tiers and any( item.destination_ref is not None for item in ledger.obligations ): raise ReportingReconciliationError( "MANAGED_DELIVERY_NOT_ENABLED", "reporting obligations require the managed_delivery tier", ) if ReportingTier.RECONCILED_BILLING not in tiers and any( _enum(item.reconciliation_mode) == "consumer_receipt" or _enum(item.feed_purpose) == "billing" for item in ledger.obligations ): raise ReportingReconciliationError( "RECONCILED_BILLING_NOT_ENABLED", "consumer receipts and billing reporting require the reconciled_billing tier", ) if ReportingTier.RECONCILED_BILLING not in tiers and any( item.canonical_content_digest is not None for item in ledger.revisions ): raise ReportingReconciliationError( "RECONCILED_BILLING_NOT_ENABLED", "canonical-digest reporting requires the reconciled_billing tier", ) if ReportingTier.RECONCILED_BILLING not in tiers and any( item.verification and _enum(item.verification.verification_profile) == "canonical_digest" for item in ledger.materializations ): raise ReportingReconciliationError( "RECONCILED_BILLING_NOT_ENABLED", "canonical-digest verification requires the reconciled_billing tier", ) submitted: list[ReportingReceipt] = [] _, _, index = _validate_read_ledger(ledger) for obligation in ledger.obligations: if _enum(obligation.reconciliation_mode) != "consumer_receipt": continue revision, materialization, reasons = index.selections[obligation.reporting_obligation_id] if not revision or not materialization or reasons: continue if any(_receipt_matches(item, revision, materialization) for item in ledger.receipts): continue receipt = ( await checkpoint_store.get(materialization.reporting_materialization_id) if checkpoint_store else None ) checkpoint_is_recorded = bool( receipt and any( item.reporting_receipt_id == receipt.reporting_receipt_id for item in ledger.receipts ) ) if ( not receipt or checkpoint_is_recorded or not _receipt_targets(receipt, obligation, revision, materialization) ): if inspect is None: raise ReportingReconciliationError( "INSPECTOR_REQUIRED", "consumer-receipt reconciliation requires inspect or resource_reader", ) last_error: Exception | None = None observation = None inspection_context = ReportingInspectionContext(obligation, revision, materialization) for attempt in range(max_inspection_attempts): try: observation = await asyncio.wait_for( inspect(inspection_context), timeout=inspection_timeout_seconds ) break except Exception as error: # destination SDKs define their own transient errors if getattr(error, "retryable", None) is False: code = getattr(error, "code", "INSPECTION_FAILED") raise ReportingReconciliationError(_enum(code), str(error)) from error last_error = ( TimeoutError( "materialization inspection timed out after " f"{inspection_timeout_seconds:g} seconds" ) if isinstance(error, asyncio.TimeoutError) else error ) if attempt + 1 < max_inspection_attempts: await asyncio.sleep(inspection_retry_backoff_seconds * (2**attempt)) if observation is None: raise ReportingReconciliationError( "INSPECTION_FAILED", "materialization inspection failed after " f"{max_inspection_attempts} attempts: {last_error}", ) receipt = build_reporting_receipt(inspection_context, observation) if checkpoint_store: await checkpoint_store.put(receipt) submitted.append(receipt) if submitted: write = await client.sync_reporting_receipts( SyncReportingReceiptsRequest.model_validate( { "account": request.account, "idempotency_key": str(uuid4()), "receipts": submitted, } ) ) if not write.success or write.data is None: raise ReportingReconciliationError( "RECEIPT_WRITE_FAILED", "seller did not record reporting receipts" ) failed = [item for item in write.data.results if _enum(item.result) == "failed"] if failed: raise ReportingReconciliationError( "RECEIPT_WRITE_FAILED", f"{len(failed)} reporting receipt(s) failed" ) ledger = await load_reporting_ledger( client, request, max_snapshot_restarts=max_snapshot_restarts ) result = evaluate_reporting_ledger(ledger, expected_periods=expected_periods, now=now) result.submitted_receipts = submitted return resultReconcile a closed ledger, persist observations, and submit receipts.
Pass
resource_readerfor the built-in manifest/file inspector, orinspectas an advanced adapter for warehouses and native shares. Each inspection is time-bounded. Typed transient failures retry with exponential backoff; permanent integrity failures stop immediately. async def reconcile_reporting_core(client: ReportingStatusClient,
request: GetReportingStatusRequest,
*,
expected_periods: list[ExpectedReportingPeriod],
max_snapshot_restarts: int = 2,
now: datetime | None = None,
automated_recovery_window: timedelta | None = None,
readings: dict[str, ReportingContentReading] | None = None,
pinned_definition: ReportingPinnedDefinition | None = None,
capabilities: dict[str, Any] | None = None) ‑> adcp.reporting._reconcile.ReportingReconciliationResult-
Expand source code
async def reconcile_reporting_core( client: ReportingStatusClient, request: GetReportingStatusRequest, *, expected_periods: list[ExpectedReportingPeriod], max_snapshot_restarts: int = 2, now: datetime | None = None, automated_recovery_window: timedelta | None = None, readings: dict[str, ReportingContentReading] | None = None, pinned_definition: ReportingPinnedDefinition | None = None, capabilities: dict[str, Any] | None = None, ) -> ReportingReconciliationResult: """Reconcile the Core API-delivered tier without destination handling. A Core caller only compares obligations, revisions, coverage, finality, and the reporting clock. Destination materializations, manifests, digests, and consumer receipts are deliberately rejected rather than accidentally activating a higher tier. When ``automated_recovery_window`` is supplied the result also carries the AdCP 3.2.0-rc.3 buyer-side loop: the seller's ``obligation_counts.consumer_status_pending``, the issue lifecycle fields, ``operations_contact``, and a plan of the statements this buyer owes *by its deadline* rather than by scope close. It is opt-in because posting status is only a duty when the seller advertises ``consumer_status_task`` -- inferring a reverse endpoint from a seller that never offered one is the failure the opt-in exists to prevent. """ ledger = await load_reporting_ledger( client, request, max_snapshot_restarts=max_snapshot_restarts ) if any(obligation.destination_ref is not None for obligation in ledger.obligations): raise ReportingReconciliationError( "MANAGED_DELIVERY_NOT_ENABLED", "Core reconciliation received a managed-delivery obligation", ) if ledger.materializations or ledger.receipts: raise ReportingReconciliationError( "MANAGED_DELIVERY_NOT_ENABLED", "Core reconciliation received destination or receipt records", ) if any( _enum(obligation.reconciliation_mode) == "consumer_receipt" for obligation in ledger.obligations ): raise ReportingReconciliationError( "RECONCILED_BILLING_NOT_ENABLED", "Core reconciliation received a consumer-receipt obligation", ) result = evaluate_reporting_ledger(ledger, expected_periods=expected_periods, now=now) if automated_recovery_window is not None: result.consumer_loop = await load_consumer_loop_view( client, request, capabilities=capabilities ) result.consumer_status_plan = plan_consumer_statuses( ledger.obligations, now=now or ledger.ledger_as_of, automated_recovery_window=automated_recovery_window, # The periods the *buyer's* denominator expects and the seller's # ledger omitted. Dropping these was the whole reason # ``obligation_missing`` could never be emitted -- and it is the # one status the seller cannot derive for itself. missing_expected_periods=result.missing_expected_periods, readings=readings, definition=pinned_definition, revisions={item.reporting_revision_id: item for item in ledger.revisions}, obligation_revisions={ obligation.reporting_obligation_id: _owned_revisions(obligation, ledger) for obligation in ledger.obligations }, current_statuses=ledger.consumer_statuses, account_id=ledger.account_id, ) return resultReconcile the Core API-delivered tier without destination handling.
A Core caller only compares obligations, revisions, coverage, finality, and the reporting clock. Destination materializations, manifests, digests, and consumer receipts are deliberately rejected rather than accidentally activating a higher tier.
When
automated_recovery_windowis supplied the result also carries the AdCP 3.2.0-rc.3 buyer-side loop: the seller'sobligation_counts.consumer_status_pending, the issue lifecycle fields,operations_contact, and a plan of the statements this buyer owes by its deadline rather than by scope close. It is opt-in because posting status is only a duty when the seller advertisesconsumer_status_task– inferring a reverse endpoint from a seller that never offered one is the failure the opt-in exists to prevent. def reporting_tiers(capabilities: ReportingDeliveryCapabilities) ‑> frozenset[adcp.reporting._reconcile.ReportingTier]-
Expand source code
def reporting_tiers( capabilities: ReportingDeliveryCapabilities, ) -> frozenset[ReportingTier]: """Project reporting capability flags into their cumulative SDK tiers.""" tiers = {ReportingTier.CORE} if capabilities.managed_delivery: tiers.add(ReportingTier.MANAGED_DELIVERY) if capabilities.reconciled_billing: if not capabilities.managed_delivery: raise ReportingReconciliationError( "INVALID_REPORTING_CAPABILITIES", "reconciled_billing requires managed_delivery", ) tiers.add(ReportingTier.RECONCILED_BILLING) return frozenset(tiers)Project reporting capability flags into their cumulative SDK tiers.
async def resolve_checkpointed_leaves(checkpoints: ConsumerStatusCheckpointStore,
*,
account_id: str,
consumer_id: str | None = None,
obligations: Sequence[ReportingObligation] = (),
missing_expected_periods: Sequence[Any] = ()) ‑> dict[str, adcp.reporting._consumer.ConsumerStatusCheckpoint]-
Expand source code
async def resolve_checkpointed_leaves( checkpoints: ConsumerStatusCheckpointStore, *, account_id: str, consumer_id: str | None = None, obligations: Sequence[ReportingObligation] = (), missing_expected_periods: Sequence[Any] = (), ) -> dict[str, ConsumerStatusCheckpoint]: """Read the checkpoints a plan over these periods could need. The planner is synchronous and the store is not, so the reads happen here and the result is handed in. Chain keys are derivable before planning -- they depend only on the configuration generation, report definition, and period -- so this does not need to know what the plan will decide. """ keys: set[str] = set() for obligation in obligations: keys.add( _chain_key_of( account_id, consumer_id, ( obligation.delivery_config_id, obligation.delivery_config_version, obligation.report_definition_id, _iso(_utc(obligation.period.start)), _iso(_utc(obligation.period.end)), ), ) ) for period in missing_expected_periods: keys.add( _chain_key_of( account_id, consumer_id, ( period.delivery_config_id, period.delivery_config_version, period.report_definition_id, _iso(_parse(period.period_start)), _iso(_parse(period.period_end)), ), ) ) resolved: dict[str, ConsumerStatusCheckpoint] = {} for key in keys: checkpoint = await checkpoints.get(key) if checkpoint is not None: resolved[key] = checkpoint return resolvedRead the checkpoints a plan over these periods could need.
The planner is synchronous and the store is not, so the reads happen here and the result is handed in. Chain keys are derivable before planning – they depend only on the configuration generation, report definition, and period – so this does not need to know what the plan will decide.
async def run_reporting_adapter_conformance(adapter: ReportingAdapter,
request: ReportingSourceSliceRequestV1,
*,
clock: Callable[[], datetime] | None = None,
staging: ReportingStagingStore | None = None,
seals: ReportingSealStore | None = None) ‑> SourceBatchManifestV1-
Expand source code
async def run_reporting_adapter_conformance( adapter: ReportingAdapter, request: ReportingSourceSliceRequestV1, *, clock: Callable[[], datetime] | None = None, staging: ReportingStagingStore | None = None, seals: ReportingSealStore | None = None, ) -> SourceBatchManifestV1: """Wrap one minimal adapter and run the SDK's full replay conformance. The adapter needs only ``capabilities`` and ``fetch_slice``. This helper uses fresh memory stores by default. Supply retained memory stores or freshly constructed PostgreSQL stores to test recovery with the same request; create their schemas first. Stores and the adapter are borrowed. The same clock drives both execution timestamps and conformance deadlines. """ staging = staging if staging is not None else InMemoryStagingStore() seals = seals if seals is not None else InMemorySealStore() test_clock = clock or DeterministicReportingClock(request.period.source_read_cutoff_at) executor = InlineReportingSource( capabilities=adapter.capabilities, fetch=adapter.fetch_slice, staging=staging, seals=seals, clock=test_clock, ) return await run_reporting_source_replay_conformance( executor=executor, request=request, object_reader=staging, clock=test_clock, )Wrap one minimal adapter and run the SDK's full replay conformance.
The adapter needs only
capabilitiesandfetch_slice. This helper uses fresh memory stores by default. Supply retained memory stores or freshly constructed PostgreSQL stores to test recovery with the same request; create their schemas first. Stores and the adapter are borrowed. The same clock drives both execution timestamps and conformance deadlines.
Classes
class ConsumerLoopView (consumer_status_pending: int | None,
issues: tuple[ReportingStatusIssue, ...] = (),
operations_contact: ReportingOperationsContactView | None = None)-
Expand source code
@dataclass(frozen=True) class ConsumerLoopView: """What the seller told this buyer about its own side of the loop. Read from the *summary* view, because ``consumer_status_pending`` lives in ``obligation_counts`` and the periods view does not carry it. """ consumer_status_pending: int | None issues: tuple[ReportingStatusIssue, ...] = () operations_contact: ReportingOperationsContactView | None = None @property def mismatch_issues(self) -> tuple[ReportingStatusIssue, ...]: """Only the issues attributed to this buyer's own statements.""" return tuple( issue for issue in self.issues if str(getattr(issue.code, "value", issue.code)) == "CONSUMER_STATUS_MISMATCH" ) def escalated(self) -> tuple[ReportingStatusIssue, ...]: """Mismatches the seller has escalated to a human. ``recommended_action`` in the ``contact_*`` family is the signal, not severity: a seller past its advertised ``consumer_mismatch_escalation_seconds`` MUST switch the action, and ``wait_for_retry`` cannot survive that boundary. """ return tuple( issue for issue in self.mismatch_issues if str(getattr(issue.recommended_action, "value", issue.recommended_action)).startswith( "contact_" ) )What the seller told this buyer about its own side of the loop.
Read from the summary view, because
consumer_status_pendinglives inobligation_countsand the periods view does not carry it.Instance variables
var consumer_status_pending : int | Nonevar issues : tuple[ReportingStatusIssue, ...]prop mismatch_issues : tuple[ReportingStatusIssue, ...]-
Expand source code
@property def mismatch_issues(self) -> tuple[ReportingStatusIssue, ...]: """Only the issues attributed to this buyer's own statements.""" return tuple( issue for issue in self.issues if str(getattr(issue.code, "value", issue.code)) == "CONSUMER_STATUS_MISMATCH" )Only the issues attributed to this buyer's own statements.
var operations_contact : adcp.reporting._consumer.ReportingOperationsContactView | None
Methods
def escalated(self) ‑> tuple[ReportingStatusIssue, ...]-
Expand source code
def escalated(self) -> tuple[ReportingStatusIssue, ...]: """Mismatches the seller has escalated to a human. ``recommended_action`` in the ``contact_*`` family is the signal, not severity: a seller past its advertised ``consumer_mismatch_escalation_seconds`` MUST switch the action, and ``wait_for_retry`` cannot survive that boundary. """ return tuple( issue for issue in self.mismatch_issues if str(getattr(issue.recommended_action, "value", issue.recommended_action)).startswith( "contact_" ) )Mismatches the seller has escalated to a human.
recommended_actionin thecontact_*family is the signal, not severity: a seller past its advertisedconsumer_mismatch_escalation_secondsMUST switch the action, andwait_for_retrycannot survive that boundary.
class ConsumerStatusCheckpoint (reporting_status_id: str, supersedes_reporting_status_id: str | None = None)-
Expand source code
@dataclass(frozen=True) class ConsumerStatusCheckpoint: """The statement a buyer last posted for one logical chain. Carries what it superseded as well as its own id, because that is what a lost-response retry needs. The seller's replay fingerprint covers ``supersedes_reporting_status_id``, so re-posting the same claim with a *different* supersedes is a different statement and earns an identity conflict rather than the ``unchanged`` a retry is owed. """ reporting_status_id: str supersedes_reporting_status_id: str | None = NoneThe statement a buyer last posted for one logical chain.
Carries what it superseded as well as its own id, because that is what a lost-response retry needs. The seller's replay fingerprint covers
supersedes_reporting_status_id, so re-posting the same claim with a different supersedes is a different statement and earns an identity conflict rather than theunchangeda retry is owed.Instance variables
var reporting_status_id : strvar supersedes_reporting_status_id : str | None
class ConsumerStatusCheckpointStore (*args, **kwargs)-
Expand source code
class ConsumerStatusCheckpointStore(Protocol): """Where a buyer remembers the leaf it last posted per logical chain. A cache of the buyer's own last word, keyed by the same logical chain the seller uses: account, configuration generation, report definition, period. The seller stays authoritative: when a caller supplies ``current_statuses`` from a fresh ledger read, that wins outright. The checkpoint is what lets a buyer plan *without* re-reading the whole periods view every cycle, and in particular what makes a retry after a lost response replay as ``unchanged`` rather than being rejected for not naming the current leaf. """ async def get(self, chain: str) -> ConsumerStatusCheckpoint | None: """The statement this buyer last posted for ``chain``.""" ... async def put(self, chain: str, checkpoint: ConsumerStatusCheckpoint) -> None: """Record the statement just accepted for ``chain``.""" ...Where a buyer remembers the leaf it last posted per logical chain.
A cache of the buyer's own last word, keyed by the same logical chain the seller uses: account, configuration generation, report definition, period.
The seller stays authoritative: when a caller supplies
current_statusesfrom a fresh ledger read, that wins outright. The checkpoint is what lets a buyer plan without re-reading the whole periods view every cycle, and in particular what makes a retry after a lost response replay asunchangedrather than being rejected for not naming the current leaf.Ancestors
- typing.Protocol
- typing.Generic
Methods
async def get(self, chain: str) ‑> adcp.reporting._consumer.ConsumerStatusCheckpoint | None-
Expand source code
async def get(self, chain: str) -> ConsumerStatusCheckpoint | None: """The statement this buyer last posted for ``chain``.""" ...The statement this buyer last posted for
chain. async def put(self,
chain: str,
checkpoint: ConsumerStatusCheckpoint) ‑> None-
Expand source code
async def put(self, chain: str, checkpoint: ConsumerStatusCheckpoint) -> None: """Record the statement just accepted for ``chain``.""" ...Record the statement just accepted for
chain.
class ConsumerStatusIntent (reporting_status_id: str,
consumer_status: ConsumerStatusValue,
delivery_config_id: str,
delivery_config_version: int,
report_definition_id: str,
period_start: datetime,
period_end: datetime,
period_source_timezone: str,
status_as_of: datetime,
due_at: datetime,
supersedes_reporting_status_id: str | None = None,
reporting_obligation_id: str | None = None,
reporting_revision_id: str | None = None,
observed_revision_content_sha256: str | None = None,
failure_code: str | None = None,
mismatch_code: MismatchCode | None = None,
consumer_commit_ref: str | None = None)-
Expand source code
@dataclass(frozen=True) class ConsumerStatusIntent: """One statement the buyer owes, ready to batch into ``sync_reporting_status``. Carries ``due_at`` so a caller can log or alert on how late its own loop is without recomputing the deadline. """ reporting_status_id: str consumer_status: ConsumerStatusValue delivery_config_id: str delivery_config_version: int report_definition_id: str period_start: datetime period_end: datetime period_source_timezone: str status_as_of: datetime due_at: datetime supersedes_reporting_status_id: str | None = None reporting_obligation_id: str | None = None reporting_revision_id: str | None = None observed_revision_content_sha256: str | None = None failure_code: str | None = None mismatch_code: MismatchCode | None = None consumer_commit_ref: str | None = None @property def overdue(self) -> bool: return _utc(self.status_as_of) >= _utc(self.due_at) def to_wire(self) -> dict[str, Any]: """The ``statuses[]`` entry for ``sync_reporting_status``. Absent keys are omitted rather than sent as ``null``: the schema's conditional blocks use ``not: {required: [...]}``, and an explicit null still satisfies JSON Schema's ``required``. Sending ``"failure_code": null`` on a ``received`` statement would therefore be rejected. """ payload: dict[str, Any] = { "reporting_status_id": self.reporting_status_id, "delivery_config_id": self.delivery_config_id, "delivery_config_version": self.delivery_config_version, "report_definition_id": self.report_definition_id, "period": { "start": _iso(self.period_start), "end": _iso(self.period_end), "source_timezone": self.period_source_timezone, }, "consumer_status": self.consumer_status, "status_as_of": _iso(self.status_as_of), } optional = { "supersedes_reporting_status_id": self.supersedes_reporting_status_id, "reporting_obligation_id": self.reporting_obligation_id, "reporting_revision_id": self.reporting_revision_id, "observed_revision_content_sha256": self.observed_revision_content_sha256, "failure_code": self.failure_code, "mismatch_code": self.mismatch_code, "consumer_commit_ref": self.consumer_commit_ref, } payload.update({key: value for key, value in optional.items() if value is not None}) return payloadOne statement the buyer owes, ready to batch into
sync_reporting_status.Carries
due_atso a caller can log or alert on how late its own loop is without recomputing the deadline.Instance variables
var consumer_commit_ref : str | Nonevar consumer_status : Literal['received', 'obligation_missing', 'revision_missing', 'unreadable', 'content_mismatch']var delivery_config_id : strvar delivery_config_version : intvar due_at : datetime.datetimevar failure_code : str | Nonevar mismatch_code : Literal['scope_media_buy_missing', 'coverage_short', 'metric_missing', 'schema_nonconformant', 'currency_mismatch', 'period_mismatch'] | Nonevar observed_revision_content_sha256 : str | Noneprop overdue : bool-
Expand source code
@property def overdue(self) -> bool: return _utc(self.status_as_of) >= _utc(self.due_at) var period_end : datetime.datetimevar period_source_timezone : strvar period_start : datetime.datetimevar report_definition_id : strvar reporting_obligation_id : str | Nonevar reporting_revision_id : str | Nonevar reporting_status_id : strvar status_as_of : datetime.datetimevar supersedes_reporting_status_id : str | None
Methods
def to_wire(self) ‑> dict[str, typing.Any]-
Expand source code
def to_wire(self) -> dict[str, Any]: """The ``statuses[]`` entry for ``sync_reporting_status``. Absent keys are omitted rather than sent as ``null``: the schema's conditional blocks use ``not: {required: [...]}``, and an explicit null still satisfies JSON Schema's ``required``. Sending ``"failure_code": null`` on a ``received`` statement would therefore be rejected. """ payload: dict[str, Any] = { "reporting_status_id": self.reporting_status_id, "delivery_config_id": self.delivery_config_id, "delivery_config_version": self.delivery_config_version, "report_definition_id": self.report_definition_id, "period": { "start": _iso(self.period_start), "end": _iso(self.period_end), "source_timezone": self.period_source_timezone, }, "consumer_status": self.consumer_status, "status_as_of": _iso(self.status_as_of), } optional = { "supersedes_reporting_status_id": self.supersedes_reporting_status_id, "reporting_obligation_id": self.reporting_obligation_id, "reporting_revision_id": self.reporting_revision_id, "observed_revision_content_sha256": self.observed_revision_content_sha256, "failure_code": self.failure_code, "mismatch_code": self.mismatch_code, "consumer_commit_ref": self.consumer_commit_ref, } payload.update({key: value for key, value in optional.items() if value is not None}) return payloadThe
statuses[]entry forsync_reporting_status.Absent keys are omitted rather than sent as
null: the schema's conditional blocks usenot: {required: [...]}, and an explicit null still satisfies JSON Schema'srequired. Sending"failure_code": null<code> on a </code>receivedstatement would therefore be rejected.
class ConsumerStatusPlanError (*args, **kwargs)-
Expand source code
class ConsumerStatusPlanError(ValueError): """The caller's inputs cannot support an honest statement. Raised rather than guessing. Every guess this function could make is a claim attributed to the buyer about the seller's behaviour, and filing the wrong one is worse than refusing to file: ``revision_missing`` accuses the seller of publishing nothing, ``received`` asserts bytes were consumed. """The caller's inputs cannot support an honest statement.
Raised rather than guessing. Every guess this function could make is a claim attributed to the buyer about the seller's behaviour, and filing the wrong one is worse than refusing to file:
revision_missingaccuses the seller of publishing nothing,receivedasserts bytes were consumed.Ancestors
- builtins.ValueError
- builtins.Exception
- builtins.BaseException
class ConsumerStatusPostError (*args, **kwargs)-
Expand source code
class ConsumerStatusPostError(RuntimeError): """One or more statements in a posting pass were rejected."""One or more statements in a posting pass were rejected.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class ConsumerStatusPostResult (recorded: tuple[ConsumerStatusIntent, ...] = (),
unchanged: tuple[ConsumerStatusIntent, ...] = (),
failed: tuple[tuple[ConsumerStatusIntent, str, str], ...] = ())-
Expand source code
@dataclass(frozen=True) class ConsumerStatusPostResult: """What one posting pass actually achieved. Failures are carried rather than raised: ``partial_results`` makes each statement independent, so one stale supersession pointer must not discard four good statements. A caller that wants to fail loudly checks :attr:`failed`. """ recorded: tuple[ConsumerStatusIntent, ...] = () unchanged: tuple[ConsumerStatusIntent, ...] = () failed: tuple[tuple[ConsumerStatusIntent, str, str], ...] = () @property def posted(self) -> int: """Statements the seller accepted, new or replayed.""" return len(self.recorded) + len(self.unchanged) @property def ok(self) -> bool: return not self.failed def raise_for_failures(self) -> None: """Raise when any statement failed, after the successes are recorded. For a caller that would rather crash a scheduled job than let a reporting gap accumulate quietly. """ if not self.failed: return detail = "; ".join( f"{intent.reporting_status_id}: {code} {message}" for intent, code, message in self.failed ) raise ConsumerStatusPostError(f"{len(self.failed)} consumer status(es) failed: {detail}")What one posting pass actually achieved.
Failures are carried rather than raised:
partial_resultsmakes each statement independent, so one stale supersession pointer must not discard four good statements. A caller that wants to fail loudly checks :attr:failed.Instance variables
var failed : tuple[tuple[adcp.reporting._consumer.ConsumerStatusIntent, str, str], ...]prop ok : bool-
Expand source code
@property def ok(self) -> bool: return not self.failed prop posted : int-
Expand source code
@property def posted(self) -> int: """Statements the seller accepted, new or replayed.""" return len(self.recorded) + len(self.unchanged)Statements the seller accepted, new or replayed.
var recorded : tuple[adcp.reporting._consumer.ConsumerStatusIntent, ...]var unchanged : tuple[adcp.reporting._consumer.ConsumerStatusIntent, ...]
Methods
def raise_for_failures(self) ‑> None-
Expand source code
def raise_for_failures(self) -> None: """Raise when any statement failed, after the successes are recorded. For a caller that would rather crash a scheduled job than let a reporting gap accumulate quietly. """ if not self.failed: return detail = "; ".join( f"{intent.reporting_status_id}: {code} {message}" for intent, code, message in self.failed ) raise ConsumerStatusPostError(f"{len(self.failed)} consumer status(es) failed: {detail}")Raise when any statement failed, after the successes are recorded.
For a caller that would rather crash a scheduled job than let a reporting gap accumulate quietly.
class DeterministicReportingClock (now: datetime)-
Expand source code
class DeterministicReportingClock: """A small aware UTC clock that tests can advance without sleeping.""" def __init__(self, now: datetime) -> None: self.set(now) def __call__(self) -> datetime: return self._now def set(self, value: datetime) -> datetime: if value.tzinfo is None or value.utcoffset() is None: raise ValueError("DeterministicReportingClock requires an aware datetime") self._now = value.astimezone(timezone.utc) return self._now def advance(self, delta: timedelta) -> datetime: if delta < timedelta(0): raise ValueError("a deterministic reporting clock cannot move backward") return self.set(self._now + delta)A small aware UTC clock that tests can advance without sleeping.
Methods
def advance(self, delta: timedelta) ‑> datetime.datetime-
Expand source code
def advance(self, delta: timedelta) -> datetime: if delta < timedelta(0): raise ValueError("a deterministic reporting clock cannot move backward") return self.set(self._now + delta) def set(self, value: datetime) ‑> datetime.datetime-
Expand source code
def set(self, value: datetime) -> datetime: if value.tzinfo is None or value.utcoffset() is None: raise ValueError("DeterministicReportingClock requires an aware datetime") self._now = value.astimezone(timezone.utc) return self._now
class ExpectedReportingPeriod (delivery_config_id: str,
delivery_config_version: int,
report_definition_id: str,
feed_purpose: str,
reporting_profile: str,
media_buy_ids: tuple[str, ...],
period_start: str,
period_end: str,
source_timezone: str = 'UTC',
expected_at: str | None = None)-
Expand source code
@dataclass(frozen=True) class ExpectedReportingPeriod: delivery_config_id: str delivery_config_version: int report_definition_id: str feed_purpose: str reporting_profile: str media_buy_ids: tuple[str, ...] period_start: str period_end: str #: Required to *state* a period the seller omitted: the consumer-status #: ``period`` object requires it, and a buyer filing #: ``obligation_missing`` has no seller obligation to read it off. Defaults #: to UTC because that is what an omitted schedule alignment resolves to; #: a configuration with a civil-time alignment must set it explicitly or #: the statement describes a different period than the one expected. source_timezone: str = "UTC" #: The buyer's independently derived ``expected_at``. Needed because #: ``obligation_missing`` is only valid at or after it, and the posting #: deadline is measured from it -- neither is knowable from the seller's #: ledger, which is precisely the point when the period is absent from it. expected_at: str | None = NoneExpectedReportingPeriod(delivery_config_id: 'str', delivery_config_version: 'int', report_definition_id: 'str', feed_purpose: 'str', reporting_profile: 'str', media_buy_ids: 'tuple[str, …]', period_start: 'str', period_end: 'str', source_timezone: 'str' = 'UTC', expected_at: 'str | None' = None)
Instance variables
var delivery_config_id : strvar delivery_config_version : intvar expected_at : str | None-
The buyer's independently derived
expected_at. Needed becauseobligation_missingis only valid at or after it, and the posting deadline is measured from it – neither is knowable from the seller's ledger, which is precisely the point when the period is absent from it. var feed_purpose : strvar media_buy_ids : tuple[str, ...]var period_end : strvar period_start : strvar report_definition_id : strvar reporting_profile : strvar source_timezone : str-
Required to state a period the seller omitted: the consumer-status
periodobject requires it, and a buyer filingobligation_missinghas no seller obligation to read it off. Defaults to UTC because that is what an omitted schedule alignment resolves to; a configuration with a civil-time alignment must set it explicitly or the statement describes a different period than the one expected.
class FaultInjectingReportingSealStore (store: ReportingSealStore,
failures: ReportingFailurePlan)-
Expand source code
class FaultInjectingReportingSealStore: """Borrow a seal store with ``seal.get`` and ``seal`` before/after points. In particular, ``seal.after`` models a lost acknowledgement: replay must discover the original seal instead of refetching or publishing new bytes. """ def __init__(self, store: ReportingSealStore, failures: ReportingFailurePlan) -> None: self._store = store self._failures = failures async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None: await self._failures.hit("seal.get.before") result = await self._store.get( account_id=account_id, source_execution_key=source_execution_key ) await self._failures.hit("seal.get.after") return result async def put( self, *, account_id: str, source_execution_key: str, sealed: SealedSlice ) -> SealedSlice: await self._failures.hit("seal.before") result = await self._store.put( account_id=account_id, source_execution_key=source_execution_key, sealed=sealed ) await self._failures.hit("seal.after") return resultBorrow a seal store with
seal.getandsealbefore/after points.In particular,
seal.aftermodels a lost acknowledgement: replay must discover the original seal instead of refetching or publishing new bytes.Methods
async def get(self, *, account_id: str, source_execution_key: str) ‑> SealedSlice | None-
Expand source code
async def get(self, *, account_id: str, source_execution_key: str) -> SealedSlice | None: await self._failures.hit("seal.get.before") result = await self._store.get( account_id=account_id, source_execution_key=source_execution_key ) await self._failures.hit("seal.get.after") return result async def put(self, *, account_id: str, source_execution_key: str, sealed: SealedSlice) ‑> SealedSlice-
Expand source code
async def put( self, *, account_id: str, source_execution_key: str, sealed: SealedSlice ) -> SealedSlice: await self._failures.hit("seal.before") result = await self._store.put( account_id=account_id, source_execution_key=source_execution_key, sealed=sealed ) await self._failures.hit("seal.after") return result
class FaultInjectingReportingStagingStore (store: ReportingStagingStore,
failures: ReportingFailurePlan)-
Expand source code
class FaultInjectingReportingStagingStore: """Borrow a memory or durable store with ``stage`` / ``read`` fault points. Each operation visits ``<operation>.before`` and ``<operation>.after``. A fault after staging retains the real store's write. Schema creation and resource ownership stay with the caller; this wrapper copies no state. """ def __init__(self, store: ReportingStagingStore, failures: ReportingFailurePlan) -> None: self._store = store self._failures = failures async def stage( self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes ) -> tuple[str, str]: await self._failures.hit("stage.before") result = await self._store.stage( account_id=account_id, source_execution_key=source_execution_key, ordinal=ordinal, payload=payload, ) await self._failures.hit("stage.after") return result async def read( self, *, object_ref: str, object_generation: str, account_id: str, source_scope: Mapping[str, Any], cancel: asyncio.Event, ) -> bytes: await self._failures.hit("read.before") result = await self._store.read( object_ref=object_ref, object_generation=object_generation, account_id=account_id, source_scope=source_scope, cancel=cancel, ) await self._failures.hit("read.after") return resultBorrow a memory or durable store with
stage/readfault points.Each operation visits
<operation>.beforeand<operation>.after. A fault after staging retains the real store's write. Schema creation and resource ownership stay with the caller; this wrapper copies no state.Methods
async def read(self,
*,
object_ref: str,
object_generation: str,
account_id: str,
source_scope: Mapping[str, Any],
cancel: asyncio.Event) ‑> bytes-
Expand source code
async def read( self, *, object_ref: str, object_generation: str, account_id: str, source_scope: Mapping[str, Any], cancel: asyncio.Event, ) -> bytes: await self._failures.hit("read.before") result = await self._store.read( object_ref=object_ref, object_generation=object_generation, account_id=account_id, source_scope=source_scope, cancel=cancel, ) await self._failures.hit("read.after") return result async def stage(self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes) ‑> tuple[str, str]-
Expand source code
async def stage( self, *, account_id: str, source_execution_key: str, ordinal: int, payload: bytes ) -> tuple[str, str]: await self._failures.hit("stage.before") result = await self._store.stage( account_id=account_id, source_execution_key=source_execution_key, ordinal=ordinal, payload=payload, ) await self._failures.hit("stage.after") return result
class InMemoryConsumerStatusCheckpoints-
Expand source code
class InMemoryConsumerStatusCheckpoints: """Process-local checkpoints. Correct, and not durable. A real buyer persists these. Losing them is recoverable -- the next plan that passes ``current_statuses`` from a ledger read re-establishes the leaf -- but a buyer that loses them *and* plans without a ledger read will have its next statement rejected for not superseding the current leaf. """ def __init__(self) -> None: self._leaves: dict[str, ConsumerStatusCheckpoint] = {} async def get(self, chain: str) -> ConsumerStatusCheckpoint | None: return self._leaves.get(chain) async def put(self, chain: str, checkpoint: ConsumerStatusCheckpoint) -> None: self._leaves[chain] = checkpointProcess-local checkpoints. Correct, and not durable.
A real buyer persists these. Losing them is recoverable – the next plan that passes
current_statusesfrom a ledger read re-establishes the leaf – but a buyer that loses them and plans without a ledger read will have its next statement rejected for not superseding the current leaf.Methods
async def get(self, chain: str) ‑> adcp.reporting._consumer.ConsumerStatusCheckpoint | None-
Expand source code
async def get(self, chain: str) -> ConsumerStatusCheckpoint | None: return self._leaves.get(chain) async def put(self,
chain: str,
checkpoint: ConsumerStatusCheckpoint) ‑> None-
Expand source code
async def put(self, chain: str, checkpoint: ConsumerStatusCheckpoint) -> None: self._leaves[chain] = checkpoint
class ObligationReconciliation (reporting_obligation_id: str,
definitive: bool,
reporting_revision_id: str | None = None,
reporting_materialization_id: str | None = None,
reasons: tuple[str, ...] = ())-
Expand source code
@dataclass(frozen=True) class ObligationReconciliation: reporting_obligation_id: str definitive: bool reporting_revision_id: str | None = None reporting_materialization_id: str | None = None reasons: tuple[str, ...] = ()ObligationReconciliation(reporting_obligation_id: 'str', definitive: 'bool', reporting_revision_id: 'str | None' = None, reporting_materialization_id: 'str | None' = None, reasons: 'tuple[str, …]' = ())
Instance variables
var definitive : boolvar reasons : tuple[str, ...]var reporting_materialization_id : str | Nonevar reporting_obligation_id : strvar reporting_revision_id : str | None
class ReliableReportingConfigurationError (message: str,
*,
kind: "Literal['misconfiguration', 'auth_required', 'invalid_request']" = 'misconfiguration')-
Expand source code
class ReliableReportingConfigurationError(ValueError): """A composition failure, missing caller, or invalid buyer configuration. Composition failures default to terminal INTERNAL_ERROR. Only explicit request-time failures may ask the buyer to correct authentication/input. """ def __init__( self, message: str, *, kind: Literal["misconfiguration", "auth_required", "invalid_request"] = "misconfiguration", ) -> None: super().__init__(message) self.kind = kindA composition failure, missing caller, or invalid buyer configuration.
Composition failures default to terminal INTERNAL_ERROR. Only explicit request-time failures may ask the buyer to correct authentication/input.
Ancestors
- builtins.ValueError
- builtins.Exception
- builtins.BaseException
class ReliableReportingService (*,
store: ReportingLedgerStore,
account_context: ReportingContextResolver,
caller_resolver: ReportingCallerResolver | None = None,
consumer_status_enabled: bool = False,
escalation: ReportingDeliveryEscalation | None = None,
clock: Callable[[], datetime] | None = None,
worker_interval: timedelta | None = None,
producer_factory: ProducerFactory = adcp.reporting.ledger.producer.ReportingProducer,
materialization_worker: ReportingBackgroundWorker | None = None,
notification_worker: ReportingBackgroundWorker | None = None,
notification_attempt_store: Any | None = None,
receipt_handler: ReportingReceiptHandler | None = None,
reconciled_billing: bool = False,
worker_error_handler: ReportingWorkerErrorHandler | None = None,
owned_resources: Sequence[ReportingServiceResource] = (),
production: ReportingProductionOptions | None = None,
automated_recovery_window: timedelta = datetime.timedelta(seconds=21600),
status_retention_days: int = 400)-
Expand source code
class ReliableReportingService: """High-level owner of adapters, ledger handlers, workers, and lifecycle. Admitted calls own a task so transport cancellation cannot interrupt their cleanup. Invoke them outside caller-owned ledger transactions; use the low-level store API when composing an ambient transaction batch. Injected resources are borrowed unless explicitly transferred via ``owned_resources``. """ def __init__( self, *, store: ReportingLedgerStore, account_context: ReportingContextResolver, caller_resolver: ReportingCallerResolver | None = None, consumer_status_enabled: bool = False, escalation: ReportingDeliveryEscalation | None = None, clock: Callable[[], datetime] | None = None, worker_interval: timedelta | None = None, producer_factory: ProducerFactory = ReportingProducer, materialization_worker: ReportingBackgroundWorker | None = None, notification_worker: ReportingBackgroundWorker | None = None, notification_attempt_store: Any | None = None, receipt_handler: ReportingReceiptHandler | None = None, reconciled_billing: bool = False, worker_error_handler: ReportingWorkerErrorHandler | None = None, owned_resources: Sequence[ReportingServiceResource] = (), production: ReportingProductionOptions | None = None, automated_recovery_window: timedelta = timedelta(hours=6), status_retention_days: int = 400, ) -> None: if production is not None and ( caller_resolver is not None or producer_factory is not ReportingProducer or worker_interval is not None or materialization_worker is not None or notification_worker is not None or notification_attempt_store is not None or receipt_handler is not None or reconciled_billing or worker_error_handler is not None ): raise ReliableReportingConfigurationError( "production options compose their own workers, admission and receipt handlers" ) effective_clock = clock or (lambda: datetime.now(timezone.utc)) self.store = store self._production: ReportingProductionSupport | None = None self._production_options = production self._initialize_source_storage: Callable[[], Awaitable[None]] | None = None self.sources = ReportingAdapterRegistry(clock=effective_clock) self._context_resolver = account_context self._caller_resolver = caller_resolver or self._default_caller self._consumer_status_enabled = consumer_status_enabled self._escalation = escalation or ReportingDeliveryEscalation() self._clock = effective_clock if automated_recovery_window < timedelta(0) or status_retention_days < 1: raise ReliableReportingConfigurationError("invalid reporting recovery/retention policy") self._automated_recovery_window = automated_recovery_window self._status_retention_days = status_retention_days self._worker_interval = worker_interval self._producer_factory = producer_factory self._materialization_worker = materialization_worker self._notification_worker = notification_worker self._notification_attempt_store = notification_attempt_store self._receipt_handler = receipt_handler self._reconciled_billing = reconciled_billing self._worker_error_handler = worker_error_handler self._bindings: dict[ReportingConfigurationGenerationKey, _Binding] = {} self._pending_configurations: dict[ ReportingConfigurationGenerationKey, ReportingConfiguration ] = {} self._initialized = False self._configuration_lock = asyncio.Lock() self._turn_lock = asyncio.Lock() self._lifecycle = _ServiceLifecycle( self._initialize, resources=( (*owned_resources, ReportingServiceResource(close=self._close_production)) if production is not None else owned_resources ), configuration_error=ReliableReportingConfigurationError, ) @classmethod def from_production( cls, production: ReportingProductionSupport, *, owned_resources: Sequence[ReportingServiceResource] = (), ) -> ReliableReportingService: """Own a preassembled B2 graph with durable fixed-profile admission. Mount the exact handler returned by ``install(application)`` before startup. This explicit bridge requires the real production components and a source registry; it is not an adapter-first component factory. Production tasks use their existing indexed leases. Pools and providers stay borrowed unless explicitly transferred in ``owned_resources``. """ from adcp.reporting.production.service import ReportingProductionSupport if ( type(production) is not ReportingProductionSupport or production.source_registry is None or production._task is not None or production._closed or production._service_lifecycle is not None ): raise ReliableReportingConfigurationError( "requires an unstarted registered production graph" ) def unavailable(_: ReportingConfiguration) -> ReportingAccountContext: raise ReliableReportingConfigurationError( "production requires typed sync_accounts admission" ) service = cls( store=production.store, account_context=unavailable, owned_resources=(*owned_resources, ReportingServiceResource(close=production.aclose)), ) service._production = production production._service_lifecycle = service._lifecycle return service @classmethod def memory( cls, *, account_context: ReportingContextResolver, caller_resolver: ReportingCallerResolver | None = None, clock: Callable[[], datetime] | None = None, production: ReportingProductionOptions | None = None, **kwargs: Any, ) -> ReliableReportingService: """Create a deterministic in-memory service for tests and pilots.""" from adcp.reporting.production.memory import InMemoryReportingProductionStore store = ( InMemoryReportingProductionStore(clock=clock, notifications=production.notifications) if production is not None else InMemoryReportingLedgerStore(clock=clock) ) return cls( store=store, account_context=account_context, caller_resolver=caller_resolver, clock=clock, production=production, **kwargs, ) @classmethod def postgres( cls, *, pool: Any, account_context: ReportingContextResolver, caller_resolver: ReportingCallerResolver | None = None, clock: Callable[[], datetime] | None = None, production: ReportingProductionOptions | None = None, **kwargs: Any, ) -> ReliableReportingService: """Create durable ledger, staging and seal stores over a borrowed pool. ``production`` additionally composes the managed pipeline, receipts and optional signed notifications when ``install`` freezes registered sources. Startup installs all SDK-owned schemas. Explicit source-store overrides remain the adopter's responsibility. No pool is opened or closed here. """ from adcp.reporting.inline_storage import PgReportingSealStore, PgReportingStagingStore from adcp.reporting.ledger.pg import PgReportingLedgerStore from adcp.reporting.production.pg import PgReportingProductionStore store = ( PgReportingProductionStore( pool=pool, clock=clock, notifications=production.notifications ) if production is not None else PgReportingLedgerStore(pool=pool, clock=clock) ) service = cls( store=store, account_context=account_context, caller_resolver=caller_resolver, clock=clock, production=production, **kwargs, ) staging = PgReportingStagingStore(pool=pool) service.sources = ReportingAdapterRegistry( clock=service._clock, staging=staging, seals=PgReportingSealStore(pool=pool) ) service._initialize_source_storage = staging.create_schema return service def _compose_production(self) -> None: if self._production_options is None or self._production is not None: return from adcp.reporting.production.memory import InMemoryReportingProductionStore from adcp.reporting.production.pg import PgReportingProductionStore if not isinstance( self.store, (InMemoryReportingProductionStore, PgReportingProductionStore) ): raise ReliableReportingConfigurationError( "production options require a memory or postgres production store" ) self._production = self._production_options._compose( store=self.store, sources=self.sources, account_context=self._context_resolver, clock=self._clock, escalation=self._escalation, consumer_status_enabled=self._consumer_status_enabled, ) self._production._service_lifecycle = self._lifecycle self.sources.freeze() async def _close_production(self) -> None: if self._production is not None: await self._production.aclose() async def configure(self, configuration: ReportingConfiguration) -> None: """Resolve trusted account facts once and freeze this generation's route.""" if self._production is not None or self._production_options is not None: raise ReliableReportingConfigurationError( "production requires typed sync_accounts admission" ) async def configure() -> None: async with self._configuration_lock: await self._configure(configuration) await self._lifecycle.call(configure, before_start=True) async def _configure( self, configuration: ReportingConfiguration, *, persist: bool = True ) -> None: key = configuration.generation_key existing = self._bindings.get(key) if existing is not None: if existing.configuration != configuration: raise ReliableReportingConfigurationError( f"configuration {key.delivery_config_id}@{key.delivery_config_version} " f"for account {key.account_id!r} is already frozen with other facts" ) return context = await _resolve(self._context_resolver(configuration)) if not isinstance(context, ReportingAccountContext): raise TypeError("account_context must resolve ReportingAccountContext") if context.account_id != configuration.account_id: raise ReliableReportingConfigurationError( "resolved account context does not match the configuration account" ) if context.account_timezone != configuration.account_timezone: raise ReliableReportingConfigurationError( "resolved account timezone does not match the configuration timezone" ) registration = self.sources.get(context.adapter) self._validate_offerings(configuration, context, registration) producer = self._producer_factory( source=registration.executor, offerings=context.producer_offerings(), store=self.store, object_reader=registration.object_reader, escalation=self._escalation, worker_id=f"reporting-service:{context.adapter}", clock=self._clock, ) binding = _Binding(configuration=configuration, context=context, producer=producer) if persist: if self._initialized: await self.store.put_configuration(configuration) else: self._pending_configurations[key] = configuration self._bindings[key] = binding def _validate_offerings( self, configuration: ReportingConfiguration, context: ReportingAccountContext, registration: AdapterRegistration, ) -> None: capabilities = registration.executor.capabilities if capabilities.scope != "effective_account": raise ReliableReportingConfigurationError( "registered source capabilities must be scoped to an effective account" ) if capabilities.source_scope != _thaw(context.source_scope): raise ReliableReportingConfigurationError( "resolved source_scope does not match the registered adapter capabilities" ) selected: list[ProvisionalSnapshotOfferingV1 | AuthoritativeOfferingV1] = [] if context.snapshot_offering_id is not None: try: snapshot = capabilities.offering(context.snapshot_offering_id) except KeyError as error: raise ReliableReportingConfigurationError( f"snapshot offering {context.snapshot_offering_id!r} is not registered" ) from error if not isinstance(snapshot, ProvisionalSnapshotOfferingV1): raise ReliableReportingConfigurationError( "snapshot_offering_id must name a provisional snapshot offering" ) selected.append(snapshot) if context.official_offering_id is not None: try: official = capabilities.offering(context.official_offering_id) except KeyError as error: raise ReliableReportingConfigurationError( f"official offering {context.official_offering_id!r} is not registered" ) from error if not isinstance(official, AuthoritativeOfferingV1): raise ReliableReportingConfigurationError( "official_offering_id must name an authoritative offering" ) selected.append(official) for offering in selected: if ( offering.contract.report_definition_id != configuration.report_definition_id or offering.contract.reporting_profile != configuration.reporting_profile ): raise ReliableReportingConfigurationError( "source offering contract does not match the reporting configuration" ) required = ( context.official_offering_id if configuration.required_finality == "official" else context.snapshot_offering_id ) if required is None: raise ReliableReportingConfigurationError( f"configuration requires {configuration.required_finality} but its account " "context declares no matching source offering" ) declared_id = context.capability_offering.get("offering_id") if not declared_id: raise ReliableReportingConfigurationError( "capability_offering must contain the advertised offering_id" ) if context.capability_offering.get("feed_purpose") != configuration.feed_purpose: raise ReliableReportingConfigurationError( "advertised capability offering feed_purpose does not match the configuration" ) if ( context.capability_offering.get("report_definition_id") != configuration.report_definition_id ): raise ReliableReportingConfigurationError( "advertised capability offering report definition does not match the configuration" ) supported_finality = context.capability_offering.get("supported_finality", ()) if configuration.required_finality not in supported_finality: raise ReliableReportingConfigurationError( "advertised capability offering does not support the required finality" ) if registration.capability_offerings: public = ReportingDeliveryOffering.model_validate( _thaw(context.capability_offering) ).model_dump(mode="json", exclude_none=True) if not any(_thaw(item) == public for item in registration.capability_offerings): raise ReliableReportingConfigurationError( "account capability offering must match a registered public offering" ) def validate(self) -> None: """Fail fast on tier combinations the configured components cannot honor.""" self._compose_production() if self._production is not None and self._production_options is None and self.sources.names: raise ReliableReportingConfigurationError( "production adapters must belong to its frozen source registry" ) if self._worker_interval is not None and self._worker_interval <= timedelta(0): raise ReliableReportingConfigurationError("worker_interval must be greater than zero") if self._notification_worker is not None and self._notification_attempt_store is None: raise ReliableReportingConfigurationError( "a notification worker requires durable notification attempt storage" ) if self._receipt_handler is not None and self._materialization_worker is None: raise ReliableReportingConfigurationError( "receipt handling requires managed delivery/materialization" ) if self._reconciled_billing and ( self._receipt_handler is None or self._materialization_worker is None ): raise ReliableReportingConfigurationError( "reconciled billing requires managed delivery and a receipt handler" ) if self._production is None: for name in self.sources.names: for offering in self.sources.get(name).capability_offerings: if offering.get("method") and self._materialization_worker is None: raise ReliableReportingConfigurationError( "managed capability offerings require a materialization worker" ) if ( offering["reconciliation_mode"] == "consumer_receipt" and not self._reconciled_billing ): raise ReliableReportingConfigurationError( "consumer_receipt offerings require reconciled billing" ) async def initialize(self) -> None: """Prepare resources once; start/run_worker opens reporting admission.""" await self._lifecycle.initialize() async def _initialize(self) -> None: self.validate() async with self._configuration_lock: await self.store.create_schema() if self._initialize_source_storage is not None: await self._initialize_source_storage() for configuration in self._pending_configurations.values(): await self.store.put_configuration(configuration) self._pending_configurations.clear() if self._production is None and self.sources.names: # Older custom stores can retain explicit configure() startup. # Built-in stores recover all generations in one scan, without # rewriting their activation, retirement or producer progress. enumerate_configurations = getattr(self.store, "list_all_configurations", None) if callable(enumerate_configurations): for configuration in await enumerate_configurations(): await self._configure(configuration, persist=False) self.sources.freeze() if self._production is not None: await self._production.start() self._initialized = True async def start(self) -> None: await self.initialize() await self._lifecycle.activate( self._monitor_production if self._production is not None else (self._worker_loop if self._worker_interval is not None else None) ) async def _monitor_production(self) -> None: assert self._production is not None while not self._lifecycle.stopping: self._production._assert_components() stop = asyncio.create_task(self._lifecycle.wait_for_stop(60)) workers = [ task for task in (self._production._task, self._production._notification_task) if task is not None ] try: await asyncio.wait([stop, *workers], return_when=asyncio.FIRST_COMPLETED) finally: stop.cancel() await asyncio.gather(stop, return_exceptions=True) @property def state(self) -> ReliableReportingState: """Current process lifecycle; STOPPING retains all unsettled ownership.""" return self._lifecycle.state @property def ready(self) -> bool: """Whether this process admits work (not a durable-tier graph proof).""" return self._lifecycle.ready @property def failure(self) -> ReliableReportingServiceError | None: """Sanitized first terminal lifecycle failure, if any.""" return self._lifecycle.failure async def close(self, *, timeout: float | None = None) -> None: """Reject new work, drain admitted operations, then close owned resources. A timeout or cancelled waiter leaves the shared shutdown running and the service STOPPING until settlement. Injected components are borrowed; transfer lifetime explicitly with ``owned_resources`` to close them here. """ await self._lifecycle.close(timeout=timeout) async def wait(self) -> None: """Wait for shutdown and raise a typed failure if supervision failed.""" await self._lifecycle.wait() async def __aenter__(self) -> ReliableReportingService: await self.start() return self async def __aexit__(self, *_exc: object) -> None: await self.close() async def _worker_loop(self) -> None: assert self._worker_interval is not None while not self._lifecycle.stopping: try: await self.run_worker() except Exception as error: # Configuration and extension failures are isolated inside a # turn. Only a failure of the scheduler itself stops service. if not self._lifecycle.stopping: await self._lifecycle.fail("service") await self._report_worker_error("service", error) return await self._lifecycle.wait_for_stop(self._worker_interval.total_seconds()) async def _report_worker_error(self, component: str, error: BaseException) -> None: # Preserve the adopter callback's detailed identity and original error, # without writing account identifiers or provider bodies to SDK logs. logger.error("Reliable Reporting worker component %s failed", component.split(":", 1)[0]) if self._worker_error_handler is None: return try: await _resolve(self._worker_error_handler(component, error)) except Exception: logger.error("Reliable Reporting worker error handler failed") async def run_worker(self, *, now: datetime | None = None) -> ReliableReportingTurn: """Route one turn across every frozen configuration generation.""" if self._production is not None or self._production_options is not None: raise ReliableReportingConfigurationError("production workers are owned by start/close") await self.initialize() await self._lifecycle.activate() return await self._lifecycle.call(lambda: self._run_worker(now=now)) async def _run_worker(self, *, now: datetime | None) -> ReliableReportingTurn: async with self._turn_lock: # A turn admitted before STOPPING may have waited behind another # turn. It must not start fresh work after that turn has drained. self._lifecycle.require_ready() turn = ReliableReportingTurn() with source_turn(): for key, binding in sorted( self._bindings.items(), key=lambda item: ( item[0].account_id, item[0].delivery_config_id, item[0].delivery_config_version, ), ): if self._lifecycle.stopping: break try: turn.configurations[key] = await binding.producer.run_configuration( binding.configuration, now=now ) except Exception as error: turn.configuration_errors[key] = error await self._report_worker_error( "configuration:" f"{key.account_id}:{key.delivery_config_id}@" f"{key.delivery_config_version}", error, ) extensions = ( ("materialization", self._materialization_worker), ("notification", self._notification_worker), ) for name, extension in extensions: if self._lifecycle.stopping: break if extension is not None: try: result = await extension.run_once() except Exception as error: turn.extension_errors[name] = error await self._report_worker_error(name, error) continue if result is not None: turn.extension_results.append(result) return turn async def caller_for(self, request: Any, context: Any | None) -> ReportingStatusCaller: caller = await _resolve(self._caller_resolver(request, context)) if not isinstance(caller, ReportingStatusCaller): raise TypeError("caller_resolver must resolve ReportingStatusCaller") return caller @staticmethod def _default_caller(request: Any, context: Any | None) -> ReportingStatusCaller: del request # RequestContext is duck-typed here: importing decisioning.context at # module scope would cycle through the decisioning reporting helpers. # Its caller_identity is a per-account cache key, not a buyer identity. if hasattr(context, "account_id"): # AccountAwareToolContext.account is an opaque adopter object; # resolve_account_into_context supplies its stable id separately. account_id = getattr(context, "account_id", None) consumer_id = getattr(context, "caller_identity", None) else: from adcp.server.auth import current_principal account_id = getattr(getattr(context, "account", None), "id", None) consumer_id = current_principal.get() if consumer_id is None: # Signed requests carry their verified identity on AuthInfo, # without populating the bearer middleware's ContextVar. auth_principal = getattr(context, "auth_principal", None) auth_info = getattr(context, "auth_info", None) if auth_principal == getattr(auth_info, "principal", None): consumer_id = auth_principal if not account_id or not consumer_id: raise ReliableReportingConfigurationError( "default caller resolution requires AccountAwareToolContext.account_id and " "authenticated caller_identity, or RequestContext.account.id and an " "authenticated current_principal or verified auth_principal; " "provide caller_resolver for another trust model", kind="auth_required" if not consumer_id else "misconfiguration", ) return ReportingStatusCaller(account_id=account_id, consumer_id=consumer_id) async def get_reporting_status( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.get_reporting_status(request, context)) return await self._lifecycle.call(lambda: self._get_reporting_status(request, context)) async def _get_reporting_status(self, request: Any, context: Any | None) -> dict[str, Any]: caller = await self.caller_for(request, context) return await ReportingStatusHandler( self.store, consumer_status_enabled=self._consumer_status_enabled, escalation=self._escalation, ).handle(_wire(request), caller=caller) async def sync_reporting_status( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.sync_reporting_status(request, context)) return await self._lifecycle.call(lambda: self._sync_reporting_status(request, context)) async def _sync_reporting_status(self, request: Any, context: Any | None) -> dict[str, Any]: if not self._consumer_status_enabled: raise ReliableReportingConfigurationError("consumer status ingest is disabled") caller = await self.caller_for(request, context) return await ConsumerStatusIngest(self.store, enabled=True, clock=self._clock).handle( _wire(request), account_id=caller.account_id, consumer_id=caller.consumer_id, ) async def sync_reporting_receipts( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.sync_reporting_receipts(request, context)) return await self._lifecycle.call(lambda: self._sync_reporting_receipts(request, context)) async def _sync_reporting_receipts(self, request: Any, context: Any | None) -> dict[str, Any]: if self._receipt_handler is None: raise ReliableReportingConfigurationError("reporting receipt handling is disabled") caller = await self.caller_for(request, context) return await self._receipt_handler.handle( _wire(request), account_id=caller.account_id, consumer_id=caller.consumer_id ) async def get_revision_content( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.get_media_buy_delivery(request, context)) return await self._lifecycle.call(lambda: self._get_revision_content(request, context)) async def _get_revision_content(self, request: Any, context: Any | None) -> dict[str, Any]: payload = _wire(request) revision_id = payload.get("reporting_revision_id") if not revision_id: raise LedgerConflictError( "MISSING_REVISION_ID", "an exact revision read requires reporting_revision_id" ) caller = await self.caller_for(request, context) cursor = (payload.get("pagination") or {}).get("cursor") limit = int((payload.get("pagination") or {}).get("limit") or 500) if limit < 1 or limit > 1000: raise LedgerConflictError( "INVALID_PAGE_SIZE", "pagination.limit must be between 1 and 1000" ) if cursor: selected = decode_cursor(cursor).get("revision") if selected != revision_id: raise LedgerConflictError( "CURSOR_REVISION_MISMATCH", "this cursor belongs to a different reporting revision", ) revision_view = await ReportingStatusHandler(self.store).handle( {"view": "revision", "reporting_revision_id": revision_id}, caller=caller, ) revision = revision_view["revision"] page = await self.store.read_revision_rows( account_id=caller.account_id, reporting_revision_id=revision_id, cursor=cursor, limit=limit, ) period = revision["period"] response: dict[str, Any] = { "status": "completed", "reporting_period": {"start": period["start"], "end": period["end"]}, "media_buy_deliveries": [], "reporting_revision_binding": { "reporting_revision_id": revision_id, "row_count": revision["row_count"], "control_totals": revision["control_totals"], "content_sha256": revision["revision_content_sha256"], }, "reporting_revision": revision, "reporting_rows": [dict(row) for row in page.rows], "pagination": { "total_count": page.total_count, "has_more": page.has_more, **({"cursor": page.cursor} if page.cursor else {}), }, } if "context" in payload: response["context"] = payload["context"] return response def capability_block(self) -> dict[str, Any]: """Project only components and offerings this service actually installed.""" if self._production is not None: return dict(self._production._declared_capabilities()["reporting_delivery"]) if self._production_options is not None: return self._production_options._capability_block( store=self.store, sources=self.sources, escalation=self._escalation, consumer_status_enabled=self._consumer_status_enabled, ) offerings: dict[str, dict[str, Any]] = {} declarations = [ item for name in self.sources.names for item in self.sources.get(name).capability_offerings ] # Preserve pre-registration callers that only supplied their offering # with each account context. New declarations are ready at first boot. declarations.extend( binding.context.capability_offering for binding in self._bindings.values() ) for declaration in declarations: offering = ReportingDeliveryOffering.model_validate(_thaw(declaration)).model_dump( mode="json", exclude_none=True ) offering_id = str(offering["offering_id"]) existing = offerings.get(offering_id) if existing is not None and existing != offering: raise ReliableReportingConfigurationError( f"capability offering {offering_id!r} has conflicting declarations" ) offerings[offering_id] = offering if not offerings: return {} configurations = [binding.configuration for binding in self._bindings.values()] payload = _advertised_reporting_delivery( escalation=self._escalation, consumer_status_task=self._consumer_status_enabled, offerings=[offerings[key] for key in sorted(offerings)], automated_recovery_window=max( (item.automated_recovery_window for item in configurations), default=self._automated_recovery_window, ), status_retention_days=min( (item.status_retention_days for item in configurations), default=self._status_retention_days, ), ) # The service owns these installed-component declarations; they are not # caller-provided producer extensions. payload["managed_delivery"] = self._materialization_worker is not None payload["reconciled_billing"] = self._reconciled_billing if self._reconciled_billing: payload["receipt_task"] = "sync_reporting_receipts" return payload def inject_capabilities(self, response: Any) -> dict[str, Any]: """Merge the truthful reporting block into a base capability response.""" payload = _wire(response) block = self.capability_block() if not block: return payload media_buy = dict(payload.get("media_buy") or {}) media_buy["reporting_delivery"] = block payload["media_buy"] = media_buy protocols = [str(item) for item in payload.get("supported_protocols") or []] if "media_buy" not in protocols: protocols.append("media_buy") payload["supported_protocols"] = protocols experimental = [str(item) for item in payload.get("experimental_features") or []] if "media_buy.reporting_delivery" not in experimental: experimental.append("media_buy.reporting_delivery") payload["experimental_features"] = experimental return payload def install(self, platform: Any) -> Any: """Install reporting on an ADCPHandler or a DecisioningPlatform. Core decisioning installations retain the platform's account resolution, dispatcher and boot-time idempotency checks. Pass the returned platform to ``adcp.decisioning.serve`` as usual. Managed production graphs still bind an ADCPHandler and return their exact production handler; use ``create_adcp_server_from_platform`` first for that composition. """ self._compose_production() if self._production is not None: self._production.handler.bind_application(platform) return self._production.handler from adcp.server import ADCPHandler from adcp.server.mcp_tools import get_tools_for_handler decisioning = not isinstance(platform, ADCPHandler) if decisioning: from adcp.decisioning.platform import DecisioningPlatform if not isinstance(platform, DecisioningPlatform): raise ReliableReportingConfigurationError( "install() requires an ADCPHandler or DecisioningPlatform instance" ) if getattr(platform, "_reliable_reporting_service", None) is self: return platform if getattr(platform, "_reliable_reporting_service", None) is not None: raise ReliableReportingConfigurationError( "this platform already has another ReliableReportingService installed" ) original = type(platform) existing_tools = ( set() if decisioning else {item["name"] for item in get_tools_for_handler(platform, _include_schemas=False)} ) tools = set(existing_tools) tools.update(self.reporting_tools) service = self class ReportingInstallMixin: async def get_reporting_status(self, params: Any, context: Any | None = None) -> Any: return await installed_reporting_call( "get_reporting_status", lambda: service.get_reporting_status(params, context), decisioning=decisioning, ) async def sync_reporting_status(self, params: Any, context: Any | None = None) -> Any: if not service._consumer_status_enabled: return await _resolve( getattr(super(), "sync_reporting_status")(params, context) ) return await installed_reporting_call( "sync_reporting_status", lambda: service.sync_reporting_status(params, context), decisioning=decisioning, ) async def sync_reporting_receipts(self, params: Any, context: Any | None = None) -> Any: if service._receipt_handler is None: return await _resolve( getattr(super(), "sync_reporting_receipts")(params, context) ) return await installed_reporting_call( "sync_reporting_receipts", lambda: service.sync_reporting_receipts(params, context), decisioning=decisioning, ) async def _get_reporting_revision_content( self, params: Any, context: Any | None = None ) -> Any: return await installed_reporting_call( "get_media_buy_delivery", lambda: service.get_revision_content(params, context), decisioning=decisioning, ) attrs: dict[str, Any] = {"__module__": original.__module__} if not decisioning: async def get_adcp_capabilities( instance: Any, params: Any, context: Any | None = None ) -> Any: result = await _resolve(original.get_adcp_capabilities(instance, params, context)) return service.inject_capabilities(result) async def get_media_buy_delivery( instance: Any, params: Any, context: Any | None = None ) -> Any: if not _wire(params).get("reporting_revision_id"): return await _resolve( original.get_media_buy_delivery(instance, params, context) ) return await installed_reporting_call( "get_media_buy_delivery", lambda: service.get_revision_content(params, context), decisioning=decisioning, ) attrs.update( advertised_tools=tools, get_adcp_capabilities=get_adcp_capabilities, get_media_buy_delivery=get_media_buy_delivery, ) reporting_methods = { "get_adcp_capabilities", "get_reporting_status", "sync_reporting_status", "sync_reporting_receipts", "get_media_buy_delivery", } for name in existing_tools: if name in reporting_methods: continue method = getattr(original, name, None) if method is None: continue def passthrough( method_name: str, original_method: Callable[..., Any] ) -> Callable[..., Any]: async def delegated(instance: Any, *args: Any, **kwargs: Any) -> Any: return await _resolve(original_method(instance, *args, **kwargs)) delegated.__name__ = method_name delegated.__qualname__ = f"ReliableReporting{original.__name__}.{method_name}" return delegated attrs[name] = passthrough(name, method) installed_name = f"ReliableReporting{original.__name__}_{id(self):x}" installed = type( installed_name, (ReportingInstallMixin, original), attrs, ) try: platform.__class__ = installed except TypeError: # CPython rejects some otherwise ordinary multiple-inheritance # layout changes. The documented API returns the installed # platform, so preserve the instance state in a replacement. state = getattr(platform, "__dict__", None) if state is None: raise ReliableReportingConfigurationError( "install() requires an ordinary mutable platform instance; slot-only/native " "instances should wire the service handler methods explicitly" ) from None replacement: Any = object.__new__(installed) replacement.__dict__.update(state) platform = replacement setattr(platform, "_reliable_reporting_service", self) return platform @property def reporting_tools(self) -> frozenset[str]: """Core task implementations mounted by ``install`` (independent of readiness).""" tasks = {"get_reporting_status", "get_media_buy_delivery"} if self._consumer_status_enabled: tasks.add("sync_reporting_status") if self._receipt_handler is not None: tasks.add("sync_reporting_receipts") return frozenset(tasks)High-level owner of adapters, ledger handlers, workers, and lifecycle.
Admitted calls own a task so transport cancellation cannot interrupt their cleanup. Invoke them outside caller-owned ledger transactions; use the low-level store API when composing an ambient transaction batch. Injected resources are borrowed unless explicitly transferred via
owned_resources.Static methods
def from_production(production: ReportingProductionSupport,
*,
owned_resources: Sequence[ReportingServiceResource] = ())-
Own a preassembled B2 graph with durable fixed-profile admission.
Mount the exact handler returned by
install(application)before startup. This explicit bridge requires the real production components and a source registry; it is not an adapter-first component factory. Production tasks use their existing indexed leases. Pools and providers stay borrowed unless explicitly transferred inowned_resources. def memory(*,
account_context: ReportingContextResolver,
caller_resolver: ReportingCallerResolver | None = None,
clock: Callable[[], datetime] | None = None,
production: ReportingProductionOptions | None = None,
**kwargs: Any)-
Create a deterministic in-memory service for tests and pilots.
def postgres(*,
pool: Any,
account_context: ReportingContextResolver,
caller_resolver: ReportingCallerResolver | None = None,
clock: Callable[[], datetime] | None = None,
production: ReportingProductionOptions | None = None,
**kwargs: Any)-
Create durable ledger, staging and seal stores over a borrowed pool.
adcp.reporting.productionadditionally composes the managed pipeline, receipts and optional signed notifications wheninstallfreezes registered sources. Startup installs all SDK-owned schemas. Explicit source-store overrides remain the adopter's responsibility. No pool is opened or closed here.
Instance variables
prop failure : ReliableReportingServiceError | None-
Expand source code
@property def failure(self) -> ReliableReportingServiceError | None: """Sanitized first terminal lifecycle failure, if any.""" return self._lifecycle.failureSanitized first terminal lifecycle failure, if any.
prop ready : bool-
Expand source code
@property def ready(self) -> bool: """Whether this process admits work (not a durable-tier graph proof).""" return self._lifecycle.readyWhether this process admits work (not a durable-tier graph proof).
prop reporting_tools : frozenset[str]-
Expand source code
@property def reporting_tools(self) -> frozenset[str]: """Core task implementations mounted by ``install`` (independent of readiness).""" tasks = {"get_reporting_status", "get_media_buy_delivery"} if self._consumer_status_enabled: tasks.add("sync_reporting_status") if self._receipt_handler is not None: tasks.add("sync_reporting_receipts") return frozenset(tasks)Core task implementations mounted by
install(independent of readiness). prop state : ReliableReportingState-
Expand source code
@property def state(self) -> ReliableReportingState: """Current process lifecycle; STOPPING retains all unsettled ownership.""" return self._lifecycle.stateCurrent process lifecycle; STOPPING retains all unsettled ownership.
Methods
async def caller_for(self, request: Any, context: Any | None) ‑> ReportingStatusCaller-
Expand source code
async def caller_for(self, request: Any, context: Any | None) -> ReportingStatusCaller: caller = await _resolve(self._caller_resolver(request, context)) if not isinstance(caller, ReportingStatusCaller): raise TypeError("caller_resolver must resolve ReportingStatusCaller") return caller def capability_block(self) ‑> dict[str, typing.Any]-
Expand source code
def capability_block(self) -> dict[str, Any]: """Project only components and offerings this service actually installed.""" if self._production is not None: return dict(self._production._declared_capabilities()["reporting_delivery"]) if self._production_options is not None: return self._production_options._capability_block( store=self.store, sources=self.sources, escalation=self._escalation, consumer_status_enabled=self._consumer_status_enabled, ) offerings: dict[str, dict[str, Any]] = {} declarations = [ item for name in self.sources.names for item in self.sources.get(name).capability_offerings ] # Preserve pre-registration callers that only supplied their offering # with each account context. New declarations are ready at first boot. declarations.extend( binding.context.capability_offering for binding in self._bindings.values() ) for declaration in declarations: offering = ReportingDeliveryOffering.model_validate(_thaw(declaration)).model_dump( mode="json", exclude_none=True ) offering_id = str(offering["offering_id"]) existing = offerings.get(offering_id) if existing is not None and existing != offering: raise ReliableReportingConfigurationError( f"capability offering {offering_id!r} has conflicting declarations" ) offerings[offering_id] = offering if not offerings: return {} configurations = [binding.configuration for binding in self._bindings.values()] payload = _advertised_reporting_delivery( escalation=self._escalation, consumer_status_task=self._consumer_status_enabled, offerings=[offerings[key] for key in sorted(offerings)], automated_recovery_window=max( (item.automated_recovery_window for item in configurations), default=self._automated_recovery_window, ), status_retention_days=min( (item.status_retention_days for item in configurations), default=self._status_retention_days, ), ) # The service owns these installed-component declarations; they are not # caller-provided producer extensions. payload["managed_delivery"] = self._materialization_worker is not None payload["reconciled_billing"] = self._reconciled_billing if self._reconciled_billing: payload["receipt_task"] = "sync_reporting_receipts" return payloadProject only components and offerings this service actually installed.
async def close(self, *, timeout: float | None = None) ‑> None-
Expand source code
async def close(self, *, timeout: float | None = None) -> None: """Reject new work, drain admitted operations, then close owned resources. A timeout or cancelled waiter leaves the shared shutdown running and the service STOPPING until settlement. Injected components are borrowed; transfer lifetime explicitly with ``owned_resources`` to close them here. """ await self._lifecycle.close(timeout=timeout)Reject new work, drain admitted operations, then close owned resources.
A timeout or cancelled waiter leaves the shared shutdown running and the service STOPPING until settlement. Injected components are borrowed; transfer lifetime explicitly with
owned_resourcesto close them here. async def configure(self, configuration: ReportingConfiguration) ‑> None-
Expand source code
async def configure(self, configuration: ReportingConfiguration) -> None: """Resolve trusted account facts once and freeze this generation's route.""" if self._production is not None or self._production_options is not None: raise ReliableReportingConfigurationError( "production requires typed sync_accounts admission" ) async def configure() -> None: async with self._configuration_lock: await self._configure(configuration) await self._lifecycle.call(configure, before_start=True)Resolve trusted account facts once and freeze this generation's route.
async def get_reporting_status(self, request: Any, context: Any | None = None) ‑> dict[str, typing.Any]-
Expand source code
async def get_reporting_status( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.get_reporting_status(request, context)) return await self._lifecycle.call(lambda: self._get_reporting_status(request, context)) async def get_revision_content(self, request: Any, context: Any | None = None) ‑> dict[str, typing.Any]-
Expand source code
async def get_revision_content( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.get_media_buy_delivery(request, context)) return await self._lifecycle.call(lambda: self._get_revision_content(request, context)) async def initialize(self) ‑> None-
Expand source code
async def initialize(self) -> None: """Prepare resources once; start/run_worker opens reporting admission.""" await self._lifecycle.initialize()Prepare resources once; start/run_worker opens reporting admission.
def inject_capabilities(self, response: Any) ‑> dict[str, typing.Any]-
Expand source code
def inject_capabilities(self, response: Any) -> dict[str, Any]: """Merge the truthful reporting block into a base capability response.""" payload = _wire(response) block = self.capability_block() if not block: return payload media_buy = dict(payload.get("media_buy") or {}) media_buy["reporting_delivery"] = block payload["media_buy"] = media_buy protocols = [str(item) for item in payload.get("supported_protocols") or []] if "media_buy" not in protocols: protocols.append("media_buy") payload["supported_protocols"] = protocols experimental = [str(item) for item in payload.get("experimental_features") or []] if "media_buy.reporting_delivery" not in experimental: experimental.append("media_buy.reporting_delivery") payload["experimental_features"] = experimental return payloadMerge the truthful reporting block into a base capability response.
def install(self, platform: Any) ‑> Any-
Expand source code
def install(self, platform: Any) -> Any: """Install reporting on an ADCPHandler or a DecisioningPlatform. Core decisioning installations retain the platform's account resolution, dispatcher and boot-time idempotency checks. Pass the returned platform to ``adcp.decisioning.serve`` as usual. Managed production graphs still bind an ADCPHandler and return their exact production handler; use ``create_adcp_server_from_platform`` first for that composition. """ self._compose_production() if self._production is not None: self._production.handler.bind_application(platform) return self._production.handler from adcp.server import ADCPHandler from adcp.server.mcp_tools import get_tools_for_handler decisioning = not isinstance(platform, ADCPHandler) if decisioning: from adcp.decisioning.platform import DecisioningPlatform if not isinstance(platform, DecisioningPlatform): raise ReliableReportingConfigurationError( "install() requires an ADCPHandler or DecisioningPlatform instance" ) if getattr(platform, "_reliable_reporting_service", None) is self: return platform if getattr(platform, "_reliable_reporting_service", None) is not None: raise ReliableReportingConfigurationError( "this platform already has another ReliableReportingService installed" ) original = type(platform) existing_tools = ( set() if decisioning else {item["name"] for item in get_tools_for_handler(platform, _include_schemas=False)} ) tools = set(existing_tools) tools.update(self.reporting_tools) service = self class ReportingInstallMixin: async def get_reporting_status(self, params: Any, context: Any | None = None) -> Any: return await installed_reporting_call( "get_reporting_status", lambda: service.get_reporting_status(params, context), decisioning=decisioning, ) async def sync_reporting_status(self, params: Any, context: Any | None = None) -> Any: if not service._consumer_status_enabled: return await _resolve( getattr(super(), "sync_reporting_status")(params, context) ) return await installed_reporting_call( "sync_reporting_status", lambda: service.sync_reporting_status(params, context), decisioning=decisioning, ) async def sync_reporting_receipts(self, params: Any, context: Any | None = None) -> Any: if service._receipt_handler is None: return await _resolve( getattr(super(), "sync_reporting_receipts")(params, context) ) return await installed_reporting_call( "sync_reporting_receipts", lambda: service.sync_reporting_receipts(params, context), decisioning=decisioning, ) async def _get_reporting_revision_content( self, params: Any, context: Any | None = None ) -> Any: return await installed_reporting_call( "get_media_buy_delivery", lambda: service.get_revision_content(params, context), decisioning=decisioning, ) attrs: dict[str, Any] = {"__module__": original.__module__} if not decisioning: async def get_adcp_capabilities( instance: Any, params: Any, context: Any | None = None ) -> Any: result = await _resolve(original.get_adcp_capabilities(instance, params, context)) return service.inject_capabilities(result) async def get_media_buy_delivery( instance: Any, params: Any, context: Any | None = None ) -> Any: if not _wire(params).get("reporting_revision_id"): return await _resolve( original.get_media_buy_delivery(instance, params, context) ) return await installed_reporting_call( "get_media_buy_delivery", lambda: service.get_revision_content(params, context), decisioning=decisioning, ) attrs.update( advertised_tools=tools, get_adcp_capabilities=get_adcp_capabilities, get_media_buy_delivery=get_media_buy_delivery, ) reporting_methods = { "get_adcp_capabilities", "get_reporting_status", "sync_reporting_status", "sync_reporting_receipts", "get_media_buy_delivery", } for name in existing_tools: if name in reporting_methods: continue method = getattr(original, name, None) if method is None: continue def passthrough( method_name: str, original_method: Callable[..., Any] ) -> Callable[..., Any]: async def delegated(instance: Any, *args: Any, **kwargs: Any) -> Any: return await _resolve(original_method(instance, *args, **kwargs)) delegated.__name__ = method_name delegated.__qualname__ = f"ReliableReporting{original.__name__}.{method_name}" return delegated attrs[name] = passthrough(name, method) installed_name = f"ReliableReporting{original.__name__}_{id(self):x}" installed = type( installed_name, (ReportingInstallMixin, original), attrs, ) try: platform.__class__ = installed except TypeError: # CPython rejects some otherwise ordinary multiple-inheritance # layout changes. The documented API returns the installed # platform, so preserve the instance state in a replacement. state = getattr(platform, "__dict__", None) if state is None: raise ReliableReportingConfigurationError( "install() requires an ordinary mutable platform instance; slot-only/native " "instances should wire the service handler methods explicitly" ) from None replacement: Any = object.__new__(installed) replacement.__dict__.update(state) platform = replacement setattr(platform, "_reliable_reporting_service", self) return platformInstall reporting on an ADCPHandler or a DecisioningPlatform.
Core decisioning installations retain the platform's account resolution, dispatcher and boot-time idempotency checks. Pass the returned platform to
serve()as usual. Managed production graphs still bind an ADCPHandler and return their exact production handler; usecreate_adcp_server_from_platformfirst for that composition. async def run_worker(self, *, now: datetime | None = None) ‑> ReliableReportingTurn-
Expand source code
async def run_worker(self, *, now: datetime | None = None) -> ReliableReportingTurn: """Route one turn across every frozen configuration generation.""" if self._production is not None or self._production_options is not None: raise ReliableReportingConfigurationError("production workers are owned by start/close") await self.initialize() await self._lifecycle.activate() return await self._lifecycle.call(lambda: self._run_worker(now=now))Route one turn across every frozen configuration generation.
async def start(self) ‑> None-
Expand source code
async def start(self) -> None: await self.initialize() await self._lifecycle.activate( self._monitor_production if self._production is not None else (self._worker_loop if self._worker_interval is not None else None) ) async def sync_reporting_receipts(self, request: Any, context: Any | None = None) ‑> dict[str, typing.Any]-
Expand source code
async def sync_reporting_receipts( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.sync_reporting_receipts(request, context)) return await self._lifecycle.call(lambda: self._sync_reporting_receipts(request, context)) async def sync_reporting_status(self, request: Any, context: Any | None = None) ‑> dict[str, typing.Any]-
Expand source code
async def sync_reporting_status( self, request: Any, context: Any | None = None ) -> dict[str, Any]: if self._production is not None: return _wire(await self._production.handler.sync_reporting_status(request, context)) return await self._lifecycle.call(lambda: self._sync_reporting_status(request, context)) def validate(self) ‑> None-
Expand source code
def validate(self) -> None: """Fail fast on tier combinations the configured components cannot honor.""" self._compose_production() if self._production is not None and self._production_options is None and self.sources.names: raise ReliableReportingConfigurationError( "production adapters must belong to its frozen source registry" ) if self._worker_interval is not None and self._worker_interval <= timedelta(0): raise ReliableReportingConfigurationError("worker_interval must be greater than zero") if self._notification_worker is not None and self._notification_attempt_store is None: raise ReliableReportingConfigurationError( "a notification worker requires durable notification attempt storage" ) if self._receipt_handler is not None and self._materialization_worker is None: raise ReliableReportingConfigurationError( "receipt handling requires managed delivery/materialization" ) if self._reconciled_billing and ( self._receipt_handler is None or self._materialization_worker is None ): raise ReliableReportingConfigurationError( "reconciled billing requires managed delivery and a receipt handler" ) if self._production is None: for name in self.sources.names: for offering in self.sources.get(name).capability_offerings: if offering.get("method") and self._materialization_worker is None: raise ReliableReportingConfigurationError( "managed capability offerings require a materialization worker" ) if ( offering["reconciliation_mode"] == "consumer_receipt" and not self._reconciled_billing ): raise ReliableReportingConfigurationError( "consumer_receipt offerings require reconciled billing" )Fail fast on tier combinations the configured components cannot honor.
async def wait(self) ‑> None-
Expand source code
async def wait(self) -> None: """Wait for shutdown and raise a typed failure if supervision failed.""" await self._lifecycle.wait()Wait for shutdown and raise a typed failure if supervision failed.
class ReliableReportingServiceError (component: FailureComponent)-
Expand source code
class ReliableReportingServiceError(RuntimeError): """Closed diagnostic for an unexpected failure; contains no provider body.""" def __init__(self, component: FailureComponent) -> None: if component not in ( "startup", "configuration", "materialization", "notification", "service", "shutdown", ): raise ValueError("unknown reporting service failure component") self.component = component super().__init__(f"reporting service failed ({component})")Closed diagnostic for an unexpected failure; contains no provider body.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class ReliableReportingShutdownTimeoutError (*args, **kwargs)-
Expand source code
class ReliableReportingShutdownTimeoutError(TimeoutError): """A close waiter timed out; the service still owns the unsettled work."""A close waiter timed out; the service still owns the unsettled work.
Ancestors
- builtins.TimeoutError
- builtins.OSError
- builtins.Exception
- builtins.BaseException
class ReliableReportingState (*args, **kwds)-
Expand source code
class ReliableReportingState(str, Enum): NEW = "new" STARTING = "starting" INITIALIZED = "initialized" READY = "ready" STOPPING = "stopping" CLOSED = "closed" FAILED = "failed"str(object='') -> str str(bytes_or_buffer[, encoding[, errors]]) -> str
Create a new string object from the given object. If encoding or errors is specified, then the object must expose a data buffer that will be decoded using the given encoding and error handler. Otherwise, returns the result of object.str() (if defined) or repr(object). encoding defaults to sys.getdefaultencoding(). errors defaults to 'strict'.
Ancestors
- builtins.str
- enum.Enum
Class variables
var CLOSEDvar FAILEDvar INITIALIZEDvar NEWvar READYvar STARTINGvar STOPPING
-
Expand source code
class ReliableReportingUnavailableError(RuntimeError): """The service is not admitting new reporting work.""" def __init__(self, state: ReliableReportingState) -> None: self.state = state super().__init__(f"reporting service is unavailable ({state.value})")The service is not admitting new reporting work.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class ReportingAccountContext (account_id: str,
adapter: str,
currency: str,
source_scope: Mapping[str, Any],
account_timezone: str = 'UTC',
snapshot_offering_id: str | None = None,
official_offering_id: str | None = None,
requested_metrics: tuple[str, ...] = ('impressions', 'spend'),
requested_dimensions: tuple[str, ...] = (),
capability_offering: Mapping[str, Any] = <factory>,
publication_namespace: str = 'reporting-source:default',
slice_timeout: timedelta = datetime.timedelta(seconds=600))-
Expand source code
@dataclass(frozen=True) class ReportingAccountContext: """Trusted source facts frozen for one configuration generation.""" account_id: str adapter: str currency: str source_scope: Mapping[str, Any] account_timezone: str = "UTC" snapshot_offering_id: str | None = None official_offering_id: str | None = None requested_metrics: tuple[str, ...] = ("impressions", "spend") requested_dimensions: tuple[str, ...] = () capability_offering: Mapping[str, Any] = field(default_factory=dict) publication_namespace: str = "reporting-source:default" slice_timeout: timedelta = timedelta(minutes=10) def __post_init__(self) -> None: object.__setattr__(self, "requested_metrics", tuple(self.requested_metrics)) object.__setattr__(self, "requested_dimensions", tuple(self.requested_dimensions)) if not self.account_id: raise ReliableReportingConfigurationError("account_id is required") if not _ROUTE_RE.fullmatch(self.adapter): raise ReliableReportingConfigurationError( "adapter must be a stable identifier containing only letters, digits, ._:-" ) if not re.fullmatch(r"[A-Z]{3}", self.currency): raise ReliableReportingConfigurationError("currency must be an ISO 4217 alpha-3 code") try: ZoneInfo(self.account_timezone) except ZoneInfoNotFoundError as error: raise ReliableReportingConfigurationError( f"unknown account_timezone {self.account_timezone!r}" ) from error if self.snapshot_offering_id is None and self.official_offering_id is None: raise ReliableReportingConfigurationError( "at least one snapshot or official source offering is required" ) if not self.requested_metrics: raise ReliableReportingConfigurationError("requested_metrics must not be empty") if self.slice_timeout <= timedelta(0): raise ReliableReportingConfigurationError("slice_timeout must be greater than zero") object.__setattr__(self, "source_scope", _freeze(self.source_scope)) object.__setattr__(self, "capability_offering", _freeze(self.capability_offering)) def producer_offerings(self) -> ProducerOfferings: return ProducerOfferings( snapshot_offering_id=self.snapshot_offering_id, official_offering_id=self.official_offering_id, publication_namespace=self.publication_namespace, requested_metrics=self.requested_metrics, requested_dimensions=self.requested_dimensions, currency=self.currency, source_scope=_thaw(self.source_scope), slice_timeout=self.slice_timeout, )Trusted source facts frozen for one configuration generation.
Instance variables
var account_id : strvar account_timezone : strvar adapter : strvar capability_offering : Mapping[str, typing.Any]var currency : strvar official_offering_id : str | Nonevar publication_namespace : strvar requested_dimensions : tuple[str, ...]var requested_metrics : tuple[str, ...]var slice_timeout : datetime.timedeltavar snapshot_offering_id : str | Nonevar source_scope : Mapping[str, typing.Any]
Methods
def producer_offerings(self) ‑> ProducerOfferings-
Expand source code
def producer_offerings(self) -> ProducerOfferings: return ProducerOfferings( snapshot_offering_id=self.snapshot_offering_id, official_offering_id=self.official_offering_id, publication_namespace=self.publication_namespace, requested_metrics=self.requested_metrics, requested_dimensions=self.requested_dimensions, currency=self.currency, source_scope=_thaw(self.source_scope), slice_timeout=self.slice_timeout, )
class ReportingAdapter (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingAdapter(Protocol): """The minimum source adapter: declaration plus one normalized slice read. ``fetch_slice`` contains no MCP, A2A, or AdCP task-handler concerns. It may be synchronous (the SDK moves it to a worker thread) or asynchronous. """ @property def capabilities(self) -> ReportingSourceCapabilitiesV1: ... def fetch_slice( self, request: ReportingSourceSliceRequestV1 ) -> InlineFetchResult | Sequence[Mapping[str, Any]] | None | Awaitable[Any]: ...The minimum source adapter: declaration plus one normalized slice read.
fetch_slicecontains no MCP, A2A, or AdCP task-handler concerns. It may be synchronous (the SDK moves it to a worker thread) or asynchronous.Ancestors
- typing.Protocol
- typing.Generic
Instance variables
prop capabilities : ReportingSourceCapabilitiesV1-
Expand source code
@property def capabilities(self) -> ReportingSourceCapabilitiesV1: ...
Methods
def fetch_slice(self, request: ReportingSourceSliceRequestV1) ‑> InlineFetchResult | collections.abc.Sequence[collections.abc.Mapping[str, typing.Any]] | None | collections.abc.Awaitable[typing.Any]-
Expand source code
def fetch_slice( self, request: ReportingSourceSliceRequestV1 ) -> InlineFetchResult | Sequence[Mapping[str, Any]] | None | Awaitable[Any]: ...
class ReportingCheckpointStore (*args, **kwargs)-
Expand source code
class ReportingCheckpointStore(Protocol): async def get(self, reporting_materialization_id: str) -> ReportingReceipt | None: raise NotImplementedError async def put(self, receipt: ReportingReceipt) -> None: raise NotImplementedErrorBase class for protocol classes.
Protocol classes are defined as::
class Proto(Protocol): def meth(self) -> int: ...Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).
For example::
class C: def meth(self) -> int: return 0 def func(x: Proto) -> int: return x.meth() func(C()) # Passes static type checkSee PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::
class GenProto(Protocol[T]): def meth(self) -> T: ...Ancestors
- typing.Protocol
- typing.Generic
Methods
async def get(self, reporting_materialization_id: str) ‑> ReportingReceipt | None-
Expand source code
async def get(self, reporting_materialization_id: str) -> ReportingReceipt | None: raise NotImplementedError async def put(self, receipt: ReportingReceipt) ‑> None-
Expand source code
async def put(self, receipt: ReportingReceipt) -> None: raise NotImplementedError
class ReportingContentReading (reporting_revision_id: str,
media_buy_ids: tuple[str, ...] = (),
package_ids: tuple[str, ...] = (),
metric_names: tuple[str, ...] = (),
metric_units: Mapping[str, str] = <factory>,
control_total_units: Mapping[str, str] = <factory>,
time_dimension_values: tuple[datetime, ...] = (),
schema_conformant: bool = True,
observed_revision_content_sha256: str | None = None,
failure_code: ReportingFailureCode | None = None,
first_consumable_at: datetime | None = None)-
Expand source code
@dataclass(frozen=True) class ReportingContentReading: """What the buyer actually read out of one revision. Deliberately *not* the revision record the seller published: the whole point of ``content_mismatch`` is that the two disagree. A buyer fills this in from its own parse of the rows, which is why every field is what was observed rather than what was promised. """ reporting_revision_id: str #: Media buys with at least one row, **including** explicit zero rows. A #: buy present only as an explicit zero is covered; a buy with neither rows #: nor an explicit zero is ``scope_media_buy_missing``, because the revision #: cannot then distinguish zero delivery from an omitted buy. media_buy_ids: tuple[str, ...] = () package_ids: tuple[str, ...] = () #: Metric names present in the rows, matched against the pinned report #: definition's ``metrics[].name``. metric_names: tuple[str, ...] = () #: Observed unit per metric, matched against the unit the definition fixed. metric_units: Mapping[str, str] = field(default_factory=dict) #: Observed unit per control-total name. control_total_units: Mapping[str, str] = field(default_factory=dict) #: Values of the time dimension the pinned grain declares, if any. Used for #: the half-open period check, so a value exactly equal to ``period.end`` #: is out of range. time_dimension_values: tuple[datetime, ...] = () #: ``False`` when the rows failed structural validation against the pinned #: ``schema_uri``/``schema_sha256``. schema_conformant: bool = True #: The Core revision binding the buyer recomputed. Required to file either #: ``received`` or ``content_mismatch``: it proves which exact bytes the #: buyer is describing. observed_revision_content_sha256: str | None = None #: Set when the revision was *advertised* but its content could not be #: consumed at all. Turns the reading into ``unreadable`` rather than #: ``received``/``content_mismatch``, which need bytes to talk about. The #: two are genuinely different claims: ``unreadable`` says "I could not #: read what you published", ``content_mismatch`` says "I read it and it #: contradicts what we agreed". failure_code: ReportingFailureCode | None = None #: When the named revision first became consumable *to this consumer*. #: ``status_as_of`` for a ``received`` statement, because the spec has #: sellers use it as buyer-attributed arrival evidence rather than #: silently substituting publication time. Planning time is not an #: acceptable stand-in: it would date every statement to whenever the loop #: happened to run. first_consumable_at: datetime | None = NoneWhat the buyer actually read out of one revision.
Deliberately not the revision record the seller published: the whole point of
content_mismatchis that the two disagree. A buyer fills this in from its own parse of the rows, which is why every field is what was observed rather than what was promised.Instance variables
var control_total_units : Mapping[str, str]-
Observed unit per control-total name.
var failure_code : Literal['access_denied', 'resource_not_found', 'integrity_mismatch', 'reader_incompatible', 'transport_failed'] | None-
Set when the revision was advertised but its content could not be consumed at all. Turns the reading into
unreadablerather thanreceived/content_mismatch, which need bytes to talk about. The two are genuinely different claims:unreadablesays "I could not read what you published",content_mismatchsays "I read it and it contradicts what we agreed". var first_consumable_at : datetime.datetime | None-
When the named revision first became consumable to this consumer.
status_as_offor areceivedstatement, because the spec has sellers use it as buyer-attributed arrival evidence rather than silently substituting publication time. Planning time is not an acceptable stand-in: it would date every statement to whenever the loop happened to run. var media_buy_ids : tuple[str, ...]-
Media buys with at least one row, including explicit zero rows. A buy present only as an explicit zero is covered; a buy with neither rows nor an explicit zero is
scope_media_buy_missing, because the revision cannot then distinguish zero delivery from an omitted buy. var metric_names : tuple[str, ...]-
Metric names present in the rows, matched against the pinned report definition's
metrics[].name. var metric_units : Mapping[str, str]-
Observed unit per metric, matched against the unit the definition fixed.
var observed_revision_content_sha256 : str | None-
The Core revision binding the buyer recomputed. Required to file either
receivedorcontent_mismatch: it proves which exact bytes the buyer is describing. var package_ids : tuple[str, ...]var reporting_revision_id : strvar schema_conformant : bool-
Falsewhen the rows failed structural validation against the pinnedschema_uri/schema_sha256. var time_dimension_values : tuple[datetime.datetime, ...]-
Values of the time dimension the pinned grain declares, if any. Used for the half-open period check, so a value exactly equal to
period.endis out of range.
class ReportingFailurePlan-
Expand source code
class ReportingFailurePlan: """FIFO faults and barriers at named operation boundaries. ``at("seal.after", OSError("lost acknowledgement"))`` fails once *after* the seal store commits. ``None`` consumes one successful visit. Unplanned visits succeed. Use :meth:`assert_consumed` to catch a misspelled or unreached boundary. Adopter wrappers can use ``hit`` / ``hit_sync`` for additional destination, notification, and receipt boundaries. Queue access is thread-safe, but tests that need a particular ordering of concurrent calls must impose it with barriers. Errors are test inputs and must not contain real credentials or provider responses. """ def __init__(self) -> None: self._steps: dict[str, deque[BaseException | ReportingTestBarrier | None]] = defaultdict( deque ) self._hits: list[str] = [] self._lock = threading.Lock() def at(self, point: str, *steps: BaseException | ReportingTestBarrier | None) -> None: if not point or not steps: raise ValueError("a failure point needs a name and at least one step") if any( step is not None and not isinstance(step, (BaseException, ReportingTestBarrier)) for step in steps ): raise TypeError("failure steps must be exceptions, barriers, or None") with self._lock: self._steps[point].extend(steps) @property def hits(self) -> tuple[str, ...]: with self._lock: return tuple(self._hits) def _take(self, point: str) -> BaseException | ReportingTestBarrier | None: with self._lock: self._hits.append(point) steps = self._steps.get(point) return steps.popleft() if steps else None async def hit(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): await step.pause() def hit_sync(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): step.pause_sync() def assert_consumed(self) -> None: with self._lock: remaining = {point: len(steps) for point, steps in self._steps.items() if steps} if remaining: raise AssertionError(f"unreached reporting failure steps: {remaining}")FIFO faults and barriers at named operation boundaries.
at("seal.after", OSError("lost acknowledgement"))fails once after the seal store commits.Noneconsumes one successful visit. Unplanned visits succeed. Use :meth:assert_consumedto catch a misspelled or unreached boundary. Adopter wrappers can usehit/hit_syncfor additional destination, notification, and receipt boundaries.Queue access is thread-safe, but tests that need a particular ordering of concurrent calls must impose it with barriers. Errors are test inputs and must not contain real credentials or provider responses.
Instance variables
prop hits : tuple[str, ...]-
Expand source code
@property def hits(self) -> tuple[str, ...]: with self._lock: return tuple(self._hits)
Methods
def assert_consumed(self) ‑> None-
Expand source code
def assert_consumed(self) -> None: with self._lock: remaining = {point: len(steps) for point, steps in self._steps.items() if steps} if remaining: raise AssertionError(f"unreached reporting failure steps: {remaining}") def at(self,
point: str,
*steps: BaseException | ReportingTestBarrier | None) ‑> None-
Expand source code
def at(self, point: str, *steps: BaseException | ReportingTestBarrier | None) -> None: if not point or not steps: raise ValueError("a failure point needs a name and at least one step") if any( step is not None and not isinstance(step, (BaseException, ReportingTestBarrier)) for step in steps ): raise TypeError("failure steps must be exceptions, barriers, or None") with self._lock: self._steps[point].extend(steps) async def hit(self, point: str) ‑> None-
Expand source code
async def hit(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): await step.pause() def hit_sync(self, point: str) ‑> None-
Expand source code
def hit_sync(self, point: str) -> None: step = self._take(point) if isinstance(step, BaseException): raise step if isinstance(step, ReportingTestBarrier): step.pause_sync()
class ReportingInspectionContext (obligation: ReportingObligation,
revision: ReportingRevision,
materialization: ReportingMaterialization)-
Expand source code
@dataclass(frozen=True) class ReportingInspectionContext: obligation: ReportingObligation revision: ReportingRevision materialization: ReportingMaterializationReportingInspectionContext(obligation: 'ReportingObligation', revision: 'ReportingRevision', materialization: 'ReportingMaterialization')
Instance variables
var materialization : ReportingMaterializationvar obligation : ReportingObligationvar revision : ReportingRevision
class ReportingLedger (ledger_snapshot_id: str,
ledger_as_of: datetime,
account_id: str,
scope: BaseModel,
obligations: list[ReportingObligation],
revisions: list[ReportingRevision],
materializations: list[ReportingMaterialization],
receipts: list[ReportingReceipt],
consumer_statuses: list[Any] = <factory>,
revision_ownership: dict[str, str] | None = None,
adjustments: list[ReportingAdjustment] = <factory>,
adjustment_receipts: list[ReportingAdjustmentReceipt] = <factory>)-
Expand source code
@dataclass class ReportingLedger: ledger_snapshot_id: str ledger_as_of: datetime account_id: str scope: BaseModel obligations: list[ReportingObligation] revisions: list[ReportingRevision] materializations: list[ReportingMaterialization] receipts: list[ReportingReceipt] #: This caller's own consumer-status history, when the seller advertises #: the loop. Needed to compare a planned statement against the current #: leaf: without it a buyer cannot tell "nothing changed, do not re-file" #: from "new claim, must supersede", and re-filing the same claim under a #: new id churns the chain for no reason. consumer_statuses: list[Any] = field(default_factory=list) # None is the all-pages-absent legacy mode. An empty mapping is an explicit # new-mode empty snapshot; do not collapse those two meanings. revision_ownership: dict[str, str] | None = None adjustments: list[ReportingAdjustment] = field(default_factory=list) adjustment_receipts: list[ReportingAdjustmentReceipt] = field(default_factory=list) # Only the exhausted, non-incremental loader establishes this provenance. # A mutable/manual ledger remains useful for diagnostics, never completeness. _read_fingerprint: bytes | None = field(default=None, init=False, repr=False, compare=False) def __repr__(self) -> str: # Resources and extension fields may contain private transport values. return ( f"ReportingLedger(obligations={len(self.obligations)}, " f"revisions={len(self.revisions)}, materializations={len(self.materializations)}, " f"receipts={len(self.receipts)}, adjustments={len(self.adjustments)}, " f"adjustment_receipts={len(self.adjustment_receipts)})" )ReportingLedger(ledger_snapshot_id: 'str', ledger_as_of: 'datetime', account_id: 'str', scope: 'BaseModel', obligations: 'list[ReportingObligation]', revisions: 'list[ReportingRevision]', materializations: 'list[ReportingMaterialization]', receipts: 'list[ReportingReceipt]', consumer_statuses: 'list[Any]' =
, revision_ownership: 'dict[str, str] | None' = None, adjustments: 'list[ReportingAdjustment]' = , adjustment_receipts: 'list[ReportingAdjustmentReceipt]' = ) Instance variables
var account_id : strvar adjustment_receipts : list[ReportingAdjustmentReceipt]var adjustments : list[ReportingAdjustment]var consumer_statuses : list[typing.Any]-
This caller's own consumer-status history, when the seller advertises the loop. Needed to compare a planned statement against the current leaf: without it a buyer cannot tell "nothing changed, do not re-file" from "new claim, must supersede", and re-filing the same claim under a new id churns the chain for no reason.
var ledger_as_of : datetime.datetimevar ledger_snapshot_id : strvar materializations : list[ReportingMaterialization]var obligations : list[ReportingObligation]var receipts : list[ReportingReceipt]var revision_ownership : dict[str, str] | Nonevar revisions : list[ReportingRevision]var scope : pydantic.main.BaseModel
class ReportingLifecycleSnapshot (obligation: ReportingObligationRecord,
revisions: tuple[ReportingRevisionRecord, ...],
rows: tuple[bytes, ...])-
Expand source code
@dataclass(frozen=True) class ReportingLifecycleSnapshot: """An obligation, its exact revisions, and detached canonical row bytes. Capture a quiescent obligation after publication, then compare it after closing/rebuilding the service and retrying the same work. The assertion checks every revision, so an extra publication is a failure too. This is evidence about the supplied store, not a durability or tier certification. """ obligation: ReportingObligationRecord revisions: tuple[ReportingRevisionRecord, ...] rows: tuple[bytes, ...]An obligation, its exact revisions, and detached canonical row bytes.
Capture a quiescent obligation after publication, then compare it after closing/rebuilding the service and retrying the same work. The assertion checks every revision, so an extra publication is a failure too. This is evidence about the supplied store, not a durability or tier certification.
Instance variables
var obligation : ReportingObligationRecordvar revisions : tuple[ReportingRevisionRecord, ...]var rows : tuple[bytes, ...]
class ReportingObservation (row_count: int,
control_totals: list[ReportingControlTotal],
canonical_content_digest: ReportingCanonicalContentDigest | None = None,
manifest_sha256: str | None = None,
native_version_ref: str | None = None,
consumer_commit_ref: str | None = None)-
Expand source code
@dataclass(frozen=True) class ReportingObservation: row_count: int control_totals: list[ReportingControlTotal] canonical_content_digest: ReportingCanonicalContentDigest | None = None manifest_sha256: str | None = None native_version_ref: str | None = None consumer_commit_ref: str | None = NoneReportingObservation(row_count: 'int', control_totals: 'list[ReportingControlTotal]', canonical_content_digest: 'ReportingCanonicalContentDigest | None' = None, manifest_sha256: 'str | None' = None, native_version_ref: 'str | None' = None, consumer_commit_ref: 'str | None' = None)
Instance variables
var canonical_content_digest : ReportingCanonicalContentDigest | Nonevar consumer_commit_ref : str | Nonevar control_totals : list[ReportingControlTotal1 | ReportingControlTotal2]var manifest_sha256 : str | Nonevar native_version_ref : str | Nonevar row_count : int
class ReportingOperationsContactView (url: str | None = None, email: str | None = None)-
Expand source code
@dataclass(frozen=True) class ReportingOperationsContactView: """The seller's advertised human escalation path. Inert display metadata. Agents surface it to an operator and MUST NOT fetch the URL, send protocol traffic to it, or treat either value as a credential or a callback. It is carried as a separate type rather than as a bare dict so that intent is impossible to miss at the call site. """ url: str | None = None email: str | None = None @classmethod def from_capability( cls, payload: Mapping[str, Any] | None ) -> ReportingOperationsContactView | None: if not payload: return None contact = payload.get("operations_contact") if not isinstance(contact, Mapping): return None return cls(url=contact.get("url"), email=contact.get("email"))The seller's advertised human escalation path.
Inert display metadata. Agents surface it to an operator and MUST NOT fetch the URL, send protocol traffic to it, or treat either value as a credential or a callback. It is carried as a separate type rather than as a bare dict so that intent is impossible to miss at the call site.
Static methods
def from_capability(payload: Mapping[str, Any] | None) ‑> adcp.reporting._consumer.ReportingOperationsContactView | None
Instance variables
var email : str | Nonevar url : str | None
class ReportingPinnedDefinition (metric_names: tuple[str, ...] = (),
metric_units: Mapping[str, str] = <factory>,
control_total_units: Mapping[str, str] = <factory>)-
Expand source code
@dataclass(frozen=True) class ReportingPinnedDefinition: """The promises the accepted configuration generation already fixed. A buyer resolves this once from the retained report definition and reuses it, rather than re-reading the seller's current definition -- which is the point of pinning. Comparing against a definition the seller can change would make every mismatch unfalsifiable. """ metric_names: tuple[str, ...] = () #: Unit fixed for each metric, e.g. ``{"spend": "USD"}``. metric_units: Mapping[str, str] = field(default_factory=dict) #: Unit the profile defines for each control-total name. control_total_units: Mapping[str, str] = field(default_factory=dict)The promises the accepted configuration generation already fixed.
A buyer resolves this once from the retained report definition and reuses it, rather than re-reading the seller's current definition – which is the point of pinning. Comparing against a definition the seller can change would make every mismatch unfalsifiable.
Instance variables
var control_total_units : Mapping[str, str]-
Unit the profile defines for each control-total name.
var metric_names : tuple[str, ...]var metric_units : Mapping[str, str]-
Unit fixed for each metric, e.g.
{"spend": "USD"}.
class ReportingProductionOptions (offerings: tuple[ReportingServiceOffering, ...],
destination: ReportingProductionDestination,
registry: ReportingRevisionVerifierRegistry,
configuration_task: ReportingProductionConfigurationTask,
resolve_account: ReceiptAccountResolver,
buyer_agents: BuyerAgentRegistry | None = None,
notifications: bool = False,
subscriptions: ReportingSubscriptionResolver | None = None,
cipher: ReportingEnvelopeCipher | None = None,
signing: ReportingProductionSigning | None = None,
automated_recovery_window: timedelta = datetime.timedelta(seconds=21600),
status_retention_days: int = 400,
poll_seconds: float = 0.25)-
Expand source code
@dataclass(frozen=True) class ReportingProductionOptions: """Trusted inputs for ``ReliableReportingService.memory/postgres``. Register sources, then ``install(application)`` to compose and validate the graph. Mount the returned handler before starting the service. The factory owns its workers; the database pool and adopter providers remain borrowed. Receipt handlers and destination authorization use the existing production path. Push requires all three notification inputs and enabled notifications. """ offerings: tuple[ReportingServiceOffering, ...] destination: ReportingProductionDestination registry: ReportingRevisionVerifierRegistry configuration_task: ReportingProductionConfigurationTask resolve_account: ReceiptAccountResolver buyer_agents: BuyerAgentRegistry | None = None notifications: bool = False subscriptions: ReportingSubscriptionResolver | None = None cipher: ReportingEnvelopeCipher | None = None signing: ReportingProductionSigning | None = None automated_recovery_window: timedelta = timedelta(hours=6) status_retention_days: int = 400 poll_seconds: float = 0.25 def __post_init__(self) -> None: object.__setattr__(self, "offerings", tuple(self.offerings)) if not self.offerings: raise ValueError("production requires at least one source offering") notification_inputs = (self.subscriptions, self.cipher, self.signing) if any(value is not None for value in notification_inputs) and ( not self.notifications or any(value is None for value in notification_inputs) ): raise ValueError("push requires notifications, subscriptions, cipher and signing") def _capability_block( self, *, store: ReportingLedgerStore, sources: ReportingAdapterRegistry, escalation: ReportingDeliveryEscalation, consumer_status_enabled: bool, ) -> dict[str, Any]: from adcp.reporting.production._declaration import declared_reporting_delivery from adcp.reporting.production.pg import PgReportingProductionStore from adcp.reporting.service import _wire if set(sources.names) != {item.adapter for item in self.offerings}: return {} return declared_reporting_delivery( offerings=[_wire(item.offering) for item in self.offerings], destination=self.destination, durable=type(store) is PgReportingProductionStore, escalation=escalation, consumer_status_enabled=consumer_status_enabled, automated_recovery_window=self.automated_recovery_window, status_retention_days=self.status_retention_days, notifications=self.signing is not None, ) def _compose( self, *, store: InMemoryReportingProductionStore | PgReportingProductionStore, sources: ReportingAdapterRegistry, account_context: ReportingContextResolver, clock: Callable[[], datetime], escalation: ReportingDeliveryEscalation, consumer_status_enabled: bool, ) -> ReportingProductionSupport: from adcp.reporting.production.pg import PgReportingProductionStore if set(sources.names) != {item.adapter for item in self.offerings}: raise ValueError("every registered production adapter must have an offering") source_registry = ReportingProductionSourceRegistry(account_context=account_context) profiles: dict[str, ReportingServiceOffering] = {} producers: dict[str, ReportingProducer] = {} offerings: list[ReportingProductionOffering] = [] for item in self.offerings: registered = sources.get(item.adapter) previous = profiles.get(item.adapter) if previous is not None and ( previous.profile != item.profile or previous.verification_key != item.verification_key ): raise ValueError("one adapter must use one fixed source profile and verifier") if previous is None: producer = ReportingProducer( source=registered.executor, object_reader=registered.object_reader, offerings=item.profile, store=store, revision_verifier=self.registry.require(item.verification_key), escalation=escalation, clock=clock, worker_id=f"reporting-service:{item.adapter}", ) profiles[item.adapter] = item producers[item.adapter] = producer source_registry.register(item.adapter, producer) offerings.append( ReportingProductionOffering( item.offering, producers[item.adapter], item.verification_key, item.source_offering_id, ) ) projection: InMemoryReportingStatusProjection | PgReportingStatusProjection if type(store) is PgReportingProductionStore: projection = PgReportingStatusProjection( store, consumer_status_enabled=consumer_status_enabled, escalation=escalation, revision_ownership=True, ) elif type(store) is InMemoryReportingProductionStore: projection = InMemoryReportingStatusProjection( store, consumer_status_enabled=consumer_status_enabled, escalation=escalation, revision_ownership=True, ) else: raise ValueError("production requires the factory's matching production store") workers: tuple[ReportingNotificationWorker, ...] = () if self.subscriptions is not None and self.cipher is not None and self.signing is not None: workers = production_notification_workers( store, projection, subscriptions=self.subscriptions, cipher=self.cipher, signing=self.signing, ) return ReportingProductionSupport( ReportingMaterializerService( store, ReportingDestinationIO(self.registry, self.destination), self.destination ), projection, offerings=tuple(offerings), source_registry=source_registry, configuration_task=self.configuration_task, resolve_account=self.resolve_account, buyer_agents=self.buyer_agents, automated_recovery_window=self.automated_recovery_window, status_retention_days=self.status_retention_days, notification_workers=workers, poll_seconds=self.poll_seconds, )Trusted inputs for
ReliableReportingService.memory/postgres.Register sources, then
install(application)to compose and validate the graph. Mount the returned handler before starting the service. The factory owns its workers; the database pool and adopter providers remain borrowed. Receipt handlers and destination authorization use the existing production path. Push requires all three notification inputs and enabled notifications.Instance variables
var automated_recovery_window : timedeltavar buyer_agents : BuyerAgentRegistry | Nonevar cipher : ReportingEnvelopeCipher | Nonevar configuration_task : ReportingProductionConfigurationTaskvar destination : ReportingProductionDestinationvar notifications : boolvar offerings : tuple[ReportingServiceOffering, ...]var poll_seconds : floatvar registry : ReportingRevisionVerifierRegistryvar resolve_account : ReceiptAccountResolvervar signing : ReportingProductionSigning | Nonevar status_retention_days : intvar subscriptions : ReportingSubscriptionResolver | None
class ReportingReconciliationClient (*args, **kwargs)-
Expand source code
class ReportingReconciliationClient(ReportingStatusClient, Protocol): async def sync_reporting_receipts( self, request: SyncReportingReceiptsRequest ) -> TaskResult[SyncReportingReceiptsResponse]: raise NotImplementedErrorBase class for protocol classes.
Protocol classes are defined as::
class Proto(Protocol): def meth(self) -> int: ...Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).
For example::
class C: def meth(self) -> int: return 0 def func(x: Proto) -> int: return x.meth() func(C()) # Passes static type checkSee PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::
class GenProto(Protocol[T]): def meth(self) -> T: ...Ancestors
- adcp.reporting._reconcile.ReportingStatusClient
- typing.Protocol
- typing.Generic
Methods
async def sync_reporting_receipts(self, request: SyncReportingReceiptsRequest) ‑> TaskResult[SyncReportingReceiptsResponse]-
Expand source code
async def sync_reporting_receipts( self, request: SyncReportingReceiptsRequest ) -> TaskResult[SyncReportingReceiptsResponse]: raise NotImplementedError
class ReportingReconciliationError (code: str, message: str)-
Expand source code
class ReportingReconciliationError(RuntimeError): def __init__(self, code: str, message: str) -> None: super().__init__(message) self.code = codeUnspecified run-time error.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class ReportingReconciliationResult (definitive: bool,
ledger: ReportingLedger,
obligations: list[ObligationReconciliation],
missing_expected_periods: list[ExpectedReportingPeriod],
submitted_receipts: list[ReportingReceipt] = <factory>,
totals_by_revision: list[tuple[str, int, list[ReportingControlTotal]]] = <factory>,
consumer_loop: ConsumerLoopView | None = None,
consumer_status_plan: list[ConsumerStatusIntent] = <factory>)-
Expand source code
@dataclass class ReportingReconciliationResult: definitive: bool ledger: ReportingLedger obligations: list[ObligationReconciliation] missing_expected_periods: list[ExpectedReportingPeriod] submitted_receipts: list[ReportingReceipt] = field(default_factory=list) totals_by_revision: list[tuple[str, int, list[ReportingControlTotal]]] = field( default_factory=list ) #: AdCP 3.2.0-rc.3. What the seller said about *this buyer's* side of the #: loop: how many periods it owes a status for, the issue lifecycle fields #: for ageing a work item, and where to find a human. ``None`` when the #: seller does not advertise ``consumer_status_task``, which is distinct #: from an empty view -- see :class:`~adcp.reporting._consumer.ConsumerLoopView`. consumer_loop: ConsumerLoopView | None = None #: Statements the buyer owes right now, by the rc.3 deadline rather than by #: scope close. Empty when nothing is due. consumer_status_plan: list[ConsumerStatusIntent] = field(default_factory=list) @property def consumer_status_pending(self) -> int | None: """The seller's ``obligation_counts.consumer_status_pending``, if any.""" return self.consumer_loop.consumer_status_pending if self.consumer_loop else None @property def operations_contact(self) -> ReportingOperationsContactView | None: """The seller's advertised human escalation path. Never dereferenced.""" return self.consumer_loop.operations_contact if self.consumer_loop else None @property def consumer_mismatch_issues(self) -> tuple[ReportingStatusIssue, ...]: """Issues the seller raised from this buyer's own statements. Each carries ``opened_at`` (stable across re-emission, so it ages as one work item), ``issue_state``, and an inert ``external_ref``. """ return self.consumer_loop.mismatch_issues if self.consumer_loop else ()ReportingReconciliationResult(definitive: 'bool', ledger: 'ReportingLedger', obligations: 'list[ObligationReconciliation]', missing_expected_periods: 'list[ExpectedReportingPeriod]', submitted_receipts: 'list[ReportingReceipt]' =
, totals_by_revision: 'list[tuple[str, int, list[ReportingControlTotal]]]' = , consumer_loop: 'ConsumerLoopView | None' = None, consumer_status_plan: 'list[ConsumerStatusIntent]' = ) Instance variables
var consumer_loop : adcp.reporting._consumer.ConsumerLoopView | None-
AdCP 3.2.0-rc.3. What the seller said about this buyer's side of the loop: how many periods it owes a status for, the issue lifecycle fields for ageing a work item, and where to find a human.
Nonewhen the seller does not advertiseconsumer_status_task, which is distinct from an empty view – see :class:~adcp.reporting._consumer.ConsumerLoopView. prop consumer_mismatch_issues : tuple[ReportingStatusIssue, ...]-
Expand source code
@property def consumer_mismatch_issues(self) -> tuple[ReportingStatusIssue, ...]: """Issues the seller raised from this buyer's own statements. Each carries ``opened_at`` (stable across re-emission, so it ages as one work item), ``issue_state``, and an inert ``external_ref``. """ return self.consumer_loop.mismatch_issues if self.consumer_loop else ()Issues the seller raised from this buyer's own statements.
Each carries
opened_at(stable across re-emission, so it ages as one work item),issue_state, and an inertexternal_ref. prop consumer_status_pending : int | None-
Expand source code
@property def consumer_status_pending(self) -> int | None: """The seller's ``obligation_counts.consumer_status_pending``, if any.""" return self.consumer_loop.consumer_status_pending if self.consumer_loop else NoneThe seller's
obligation_counts.consumer_status_pending, if any. var consumer_status_plan : list[adcp.reporting._consumer.ConsumerStatusIntent]-
Statements the buyer owes right now, by the rc.3 deadline rather than by scope close. Empty when nothing is due.
var definitive : boolvar ledger : adcp.reporting._reconcile.ReportingLedgervar missing_expected_periods : list[adcp.reporting._reconcile.ExpectedReportingPeriod]var obligations : list[adcp.reporting._reconcile.ObligationReconciliation]prop operations_contact : ReportingOperationsContactView | None-
Expand source code
@property def operations_contact(self) -> ReportingOperationsContactView | None: """The seller's advertised human escalation path. Never dereferenced.""" return self.consumer_loop.operations_contact if self.consumer_loop else NoneThe seller's advertised human escalation path. Never dereferenced.
var submitted_receipts : list[ReportingReceipt]var totals_by_revision : list[tuple[str, int, list[ReportingControlTotal1 | ReportingControlTotal2]]]
class ReportingServiceOffering (adapter: str,
offering: ReportingDeliveryOffering,
profile: ProducerOfferings,
source_offering_id: str,
verification_key: ReportingVerificationKey)-
Expand source code
@dataclass(frozen=True) class ReportingServiceOffering: """Bind a public offering to one registered adapter's fixed source profile. Adapters may reuse local source offering IDs. Public offering IDs must be unique, and all offerings of one adapter must use the same profile and verifier. Account admission checks the resolved context against that profile and freezes the selected adapter for the configuration generation. """ adapter: str offering: ReportingDeliveryOffering profile: ProducerOfferings source_offering_id: str verification_key: ReportingVerificationKeyBind a public offering to one registered adapter's fixed source profile.
Adapters may reuse local source offering IDs. Public offering IDs must be unique, and all offerings of one adapter must use the same profile and verifier. Account admission checks the resolved context against that profile and freezes the selected adapter for the configuration generation.
Instance variables
var adapter : strvar offering : ReportingDeliveryOfferingvar profile : ProducerOfferingsvar source_offering_id : strvar verification_key : ReportingVerificationKey
class ReportingServiceResource (close: ResourceCallback, open: ResourceCallback | None = None)-
Expand source code
@dataclass(frozen=True) class ReportingServiceResource: """Explicit transfer of a resource's lifetime to one service. Resources open in declaration order before schema initialization and close once, in reverse order, after all admitted work settles. Ownership transfers at construction, so ``close`` must tolerate unopened or partially opened resources (including a service that never starts). Injected stores, adapters, workers, senders and pools are otherwise borrowed. Use async callbacks for loop-bound resources. Blocking synchronous callbacks run in a thread and are settled before cleanup can advance. Do not transfer the same resource to more than one owner. """ close: ResourceCallback = field(repr=False) open: ResourceCallback | None = field(default=None, repr=False) def __post_init__(self) -> None: if not callable(self.close) or (self.open is not None and not callable(self.open)): raise TypeError("resource callbacks must be callable")Explicit transfer of a resource's lifetime to one service.
Resources open in declaration order before schema initialization and close once, in reverse order, after all admitted work settles. Ownership transfers at construction, so
closemust tolerate unopened or partially opened resources (including a service that never starts). Injected stores, adapters, workers, senders and pools are otherwise borrowed.Use async callbacks for loop-bound resources. Blocking synchronous callbacks run in a thread and are settled before cleanup can advance. Do not transfer the same resource to more than one owner.
Instance variables
var close : Callable[[], object | collections.abc.Awaitable[object]]var open : collections.abc.Callable[[], object | collections.abc.Awaitable[object]] | None
class ReportingStatusClient (*args, **kwargs)-
Expand source code
class ReportingStatusClient(Protocol): async def get_reporting_status( self, request: GetReportingStatusRequest ) -> TaskResult[GetReportingStatusResponse]: raise NotImplementedErrorBase class for protocol classes.
Protocol classes are defined as::
class Proto(Protocol): def meth(self) -> int: ...Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).
For example::
class C: def meth(self) -> int: return 0 def func(x: Proto) -> int: return x.meth() func(C()) # Passes static type checkSee PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::
class GenProto(Protocol[T]): def meth(self) -> T: ...Ancestors
- typing.Protocol
- typing.Generic
Subclasses
- adcp.reporting._reconcile.ReportingReconciliationClient
Methods
async def get_reporting_status(self, request: GetReportingStatusRequest) ‑> TaskResult[GetReportingStatusResponse]-
Expand source code
async def get_reporting_status( self, request: GetReportingStatusRequest ) -> TaskResult[GetReportingStatusResponse]: raise NotImplementedError
class ReportingTestBarrier (*, timeout: float = 10)-
Expand source code
class ReportingTestBarrier: """Pause an async operation or a synchronous adapter running in a worker. Construct inside the test's event loop, await :meth:`wait` to observe the boundary, and always call :meth:`release` in test cleanup. The timeout is a real-time deadlock watchdog; it does not advance the reporting clock. """ def __init__(self, *, timeout: float = 10) -> None: if not math.isfinite(timeout) or timeout <= 0: raise ValueError("a reporting test barrier needs a finite positive timeout") self._loop = asyncio.get_running_loop() self._loop_thread = threading.get_ident() self._timeout = timeout self._entered = asyncio.Event() self._released = asyncio.Event() self._thread_released = threading.Event() async def wait(self) -> None: """Wait until the operation reaches this boundary, without polling.""" await asyncio.wait_for(self._entered.wait(), self._timeout) async def pause(self) -> None: self._entered.set() await asyncio.wait_for(self._released.wait(), self._timeout) def pause_sync(self) -> None: """Pause a worker thread; never block the event-loop thread.""" if threading.get_ident() == self._loop_thread: raise RuntimeError("pause_sync must run outside the barrier's event-loop thread") self._loop.call_soon_threadsafe(self._entered.set) if not self._thread_released.wait(self._timeout): raise TimeoutError("reporting test barrier was not released") def release(self) -> None: """Release all waiters, including a synchronous in-flight fetch.""" self._thread_released.set() self._loop.call_soon_threadsafe(self._released.set)Pause an async operation or a synchronous adapter running in a worker.
Construct inside the test's event loop, await :meth:
waitto observe the boundary, and always call :meth:releasein test cleanup. The timeout is a real-time deadlock watchdog; it does not advance the reporting clock.Methods
async def pause(self) ‑> None-
Expand source code
async def pause(self) -> None: self._entered.set() await asyncio.wait_for(self._released.wait(), self._timeout) def pause_sync(self) ‑> None-
Expand source code
def pause_sync(self) -> None: """Pause a worker thread; never block the event-loop thread.""" if threading.get_ident() == self._loop_thread: raise RuntimeError("pause_sync must run outside the barrier's event-loop thread") self._loop.call_soon_threadsafe(self._entered.set) if not self._thread_released.wait(self._timeout): raise TimeoutError("reporting test barrier was not released")Pause a worker thread; never block the event-loop thread.
def release(self) ‑> None-
Expand source code
def release(self) -> None: """Release all waiters, including a synchronous in-flight fetch.""" self._thread_released.set() self._loop.call_soon_threadsafe(self._released.set)Release all waiters, including a synchronous in-flight fetch.
async def wait(self) ‑> None-
Expand source code
async def wait(self) -> None: """Wait until the operation reaches this boundary, without polling.""" await asyncio.wait_for(self._entered.wait(), self._timeout)Wait until the operation reaches this boundary, without polling.
class ReportingTier (*args, **kwds)-
Expand source code
class ReportingTier(str, Enum): """Feature tiers advertised by ``media_buy.reporting_delivery``.""" CORE = "core" MANAGED_DELIVERY = "managed_delivery" RECONCILED_BILLING = "reconciled_billing"Feature tiers advertised by
media_buy.reporting_delivery.Ancestors
- builtins.str
- enum.Enum
Class variables
var COREvar MANAGED_DELIVERYvar RECONCILED_BILLING
class ScriptedReportingAdapter (capabilities: ReportingSourceCapabilitiesV1,
results: Sequence[InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException] = (),
*,
asynchronous: bool = False,
failures: ReportingFailurePlan | None = None)-
Expand source code
class ScriptedReportingAdapter: """Queue source answers or exceptions for deterministic service tests. Set ``asynchronous=True`` to exercise the async provider-SDK path. The default synchronous path is run in a worker thread by :class:`InlineReportingSource`, just like a synchronous production SDK. With ``failures``, each call visits ``fetch.before`` before consuming its queued answer and ``fetch.after`` after a successful answer. These points accept either an exception or a :class:`ReportingTestBarrier`. """ def __init__( self, capabilities: ReportingSourceCapabilitiesV1, results: Sequence[ InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException ] = (), *, asynchronous: bool = False, failures: ReportingFailurePlan | None = None, ) -> None: self._capabilities = capabilities self._results = deque(results) self._lock = threading.Lock() self._failures = failures self.asynchronous = asynchronous self.calls: list[ReportingSourceSliceRequestV1] = [] @property def capabilities(self) -> ReportingSourceCapabilitiesV1: return self._capabilities def queue( self, result: InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException, ) -> None: with self._lock: self._results.append(result) def _next(self) -> Any: with self._lock: if not self._results: raise AssertionError("ScriptedReportingAdapter has no queued source answer") result = self._results.popleft() if isinstance(result, BaseException): raise result return result def fetch_slice(self, request: ReportingSourceSliceRequestV1) -> Any: if not self.asynchronous: with self._lock: self.calls.append(request) if self._failures is not None: self._failures.hit_sync("fetch.before") result = self._next() if self._failures is not None: self._failures.hit_sync("fetch.after") return result async def fetch_async() -> Any: with self._lock: self.calls.append(request) if self._failures is not None: await self._failures.hit("fetch.before") result = self._next() if self._failures is not None: await self._failures.hit("fetch.after") return result return fetch_async()Queue source answers or exceptions for deterministic service tests.
Set
asynchronous=Trueto exercise the async provider-SDK path. The default synchronous path is run in a worker thread by :class:InlineReportingSource, just like a synchronous production SDK. Withfailures, each call visitsfetch.beforebefore consuming its queued answer andfetch.afterafter a successful answer. These points accept either an exception or a :class:ReportingTestBarrier.Instance variables
prop capabilities : ReportingSourceCapabilitiesV1-
Expand source code
@property def capabilities(self) -> ReportingSourceCapabilitiesV1: return self._capabilities
Methods
def fetch_slice(self, request: ReportingSourceSliceRequestV1) ‑> Any-
Expand source code
def fetch_slice(self, request: ReportingSourceSliceRequestV1) -> Any: if not self.asynchronous: with self._lock: self.calls.append(request) if self._failures is not None: self._failures.hit_sync("fetch.before") result = self._next() if self._failures is not None: self._failures.hit_sync("fetch.after") return result async def fetch_async() -> Any: with self._lock: self.calls.append(request) if self._failures is not None: await self._failures.hit("fetch.before") result = self._next() if self._failures is not None: await self._failures.hit("fetch.after") return result return fetch_async() def queue(self,
result: InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException) ‑> None-
Expand source code
def queue( self, result: InlineFetchResult | Sequence[Mapping[str, Any]] | None | BaseException, ) -> None: with self._lock: self._results.append(result)