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 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.

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_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.

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_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.

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 plan

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.

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.

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.

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 result

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.

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 result

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.

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 resolved

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.

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 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.

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_pending lives in obligation_counts and the periods view does not carry it.

Instance variables

var consumer_status_pending : int | None
var 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_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.

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 = None

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.

Instance variables

var reporting_status_id : str
var 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_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.

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 payload

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.

Instance variables

var consumer_commit_ref : str | None
var consumer_status : Literal['received', 'obligation_missing', 'revision_missing', 'unreadable', 'content_mismatch']
var delivery_config_id : str
var delivery_config_version : int
var due_at : datetime.datetime
var failure_code : str | None
var mismatch_code : Literal['scope_media_buy_missing', 'coverage_short', 'metric_missing', 'schema_nonconformant', 'currency_mismatch', 'period_mismatch'] | None
var observed_revision_content_sha256 : str | None
prop overdue : bool
Expand source code
@property
def overdue(self) -> bool:
    return _utc(self.status_as_of) >= _utc(self.due_at)
var period_end : datetime.datetime
var period_source_timezone : str
var period_start : datetime.datetime
var report_definition_id : str
var reporting_obligation_id : str | None
var reporting_revision_id : str | None
var reporting_status_id : str
var status_as_of : datetime.datetime
var 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 payload

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<code> on a </code>received statement 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_missing accuses the seller of publishing nothing, received asserts 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_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.

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 = None

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)

Instance variables

var delivery_config_id : str
var delivery_config_version : int
var expected_at : str | None

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.

var feed_purpose : str
var media_buy_ids : tuple[str, ...]
var period_end : str
var period_start : str
var report_definition_id : str
var reporting_profile : str
var source_timezone : 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.

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 result

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.

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 result

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.

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] = checkpoint

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.

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 : bool
var reasons : tuple[str, ...]
var reporting_materialization_id : str | None
var reporting_obligation_id : str
var 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 = kind

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.

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 in owned_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.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.

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.failure

Sanitized 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.ready

Whether 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.state

Current 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 payload

Project 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_resources to 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 payload

Merge 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 platform

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 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.

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 CLOSED
var FAILED
var INITIALIZED
var NEW
var READY
var STARTING
var STOPPING
class ReliableReportingUnavailableError (state: ReliableReportingState)
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 : str
var account_timezone : str
var adapter : str
var capability_offering : Mapping[str, typing.Any]
var currency : str
var official_offering_id : str | None
var publication_namespace : str
var requested_dimensions : tuple[str, ...]
var requested_metrics : tuple[str, ...]
var slice_timeout : datetime.timedelta
var snapshot_offering_id : str | None
var 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_slice contains 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 NotImplementedError

Base 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 check

See 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 = None

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.

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 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".

var first_consumable_at : datetime.datetime | 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.

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 received or content_mismatch: it proves which exact bytes the buyer is describing.

var package_ids : tuple[str, ...]
var reporting_revision_id : str
var schema_conformant : bool

False when the rows failed structural validation against the pinned schema_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.end is 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. 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.

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: ReportingMaterialization

ReportingInspectionContext(obligation: 'ReportingObligation', revision: 'ReportingRevision', materialization: 'ReportingMaterialization')

Instance variables

var materialization : ReportingMaterialization
var obligation : ReportingObligation
var 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 : str
var 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.datetime
var ledger_snapshot_id : str
var materializations : list[ReportingMaterialization]
var obligations : list[ReportingObligation]
var receipts : list[ReportingReceipt]
var revision_ownership : dict[str, str] | None
var 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 : ReportingObligationRecord
var 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 = None

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)

Instance variables

var canonical_content_digest : ReportingCanonicalContentDigest | None
var consumer_commit_ref : str | None
var control_totals : list[ReportingControlTotal1 | ReportingControlTotal2]
var manifest_sha256 : str | None
var native_version_ref : str | None
var 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 | None
var 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 : timedelta
var buyer_agents : BuyerAgentRegistry | None
var cipher : ReportingEnvelopeCipher | None
var configuration_task : ReportingProductionConfigurationTask
var destination : ReportingProductionDestination
var notifications : bool
var offerings : tuple[ReportingServiceOffering, ...]
var poll_seconds : float
var registry : ReportingRevisionVerifierRegistry
var resolve_account : ReceiptAccountResolver
var signing : ReportingProductionSigning | None
var status_retention_days : int
var 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 NotImplementedError

Base 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 check

See 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 = code

Unspecified 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. None when the seller does not advertise consumer_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 inert external_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 None

The 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 : bool
var ledger : adcp.reporting._reconcile.ReportingLedger
var 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 None

The 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: ReportingVerificationKey

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.

Instance variables

var adapter : str
var offering : ReportingDeliveryOffering
var profile : ProducerOfferings
var source_offering_id : str
var 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 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.

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 NotImplementedError

Base 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 check

See 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: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.

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 CORE
var MANAGED_DELIVERY
var 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=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.

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)