Module adcp.reporting.submissions

Buyer-only durable receipt submission, separate from reconciliation planning.

Public imports are additive and work without the optional PostgreSQL driver. The application supplies a trusted authorizer, a receipt client, and an explicit intent store. Use PgReportingSubmissionIntentStore for restart durability; the memory implementation is a volatile reference for tests. No checkpoint, seller service, consumer-status, or client.reporting facade behavior is changed.

Rollout: use a dedicated buyer schema/pool; explicitly call create_schema before enabling submissions. The migration adds only reporting_buyer_submission_*. Retain pending intents indefinitely and resume them after any uncertain result. Do not drop the tables on rollback or replace pending plans with new receipt IDs. Disable new planning while recovering, then resume the same stored scope with current authorization. Exact final-head review precedes facade integration.

Sub-modules

adcp.reporting.submissions.models

Immutable buyer intent records and closed diagnostics …

adcp.reporting.submissions.pg

PostgreSQL buyer intent storage with short transactions and no network locks …

adcp.reporting.submissions.store

Optional buyer intent protocol, independent of ReportingCheckpointStore.

adcp.reporting.submissions.submit

Explicit low-level receipt submission with durable uncertainty semantics.

Functions

def prepare_reporting_receipt_submission(scope: ReportingSubmissionScope,
receipts: Sequence[ReportingSubmissionReceipt]) ‑> ReportingReceiptSubmission
Expand source code
def prepare_reporting_receipt_submission(
    scope: ReportingSubmissionScope,
    receipts: Sequence[ReportingSubmissionReceipt],
) -> ReportingReceiptSubmission:
    """Freeze an already validated receipt plan, with at most 100 items per call.

    Outbound typed values are normalized once, then retained byte-for-byte.
    This is not a capture or proof of raw inbound adjustment evidence. Requests
    contain the authorized resolved account only, with no caller-provided
    account, idempotency key, context, extensions or credentials.
    """
    return _prepare_for_version(scope, receipts, _VERSION)

Freeze an already validated receipt plan, with at most 100 items per call.

Outbound typed values are normalized once, then retained byte-for-byte. This is not a capture or proof of raw inbound adjustment evidence. Requests contain the authorized resolved account only, with no caller-provided account, idempotency key, context, extensions or credentials.

async def submit_reporting_receipts(client: ReportingReceiptSubmissionClient,
*,
authorizer: ReportingSubmissionAuthorizer,
store: ReportingSubmissionIntentStore,
receipts: Sequence[ReportingSubmissionReceipt] | None = None,
timeout_seconds: float = 30.0) ‑> ReportingSubmissionResult
Expand source code
async def submit_reporting_receipts(
    client: ReportingReceiptSubmissionClient,
    *,
    authorizer: ReportingSubmissionAuthorizer,
    store: ReportingSubmissionIntentStore,
    receipts: Sequence[ReportingSubmissionReceipt] | None = None,
    timeout_seconds: float = 30.0,
) -> ReportingSubmissionResult:
    """Reserve or resume a receipt plan without replacing uncertain work.

    Omit ``receipts`` to resume the latest intent for the freshly authorized
    scope, including its completed outcomes after a lost return. Otherwise the
    supplied plan is frozen and atomically reserved before ANY network write.
    An earlier pending scope intent takes priority, with proposal_deferred=True.
    The caller must reconsider the deferred plan against fresh seller history.

    A transport failure, timeout, failed task or malformed response leaves the
    chunk uncertain and returns pending=True. Retry within the live protocol
    version uses the exact durable body and idempotency key. A pending intent
    from a retired version remains readable and byte-identical but is refused
    before transport; it is never rewritten as a new-version request.
    Cancellation propagates with the intent still retained.
    Storage failure raises a closed error; resume the stored scope because the
    last commit may have succeeded. There is no expiry/abandon/replacement path.

    Every item failure is retained alongside successes. Completion only means
    every submission outcome is confirmed; it does not mean every receipt was
    recorded, accepted, readable, or sufficient for definitive reconciliation.
    This API does not select history, build adjustment evidence, post consumer
    statuses, or modify the original reconciliation/checkpoint interfaces.
    """
    if (
        isinstance(timeout_seconds, bool)
        or not isinstance(timeout_seconds, (int, float))
        or not isfinite(timeout_seconds)
        or timeout_seconds <= 0
    ):
        raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN)
    scope = await _authorize(client, authorizer)
    proposed = (
        prepare_reporting_receipt_submission(scope, receipts) if receipts is not None else None
    )
    state = await _storage(store.reserve(proposed) if proposed is not None else store.get(scope))
    if state is None:
        raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND)
    _check_state(state, scope)
    deferred = proposed is not None and proposed.submission_id != state.submission_id
    if state.pending and state.adcp_version != _VERSION:
        # An unresolved older reservation retains its exact wire bytes, but
        # the current SDK must not issue a new request under a retired pin.
        raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN)
    while state.pending:
        await _authorize(client, authorizer, scope)
        chunk = state.confirmed_chunks
        request = state.request(chunk)
        response = None
        try:
            response = await asyncio.wait_for(
                client.sync_reporting_receipts(request), timeout=timeout_seconds
            )
        except Exception:
            response = None
        if response is None:
            return ReportingSubmissionResult(
                state, deferred, ReportingSubmissionCode.TRANSPORT_UNCERTAIN
            )
        if (
            response.success is not True
            or response.status != "completed"
            or not isinstance(response.data, SyncReportingReceiptsResponse)
        ):
            return ReportingSubmissionResult(
                state, deferred, ReportingSubmissionCode.RESPONSE_UNCONFIRMED
            )
        invalid = False
        try:
            updated = await _storage(
                store.confirm(scope, state.submission_id, chunk, response.data)
            )
        except ReportingSubmissionError as error:
            if error.code != ReportingSubmissionCode.INVALID_RESPONSE:
                raise
            invalid = True
        if invalid:
            return ReportingSubmissionResult(
                state, deferred, ReportingSubmissionCode.INVALID_RESPONSE
            )
        _check_state(updated, scope, state)
        if updated.confirmed_chunks <= chunk:
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        state = updated
    await _authorize(client, authorizer, scope)
    return ReportingSubmissionResult(state, deferred)

Reserve or resume a receipt plan without replacing uncertain work.

Omit receipts to resume the latest intent for the freshly authorized scope, including its completed outcomes after a lost return. Otherwise the supplied plan is frozen and atomically reserved before ANY network write. An earlier pending scope intent takes priority, with proposal_deferred=True. The caller must reconsider the deferred plan against fresh seller history.

A transport failure, timeout, failed task or malformed response leaves the chunk uncertain and returns pending=True. Retry within the live protocol version uses the exact durable body and idempotency key. A pending intent from a retired version remains readable and byte-identical but is refused before transport; it is never rewritten as a new-version request. Cancellation propagates with the intent still retained. Storage failure raises a closed error; resume the stored scope because the last commit may have succeeded. There is no expiry/abandon/replacement path.

Every item failure is retained alongside successes. Completion only means every submission outcome is confirmed; it does not mean every receipt was recorded, accepted, readable, or sufficient for definitive reconciliation. This API does not select history, build adjustment evidence, post consumer statuses, or modify the original reconciliation/checkpoint interfaces.

Classes

class InMemoryReportingSubmissionIntentStore
Expand source code
class InMemoryReportingSubmissionIntentStore:
    """Volatile reference implementation for tests; it cannot survive restart.

    Production callers must explicitly supply a durable implementation such as
    PgReportingSubmissionIntentStore. There is no implicit in-memory fallback.
    """

    def __init__(self) -> None:
        self._lock = asyncio.Lock()
        self._scopes: dict[str, tuple[bytes, str]] = {}
        self._submissions: dict[tuple[str, str], ReportingReceiptSubmission] = {}

    def _current(self, scope: ReportingSubmissionScope) -> ReportingReceiptSubmission | None:
        entry = self._scopes.get(scope.storage_key)
        if entry is None:
            return None
        if entry[0] != scope.canonical_identity:
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        current = self._submissions.get((scope.storage_key, entry[1]))
        if current is None or current.scope != scope:
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        validate_submission(current)
        return current

    async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission:
        validate_submission(proposed)
        if proposed.confirmed_chunks:
            raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN)
        async with self._lock:
            current = self._current(proposed.scope)
            if current is not None and current.pending:
                return current
            key = proposed.scope.storage_key, proposed.submission_id
            prior = self._submissions.get(key)
            if prior is not None:
                validate_submission(prior)
                if prior.scope != proposed.scope or (
                    prior._plan != proposed._plan
                    and not is_completed_legacy_replay(prior, proposed)
                ):
                    raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
                return prior
            self._submissions[key] = proposed
            self._scopes[key[0]] = proposed.scope.canonical_identity, proposed.submission_id
            return proposed

    async def get(
        self, scope: ReportingSubmissionScope, submission_id: str | None = None
    ) -> ReportingReceiptSubmission | None:
        async with self._lock:
            current = self._current(scope)
            if current is None or submission_id is None:
                return current
            result = self._submissions.get((scope.storage_key, submission_id))
            if result is not None:
                validate_submission(result)
                if result.scope != scope:
                    raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
            return result

    async def confirm(
        self,
        scope: ReportingSubmissionScope,
        submission_id: str,
        chunk: int,
        response: SyncReportingReceiptsResponse,
    ) -> ReportingReceiptSubmission:
        frozen = freeze_response(response)
        async with self._lock:
            current = self._current(scope)
            state = self._submissions.get((scope.storage_key, submission_id))
            if current is None or state is None:
                raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND)
            validate_submission(state)
            if state.scope != scope or (
                state.pending and state.submission_id != current.submission_id
            ):
                raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
            updated = confirm_submission(state, chunk, frozen)
            self._submissions[(scope.storage_key, submission_id)] = updated
            return updated

Volatile reference implementation for tests; it cannot survive restart.

Production callers must explicitly supply a durable implementation such as PgReportingSubmissionIntentStore. There is no implicit in-memory fallback.

Methods

async def confirm(self,
scope: ReportingSubmissionScope,
submission_id: str,
chunk: int,
response: SyncReportingReceiptsResponse) ‑> ReportingReceiptSubmission
Expand source code
async def confirm(
    self,
    scope: ReportingSubmissionScope,
    submission_id: str,
    chunk: int,
    response: SyncReportingReceiptsResponse,
) -> ReportingReceiptSubmission:
    frozen = freeze_response(response)
    async with self._lock:
        current = self._current(scope)
        state = self._submissions.get((scope.storage_key, submission_id))
        if current is None or state is None:
            raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND)
        validate_submission(state)
        if state.scope != scope or (
            state.pending and state.submission_id != current.submission_id
        ):
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        updated = confirm_submission(state, chunk, frozen)
        self._submissions[(scope.storage_key, submission_id)] = updated
        return updated
async def get(self,
scope: ReportingSubmissionScope,
submission_id: str | None = None) ‑> ReportingReceiptSubmission | None
Expand source code
async def get(
    self, scope: ReportingSubmissionScope, submission_id: str | None = None
) -> ReportingReceiptSubmission | None:
    async with self._lock:
        current = self._current(scope)
        if current is None or submission_id is None:
            return current
        result = self._submissions.get((scope.storage_key, submission_id))
        if result is not None:
            validate_submission(result)
            if result.scope != scope:
                raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        return result
async def reserve(self,
proposed: ReportingReceiptSubmission) ‑> ReportingReceiptSubmission
Expand source code
async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission:
    validate_submission(proposed)
    if proposed.confirmed_chunks:
        raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN)
    async with self._lock:
        current = self._current(proposed.scope)
        if current is not None and current.pending:
            return current
        key = proposed.scope.storage_key, proposed.submission_id
        prior = self._submissions.get(key)
        if prior is not None:
            validate_submission(prior)
            if prior.scope != proposed.scope or (
                prior._plan != proposed._plan
                and not is_completed_legacy_replay(prior, proposed)
            ):
                raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
            return prior
        self._submissions[key] = proposed
        self._scopes[key[0]] = proposed.scope.canonical_identity, proposed.submission_id
        return proposed
class PgReportingSubmissionIntentStore (*, pool: AsyncConnectionPool)
Expand source code
class PgReportingSubmissionIntentStore:
    """Production persistence for optional buyer receipt submission intents.

    Pass a trusted application-owned psycopg AsyncConnectionPool, preferably
    with a dedicated schema/search_path and bounded connection/statement waits.
    There is no credential or connection string in this object's representation.
    Run create_schema explicitly during rollout, then validate custom backends
    with the same memory/PostgreSQL state-machine vectors. Readiness requires all
    shipped, validated CHECK definitions as well as the keys and reservation
    index. Snapshot validation happens outside transactions; a mutation compares
    the complete locked state, retrying at most 16 times. Contention exhaustion
    raises STORAGE_UNAVAILABLE and leaves the durable intent available to resume.
    """

    def __init__(self, *, pool: AsyncConnectionPool) -> None:
        self._pool = pool

    @_closed_storage_errors
    async def create_schema(self) -> None:
        available = False
        try:
            import psycopg  # noqa: F401

            available = True
        except ImportError:
            # The optional driver is absent; report PG_REQUIRED outside the handler.
            pass
        if not available:
            raise ReportingSubmissionError(ReportingSubmissionCode.PG_REQUIRED)
        async with self._pool.connection() as connection, connection.transaction():
            await connection.execute("SELECT pg_advisory_xact_lock(%s)", (712071172,))
            sql = (
                files("adcp.reporting.ledger")
                .joinpath("reporting_buyer_submissions.sql")
                .read_text()
            )
            await connection.execute(sql)
            await self._check_schema_on(connection)

    async def _check_schema_on(self, connection: Any) -> None:
        # Verify storage bounds as well as concurrency invariants. IF NOT EXISTS
        # must not silently bless a pre-created table with weaker constraints.
        rows = await (
            await connection.execute(
                "SELECT c.conrelid='reporting_buyer_submission_scopes'::regclass,"
                " c.conname,c.contype,c.convalidated,c.condeferrable,c.connoinherit,"
                " pg_get_constraintdef(c.oid) FROM pg_constraint c"
                " WHERE c.conrelid IN ('reporting_buyer_submission_scopes'::regclass,"
                " 'reporting_buyer_submission_intents'::regclass)"
            )
        ).fetchall()
        constraints = {(row[0], row[1]): tuple(row[2:]) for row in rows}
        required = {
            (True, "reporting_buyer_submission_scopes_pkey"): (
                "p",
                True,
                False,
                True,
                "PRIMARY KEY (scope_sha256)",
            ),
            (False, "reporting_buyer_submission_intents_pkey"): (
                "p",
                True,
                False,
                True,
                "PRIMARY KEY (scope_sha256, submission_id)",
            ),
            (False, "reporting_buyer_submission_scope_fk"): (
                "f",
                True,
                False,
                True,
                "FOREIGN KEY (scope_sha256) REFERENCES "
                "reporting_buyer_submission_scopes(scope_sha256)",
            ),
        }
        checks = (
            (True, "scope_digest", "scope_sha256 ~ '^[a-f0-9]{64}$'::text"),
            (True, "scope_bound", "octet_length(canonical_identity) <= 32768"),
            (False, "submission_id", "submission_id ~ '^reporting-submission:[a-f0-9]{64}$'::text"),
            (False, "plan_digest", "plan_sha256 ~ '^[a-f0-9]{64}$'::text"),
            (False, "confirmed_digest", "confirmed_sha256 ~ '^[a-f0-9]{64}$'::text"),
            (False, "plan_bound", "octet_length(canonical_plan) <= 16777216"),
            (False, "confirmed_bound", "octet_length(confirmed_results) <= 16777216"),
        )
        for scopes_table, suffix, expression in checks:
            required[(scopes_table, f"reporting_buyer_{suffix}")] = (
                "c",
                True,
                False,
                False,
                f"CHECK (({expression}))",
            )
        if any(constraints.get(name) != expected for name, expected in required.items()):
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        index = await (
            await connection.execute(
                "SELECT i.indisunique,i.indisvalid,i.indisready,i.indnkeyatts,"
                " pg_get_indexdef(i.indexrelid,1,true),pg_get_expr(i.indpred,i.indrelid),"
                " i.indrelid='reporting_buyer_submission_intents'::regclass"
                " FROM pg_index i WHERE i.indexrelid='reporting_buyer_one_pending_scope'::regclass"
            )
        ).fetchone()
        if index is None or tuple(index) != (True, True, True, 1, "scope_sha256", "pending", True):
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)

    async def _snapshot_on(
        self, connection: Any, key: str, selected_id: str | None, *, lock: bool = False
    ) -> _Snapshot:
        scope_query = (
            "SELECT canonical_identity,current_submission_id"
            " FROM reporting_buyer_submission_scopes WHERE scope_sha256=%s FOR UPDATE"
            if lock
            else "SELECT canonical_identity,current_submission_id"
            " FROM reporting_buyer_submission_scopes WHERE scope_sha256=%s"
        )
        row = await (await connection.execute(scope_query, (key,))).fetchone()
        if row is None:
            return _Snapshot(None, None, None, False)
        intent_query = (
            "SELECT submission_id,canonical_plan,plan_sha256,confirmed_results,"
            " confirmed_sha256,pending FROM reporting_buyer_submission_intents"
            " WHERE scope_sha256=%s AND (submission_id=%s OR submission_id=%s) FOR UPDATE"
            if lock
            else "SELECT submission_id,canonical_plan,plan_sha256,confirmed_results,"
            " confirmed_sha256,pending FROM reporting_buyer_submission_intents"
            " WHERE scope_sha256=%s AND (submission_id=%s OR submission_id=%s)"
        )
        rows = await (await connection.execute(intent_query, (key, row[1], selected_id))).fetchall()
        by_id = {value[0]: tuple(value) for value in rows}
        has_intents = bool(rows)
        if row[1] is None:
            has_intents = (
                await (
                    await connection.execute(
                        "SELECT 1 FROM reporting_buyer_submission_intents"
                        " WHERE scope_sha256=%s LIMIT 1",
                        (key,),
                    )
                ).fetchone()
            ) is not None
        return _Snapshot(
            tuple(row),
            by_id.get(row[1]),
            by_id.get(row[1] if selected_id is None else selected_id),
            has_intents,
        )

    async def _snapshot(self, key: str, selected_id: str | None) -> _Snapshot:
        # Release the read-only MVCC snapshot and connection before CPU work.
        # No scope row lock or caller-owned validation occurs in this phase.
        async with self._pool.connection() as connection, connection.transaction():
            await connection.execute("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY")
            return await self._snapshot_on(connection, key, selected_id)

    def _validate_snapshot(
        self, scope: ReportingSubmissionScope, snapshot: _Snapshot
    ) -> tuple[ReportingReceiptSubmission | None, ReportingReceiptSubmission | None]:
        if snapshot.scope is None:
            return None, None
        if snapshot.scope[0] != scope.canonical_identity.decode() or (
            snapshot.scope[1] is None and snapshot.has_intents
        ):
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        if snapshot.scope[1] is not None and snapshot.current is None:
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        current = decode_submission(scope, snapshot.current) if snapshot.current else None
        selected = (
            current
            if snapshot.selected is snapshot.current
            else decode_submission(scope, snapshot.selected) if snapshot.selected else None
        )
        if selected is not None and selected.pending and selected != current:
            raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
        return current, selected

    @_closed_storage_errors
    async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission:
        validate_submission(proposed)
        if proposed.confirmed_chunks:
            raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN)
        scope = proposed.scope
        key, identity = scope.storage_key, scope.canonical_identity.decode()
        plan, plan_digest, confirmed, confirmed_digest = encode_submission(proposed)
        for _ in range(_MAX_CAS_ATTEMPTS):
            snapshot = await self._snapshot(key, proposed.submission_id)
            current, previous = await asyncio.to_thread(self._validate_snapshot, scope, snapshot)
            if current is not None and current.pending:
                return current
            if previous is not None:
                if previous._plan != proposed._plan and not is_completed_legacy_replay(
                    previous, proposed
                ):
                    raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
                return previous
            async with self._pool.connection() as connection, connection.transaction():
                expected = snapshot
                if snapshot.scope is None:
                    await connection.execute(
                        "INSERT INTO reporting_buyer_submission_scopes"
                        " (scope_sha256,canonical_identity) VALUES (%s,%s)"
                        " ON CONFLICT (scope_sha256) DO NOTHING",
                        (key, identity),
                    )
                    expected = _Snapshot((identity, None), None, None, False)
                actual = await self._snapshot_on(connection, key, proposed.submission_id, lock=True)
                if actual != expected:
                    continue
                await connection.execute(
                    "INSERT INTO reporting_buyer_submission_intents"
                    " (scope_sha256,submission_id,canonical_plan,plan_sha256,"
                    " confirmed_results,confirmed_sha256,pending) VALUES (%s,%s,%s,%s,%s,%s,true)",
                    (
                        key,
                        proposed.submission_id,
                        plan,
                        plan_digest,
                        confirmed,
                        confirmed_digest,
                    ),
                )
                await connection.execute(
                    "UPDATE reporting_buyer_submission_scopes SET current_submission_id=%s"
                    " WHERE scope_sha256=%s",
                    (proposed.submission_id, key),
                )
            return proposed
        raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE)

    @_closed_storage_errors
    async def get(
        self, scope: ReportingSubmissionScope, submission_id: str | None = None
    ) -> ReportingReceiptSubmission | None:
        snapshot = await self._snapshot(scope.storage_key, submission_id)
        _, selected = await asyncio.to_thread(self._validate_snapshot, scope, snapshot)
        return selected

    @_closed_storage_errors
    async def confirm(
        self,
        scope: ReportingSubmissionScope,
        submission_id: str,
        chunk: int,
        response: SyncReportingReceiptsResponse,
    ) -> ReportingReceiptSubmission:
        # Caller-owned Pydantic objects are detached before the first await.
        # Retries compare the same immutable response, even if its owner edits it.
        frozen = freeze_response(response)
        key = scope.storage_key
        for _ in range(_MAX_CAS_ATTEMPTS):
            snapshot = await self._snapshot(key, submission_id)
            _, state = await asyncio.to_thread(self._validate_snapshot, scope, snapshot)
            if state is None:
                raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND)
            updated = await asyncio.to_thread(confirm_submission, state, chunk, frozen)
            _, _, confirmed, digest = encode_submission(updated)
            pending = updated.pending
            async with self._pool.connection() as connection, connection.transaction():
                actual = await self._snapshot_on(connection, key, submission_id, lock=True)
                if actual != snapshot:
                    continue
                if updated is not state:
                    await connection.execute(
                        "UPDATE reporting_buyer_submission_intents"
                        " SET confirmed_results=%s,confirmed_sha256=%s,pending=%s"
                        " WHERE scope_sha256=%s AND submission_id=%s",
                        (confirmed, digest, pending, key, submission_id),
                    )
            return updated
        # Contention never frees the lane or guesses whether another commit won.
        raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE)

Production persistence for optional buyer receipt submission intents.

Pass a trusted application-owned psycopg AsyncConnectionPool, preferably with a dedicated schema/search_path and bounded connection/statement waits. There is no credential or connection string in this object's representation. Run create_schema explicitly during rollout, then validate custom backends with the same memory/PostgreSQL state-machine vectors. Readiness requires all shipped, validated CHECK definitions as well as the keys and reservation index. Snapshot validation happens outside transactions; a mutation compares the complete locked state, retrying at most 16 times. Contention exhaustion raises STORAGE_UNAVAILABLE and leaves the durable intent available to resume.

Methods

async def confirm(self,
scope: ReportingSubmissionScope,
submission_id: str,
chunk: int,
response: SyncReportingReceiptsResponse) ‑> ReportingReceiptSubmission
Expand source code
@_closed_storage_errors
async def confirm(
    self,
    scope: ReportingSubmissionScope,
    submission_id: str,
    chunk: int,
    response: SyncReportingReceiptsResponse,
) -> ReportingReceiptSubmission:
    # Caller-owned Pydantic objects are detached before the first await.
    # Retries compare the same immutable response, even if its owner edits it.
    frozen = freeze_response(response)
    key = scope.storage_key
    for _ in range(_MAX_CAS_ATTEMPTS):
        snapshot = await self._snapshot(key, submission_id)
        _, state = await asyncio.to_thread(self._validate_snapshot, scope, snapshot)
        if state is None:
            raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND)
        updated = await asyncio.to_thread(confirm_submission, state, chunk, frozen)
        _, _, confirmed, digest = encode_submission(updated)
        pending = updated.pending
        async with self._pool.connection() as connection, connection.transaction():
            actual = await self._snapshot_on(connection, key, submission_id, lock=True)
            if actual != snapshot:
                continue
            if updated is not state:
                await connection.execute(
                    "UPDATE reporting_buyer_submission_intents"
                    " SET confirmed_results=%s,confirmed_sha256=%s,pending=%s"
                    " WHERE scope_sha256=%s AND submission_id=%s",
                    (confirmed, digest, pending, key, submission_id),
                )
        return updated
    # Contention never frees the lane or guesses whether another commit won.
    raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE)
async def create_schema(self) ‑> None
Expand source code
@_closed_storage_errors
async def create_schema(self) -> None:
    available = False
    try:
        import psycopg  # noqa: F401

        available = True
    except ImportError:
        # The optional driver is absent; report PG_REQUIRED outside the handler.
        pass
    if not available:
        raise ReportingSubmissionError(ReportingSubmissionCode.PG_REQUIRED)
    async with self._pool.connection() as connection, connection.transaction():
        await connection.execute("SELECT pg_advisory_xact_lock(%s)", (712071172,))
        sql = (
            files("adcp.reporting.ledger")
            .joinpath("reporting_buyer_submissions.sql")
            .read_text()
        )
        await connection.execute(sql)
        await self._check_schema_on(connection)
async def get(self,
scope: ReportingSubmissionScope,
submission_id: str | None = None) ‑> ReportingReceiptSubmission | None
Expand source code
@_closed_storage_errors
async def get(
    self, scope: ReportingSubmissionScope, submission_id: str | None = None
) -> ReportingReceiptSubmission | None:
    snapshot = await self._snapshot(scope.storage_key, submission_id)
    _, selected = await asyncio.to_thread(self._validate_snapshot, scope, snapshot)
    return selected
async def reserve(self,
proposed: ReportingReceiptSubmission) ‑> ReportingReceiptSubmission
Expand source code
@_closed_storage_errors
async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission:
    validate_submission(proposed)
    if proposed.confirmed_chunks:
        raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN)
    scope = proposed.scope
    key, identity = scope.storage_key, scope.canonical_identity.decode()
    plan, plan_digest, confirmed, confirmed_digest = encode_submission(proposed)
    for _ in range(_MAX_CAS_ATTEMPTS):
        snapshot = await self._snapshot(key, proposed.submission_id)
        current, previous = await asyncio.to_thread(self._validate_snapshot, scope, snapshot)
        if current is not None and current.pending:
            return current
        if previous is not None:
            if previous._plan != proposed._plan and not is_completed_legacy_replay(
                previous, proposed
            ):
                raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT)
            return previous
        async with self._pool.connection() as connection, connection.transaction():
            expected = snapshot
            if snapshot.scope is None:
                await connection.execute(
                    "INSERT INTO reporting_buyer_submission_scopes"
                    " (scope_sha256,canonical_identity) VALUES (%s,%s)"
                    " ON CONFLICT (scope_sha256) DO NOTHING",
                    (key, identity),
                )
                expected = _Snapshot((identity, None), None, None, False)
            actual = await self._snapshot_on(connection, key, proposed.submission_id, lock=True)
            if actual != expected:
                continue
            await connection.execute(
                "INSERT INTO reporting_buyer_submission_intents"
                " (scope_sha256,submission_id,canonical_plan,plan_sha256,"
                " confirmed_results,confirmed_sha256,pending) VALUES (%s,%s,%s,%s,%s,%s,true)",
                (
                    key,
                    proposed.submission_id,
                    plan,
                    plan_digest,
                    confirmed,
                    confirmed_digest,
                ),
            )
            await connection.execute(
                "UPDATE reporting_buyer_submission_scopes SET current_submission_id=%s"
                " WHERE scope_sha256=%s",
                (proposed.submission_id, key),
            )
        return proposed
    raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE)
class ReportingReceiptFailureCode (*args, **kwds)
Expand source code
class ReportingReceiptFailureCode(str, Enum):
    """Known item failures; unrecognized seller codes map to UNKNOWN.

    Seller messages/details are deliberately not persisted or exposed. Each
    failed item and each error position is retained, including unknown codes.
    """

    UNKNOWN = "UNKNOWN"
    INVALID_REQUEST = "INVALID_REQUEST"
    UNAUTHORIZED = "UNAUTHORIZED"
    RATE_LIMITED = "RATE_LIMITED"
    NOT_SUPPORTED = "NOT_SUPPORTED"
    IDEMPOTENCY_CONFLICT = "IDEMPOTENCY_CONFLICT"
    INVALID_REPORTING_RECORD = "INVALID_REPORTING_RECORD"
    REPORTING_RECORD_UNAVAILABLE = "REPORTING_RECORD_UNAVAILABLE"
    REPORTING_HISTORY_CORRUPT = "REPORTING_HISTORY_CORRUPT"
    REPORTING_IDENTITY_CONFLICT = "REPORTING_IDENTITY_CONFLICT"
    REPORTING_TIME_INVALID = "REPORTING_TIME_INVALID"
    RECEIPTS_NOT_ENABLED = "RECEIPTS_NOT_ENABLED"
    RECEIVED_AT_READ_ONLY = "RECEIVED_AT_READ_ONLY"
    RECEIPT_PROFILE_MISMATCH = "RECEIPT_PROFILE_MISMATCH"
    RECEIPT_TOTALS_MISMATCH = "RECEIPT_TOTALS_MISMATCH"
    RECEIPT_EVIDENCE_MISMATCH = "RECEIPT_EVIDENCE_MISMATCH"
    MATERIALIZATION_UNREADABLE = "MATERIALIZATION_UNREADABLE"
    ADJUSTMENT_REQUIRES_OFFICIAL = "ADJUSTMENT_REQUIRES_OFFICIAL"
    ADJUSTMENT_ORDER_INVALID = "ADJUSTMENT_ORDER_INVALID"
    ADJUSTMENT_DIGEST_MISMATCH = "ADJUSTMENT_DIGEST_MISMATCH"
    ACCEPTED_RECEIPT_TERMINAL = "ACCEPTED_RECEIPT_TERMINAL"

Known item failures; unrecognized seller codes map to UNKNOWN.

Seller messages/details are deliberately not persisted or exposed. Each failed item and each error position is retained, including unknown codes.

Ancestors

  • builtins.str
  • enum.Enum

Class variables

var ACCEPTED_RECEIPT_TERMINAL
var ADJUSTMENT_DIGEST_MISMATCH
var ADJUSTMENT_ORDER_INVALID
var ADJUSTMENT_REQUIRES_OFFICIAL
var IDEMPOTENCY_CONFLICT
var INVALID_REPORTING_RECORD
var INVALID_REQUEST
var MATERIALIZATION_UNREADABLE
var NOT_SUPPORTED
var RATE_LIMITED
var RECEIPTS_NOT_ENABLED
var RECEIPT_EVIDENCE_MISMATCH
var RECEIPT_PROFILE_MISMATCH
var RECEIPT_TOTALS_MISMATCH
var RECEIVED_AT_READ_ONLY
var REPORTING_HISTORY_CORRUPT
var REPORTING_IDENTITY_CONFLICT
var REPORTING_RECORD_UNAVAILABLE
var REPORTING_TIME_INVALID
var UNAUTHORIZED
var UNKNOWN
class ReportingReceiptOutcome (ordinal: int,
kind: ReceiptKind,
result: ReceiptResult,
error_codes: tuple[ReportingReceiptFailureCode, ...] = ())
Expand source code
@dataclass(frozen=True)
class ReportingReceiptOutcome:
    """One confirmed item, in original input order, with fresh model views."""

    ordinal: int
    kind: ReceiptKind
    result: ReceiptResult
    error_codes: tuple[ReportingReceiptFailureCode, ...] = ()
    _submitted: bytes = field(default=b"", repr=False)
    _stored: bytes | None = field(default=None, repr=False)

    @property
    def submitted_receipt(self) -> ReportingSubmissionReceipt:
        return _receipt(self.kind, json.loads(self._submitted))

    @property
    def receipt(self) -> ReportingSubmissionReceipt | None:
        """The seller's actual stored receipt on success, including received_at."""
        return _receipt(self.kind, json.loads(self._stored)) if self._stored is not None else None

One confirmed item, in original input order, with fresh model views.

Instance variables

var error_codes : tuple[ReportingReceiptFailureCode, ...]
var kind : Literal['receipt', 'adjustment_receipt']
var ordinal : int
prop receipt : ReportingSubmissionReceipt | None
Expand source code
@property
def receipt(self) -> ReportingSubmissionReceipt | None:
    """The seller's actual stored receipt on success, including received_at."""
    return _receipt(self.kind, json.loads(self._stored)) if self._stored is not None else None

The seller's actual stored receipt on success, including received_at.

var result : Literal['recorded', 'unchanged', 'failed']
prop submitted_receipt : ReportingSubmissionReceipt
Expand source code
@property
def submitted_receipt(self) -> ReportingSubmissionReceipt:
    return _receipt(self.kind, json.loads(self._submitted))
class ReportingReceiptSubmission (scope: ReportingSubmissionScope,
submission_id: str,
_plan: bytes)
Expand source code
@dataclass(frozen=True)
class ReportingReceiptSubmission:
    """Exact reserved plan and its immutable, fully confirmed chunk prefix.

    Construct with :func:`prepare_reporting_receipt_submission`. Stores validate
    persisted records before returning them. Private bytes keep mutable caller
    models, driver results, and representations out of the retained identity.
    """

    scope: ReportingSubmissionScope = field(repr=False)
    submission_id: str = field(repr=False)
    _plan: bytes = field(repr=False)
    _confirmed: tuple[bytes, ...] = field(default=(), repr=False)

    @property
    def chunk_count(self) -> int:
        return len(_validated_requests(self))

    @property
    def confirmed_chunks(self) -> int:
        return len(self._confirmed)

    @property
    def pending(self) -> bool:
        return self.confirmed_chunks < self.chunk_count

    @property
    def adcp_version(self) -> str:
        """The immutable request version, including for completed older plans."""
        return str(json.loads(_validated_requests(self)[0])["adcp_version"])

    def request(self, ordinal: int) -> SyncReportingReceiptsRequest:
        """Return a fresh typed view of one exact persisted request body."""
        body = json.loads(_validated_requests(self)[ordinal])
        return SyncReportingReceiptsRequest.model_validate(body)

    @property
    def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]:
        plan = json.loads(self._plan)
        confirmed = {
            item["reporting_receipt_id"]: item
            for chunk in self._confirmed
            for item in json.loads(chunk)
        }
        outcomes = []
        for ordinal, item in enumerate(plan["items"]):
            result = confirmed.get(item["body"]["reporting_receipt_id"])
            if result is not None:
                stored = result.get("receipt")
                outcomes.append(
                    ReportingReceiptOutcome(
                        ordinal,
                        item["kind"],
                        result["result"],
                        tuple(ReportingReceiptFailureCode(code) for code in result["errors"]),
                        canonical_json_utf8_v1(item["body"]),
                        canonical_json_utf8_v1(stored) if stored is not None else None,
                    )
                )
        return tuple(outcomes)

Exact reserved plan and its immutable, fully confirmed chunk prefix.

Construct with :func:prepare_reporting_receipt_submission(). Stores validate persisted records before returning them. Private bytes keep mutable caller models, driver results, and representations out of the retained identity.

Instance variables

prop adcp_version : str
Expand source code
@property
def adcp_version(self) -> str:
    """The immutable request version, including for completed older plans."""
    return str(json.loads(_validated_requests(self)[0])["adcp_version"])

The immutable request version, including for completed older plans.

prop chunk_count : int
Expand source code
@property
def chunk_count(self) -> int:
    return len(_validated_requests(self))
prop confirmed_chunks : int
Expand source code
@property
def confirmed_chunks(self) -> int:
    return len(self._confirmed)
prop outcomes : tuple[ReportingReceiptOutcome, ...]
Expand source code
@property
def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]:
    plan = json.loads(self._plan)
    confirmed = {
        item["reporting_receipt_id"]: item
        for chunk in self._confirmed
        for item in json.loads(chunk)
    }
    outcomes = []
    for ordinal, item in enumerate(plan["items"]):
        result = confirmed.get(item["body"]["reporting_receipt_id"])
        if result is not None:
            stored = result.get("receipt")
            outcomes.append(
                ReportingReceiptOutcome(
                    ordinal,
                    item["kind"],
                    result["result"],
                    tuple(ReportingReceiptFailureCode(code) for code in result["errors"]),
                    canonical_json_utf8_v1(item["body"]),
                    canonical_json_utf8_v1(stored) if stored is not None else None,
                )
            )
    return tuple(outcomes)
prop pending : bool
Expand source code
@property
def pending(self) -> bool:
    return self.confirmed_chunks < self.chunk_count
var scope : ReportingSubmissionScope
var submission_id : str

Methods

def request(self, ordinal: int) ‑> SyncReportingReceiptsRequest
Expand source code
def request(self, ordinal: int) -> SyncReportingReceiptsRequest:
    """Return a fresh typed view of one exact persisted request body."""
    body = json.loads(_validated_requests(self)[ordinal])
    return SyncReportingReceiptsRequest.model_validate(body)

Return a fresh typed view of one exact persisted request body.

class ReportingReceiptSubmissionClient (*args, **kwargs)
Expand source code
class ReportingReceiptSubmissionClient(Protocol):
    async def sync_reporting_receipts(
        self, request: SyncReportingReceiptsRequest
    ) -> TaskResult[SyncReportingReceiptsResponse]: ...

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 sync_reporting_receipts(self, request: SyncReportingReceiptsRequest) ‑> TaskResult[SyncReportingReceiptsResponse]
Expand source code
async def sync_reporting_receipts(
    self, request: SyncReportingReceiptsRequest
) -> TaskResult[SyncReportingReceiptsResponse]: ...
class ReportingSubmissionAuthorizer (*args, **kwargs)
Expand source code
class ReportingSubmissionAuthorizer(Protocol):
    """Trusted application adapter resolving the exact client's current access.

    Resolve seller identity from trusted client/registry configuration, account
    from an authorized seller account lookup, and consumer from authenticated
    credentials/registry. Never source these identities from receipt payloads,
    an asserted account reference, or a debug/raw response. Aliases must already
    be resolved. Raise on revoked access or an identity disagreement.

    Called before any store read/reservation, before each send, and before
    returning completed cached outcomes. It must authorize replay afresh too.
    """

    async def __call__(
        self, client: ReportingReceiptSubmissionClient
    ) -> ReportingSubmissionScope: ...

Trusted application adapter resolving the exact client's current access.

Resolve seller identity from trusted client/registry configuration, account from an authorized seller account lookup, and consumer from authenticated credentials/registry. Never source these identities from receipt payloads, an asserted account reference, or a debug/raw response. Aliases must already be resolved. Raise on revoked access or an identity disagreement.

Called before any store read/reservation, before each send, and before returning completed cached outcomes. It must authorize replay afresh too.

Ancestors

  • typing.Protocol
  • typing.Generic
class ReportingSubmissionCode (*args, **kwds)
Expand source code
class ReportingSubmissionCode(str, Enum):
    INVALID_SCOPE = "INVALID_SUBMISSION_SCOPE"
    INVALID_PLAN = "INVALID_SUBMISSION_PLAN"
    UNAUTHORIZED = "SUBMISSION_NOT_AUTHORIZED"
    NOT_FOUND = "SUBMISSION_NOT_FOUND"
    HISTORY_CORRUPT = "SUBMISSION_HISTORY_CORRUPT"
    STORAGE_UNAVAILABLE = "SUBMISSION_STORAGE_UNAVAILABLE"
    PG_REQUIRED = "SUBMISSION_PG_REQUIRED"
    INVALID_RESPONSE = "SUBMISSION_RESPONSE_INVALID"
    TRANSPORT_UNCERTAIN = "SUBMISSION_TRANSPORT_UNCERTAIN"
    RESPONSE_UNCONFIRMED = "SUBMISSION_RESPONSE_UNCONFIRMED"

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 HISTORY_CORRUPT
var INVALID_PLAN
var INVALID_RESPONSE
var INVALID_SCOPE
var NOT_FOUND
var PG_REQUIRED
var RESPONSE_UNCONFIRMED
var STORAGE_UNAVAILABLE
var TRANSPORT_UNCERTAIN
var UNAUTHORIZED
class ReportingSubmissionError (code: ReportingSubmissionCode)
Expand source code
class ReportingSubmissionError(RuntimeError):
    """Closed code only: no provider, authentication, SQL or wire body details."""

    def __init__(self, code: ReportingSubmissionCode) -> None:
        self.code = (
            code
            if isinstance(code, ReportingSubmissionCode)
            else ReportingSubmissionCode.INVALID_PLAN
        )
        super().__init__(self.code.value)

Closed code only: no provider, authentication, SQL or wire body details.

Ancestors

  • builtins.RuntimeError
  • builtins.Exception
  • builtins.BaseException
class ReportingSubmissionIntentStore (*args, **kwargs)
Expand source code
class ReportingSubmissionIntentStore(Protocol):
    """Atomic exact-request reservation and monotonic confirmed chunk storage.

    Adopters implementing this protocol must run the shared store vectors against
    their durable backend. A scope has one unresolved reservation, never a TTL,
    lease expiry, or automatic abandonment. Each mutation commits before return.
    No transport call may run while holding a storage transaction or lock.

    The supplied scope is already resolved/authorized by a trusted adapter. The
    store is an internal persistence boundary, not a request authentication API.
    Old ReportingCheckpointStore implementations need no new methods.
    """

    async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission:
        """Atomically reserve or return the scope's earlier unresolved intent.

        A different proposal must not replace an uncertain request. An already
        completed proposal for the same frozen receipt items returns its retained
        outcomes across the rc.6/rc.7 to 3.2 upgrade. Preserve all old request bytes,
        chunk keys/order and confirmed outcomes across restart.
        """
        ...

    async def get(
        self, scope: ReportingSubmissionScope, submission_id: str | None = None
    ) -> ReportingReceiptSubmission | None:
        """Read a validated reservation, including completed work after a lost return.

        With no ID return the latest reserved intent; do not erase it when it
        completes. Never resolve an ID outside the exact trusted scope.
        """
        ...

    async def confirm(
        self,
        scope: ReportingSubmissionScope,
        submission_id: str,
        chunk: int,
        response: SyncReportingReceiptsResponse,
    ) -> ReportingReceiptSubmission:
        """Atomically acknowledge full unique ID/body coverage for the next chunk.

        Retain item successes AND failures in the original plan order. Confirmed
        chunks never change. A concurrent duplicate acknowledgement returns the
        retained state; it cannot replace an earlier confirmed outcome.
        """
        ...

Atomic exact-request reservation and monotonic confirmed chunk storage.

Adopters implementing this protocol must run the shared store vectors against their durable backend. A scope has one unresolved reservation, never a TTL, lease expiry, or automatic abandonment. Each mutation commits before return. No transport call may run while holding a storage transaction or lock.

The supplied scope is already resolved/authorized by a trusted adapter. The store is an internal persistence boundary, not a request authentication API. Old ReportingCheckpointStore implementations need no new methods.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def confirm(self,
scope: ReportingSubmissionScope,
submission_id: str,
chunk: int,
response: SyncReportingReceiptsResponse) ‑> ReportingReceiptSubmission
Expand source code
async def confirm(
    self,
    scope: ReportingSubmissionScope,
    submission_id: str,
    chunk: int,
    response: SyncReportingReceiptsResponse,
) -> ReportingReceiptSubmission:
    """Atomically acknowledge full unique ID/body coverage for the next chunk.

    Retain item successes AND failures in the original plan order. Confirmed
    chunks never change. A concurrent duplicate acknowledgement returns the
    retained state; it cannot replace an earlier confirmed outcome.
    """
    ...

Atomically acknowledge full unique ID/body coverage for the next chunk.

Retain item successes AND failures in the original plan order. Confirmed chunks never change. A concurrent duplicate acknowledgement returns the retained state; it cannot replace an earlier confirmed outcome.

async def get(self,
scope: ReportingSubmissionScope,
submission_id: str | None = None) ‑> ReportingReceiptSubmission | None
Expand source code
async def get(
    self, scope: ReportingSubmissionScope, submission_id: str | None = None
) -> ReportingReceiptSubmission | None:
    """Read a validated reservation, including completed work after a lost return.

    With no ID return the latest reserved intent; do not erase it when it
    completes. Never resolve an ID outside the exact trusted scope.
    """
    ...

Read a validated reservation, including completed work after a lost return.

With no ID return the latest reserved intent; do not erase it when it completes. Never resolve an ID outside the exact trusted scope.

async def reserve(self,
proposed: ReportingReceiptSubmission) ‑> ReportingReceiptSubmission
Expand source code
async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission:
    """Atomically reserve or return the scope's earlier unresolved intent.

    A different proposal must not replace an uncertain request. An already
    completed proposal for the same frozen receipt items returns its retained
    outcomes across the rc.6/rc.7 to 3.2 upgrade. Preserve all old request bytes,
    chunk keys/order and confirmed outcomes across restart.
    """
    ...

Atomically reserve or return the scope's earlier unresolved intent.

A different proposal must not replace an uncertain request. An already completed proposal for the same frozen receipt items returns its retained outcomes across the rc.6/rc.7 to 3.2 upgrade. Preserve all old request bytes, chunk keys/order and confirmed outcomes across restart.

class ReportingSubmissionResult (submission: ReportingReceiptSubmission,
proposal_deferred: bool = False,
diagnostic: ReportingSubmissionCode | None = None)
Expand source code
@dataclass(frozen=True)
class ReportingSubmissionResult:
    """Submission outcomes only; completion does not establish reconciliation.

    ``proposal_deferred`` means an earlier pending scope reservation was resumed
    instead. Inspect its outcomes before planning any subsequent submission.
    """

    submission: ReportingReceiptSubmission
    proposal_deferred: bool = False
    diagnostic: ReportingSubmissionCode | None = None

    @property
    def pending(self) -> bool:
        return self.submission.pending

    @property
    def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]:
        return self.submission.outcomes

    @property
    def submitted_receipts(self) -> tuple[ReportingSubmissionReceipt, ...]:
        return tuple(
            receipt for outcome in self.outcomes if (receipt := outcome.receipt) is not None
        )

Submission outcomes only; completion does not establish reconciliation.

proposal_deferred means an earlier pending scope reservation was resumed instead. Inspect its outcomes before planning any subsequent submission.

Instance variables

var diagnostic : ReportingSubmissionCode | None
prop outcomes : tuple[ReportingReceiptOutcome, ...]
Expand source code
@property
def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]:
    return self.submission.outcomes
prop pending : bool
Expand source code
@property
def pending(self) -> bool:
    return self.submission.pending
var proposal_deferred : bool
var submission : ReportingReceiptSubmission
prop submitted_receipts : tuple[ReportingSubmissionReceipt, ...]
Expand source code
@property
def submitted_receipts(self) -> tuple[ReportingSubmissionReceipt, ...]:
    return tuple(
        receipt for outcome in self.outcomes if (receipt := outcome.receipt) is not None
    )
class ReportingSubmissionScope (seller_id: str, account_id: str, consumer_id: str)
Expand source code
@dataclass(frozen=True)
class ReportingSubmissionScope:
    """Identity resolved by a trusted authorization adapter, never from a request.

    ``seller_id`` identifies the configured seller, ``account_id`` is the
    seller-resolved account and ``consumer_id`` is the canonical authenticated
    principal. Syntax checks cannot establish provenance: the required
    authorizer must bind all three to the exact client and its current access.
    No credential or natural-key account assertion belongs in these fields.
    """

    seller_id: str = field(repr=False)
    account_id: str = field(repr=False)
    consumer_id: str = field(repr=False)

    def __post_init__(self) -> None:
        valid = False
        try:
            consumer_reference(self.seller_id)
            principal_reference(self.account_id)
            canonical_consumer(self.consumer_id)
            valid = all(value == value.strip() for value in (self.seller_id, self.account_id))
        except (ValueError, TypeError, RuntimeError):
            # Map provider validation details to the closed INVALID_SCOPE error below.
            pass
        if not valid:
            raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_SCOPE)

    @property
    def canonical_identity(self) -> bytes:
        return canonical_json_utf8_v1([self.seller_id, self.account_id, self.consumer_id])

    @property
    def storage_key(self) -> str:
        # Index a fixed-size digest, not a possibly 2048-character principal.
        # Stores also compare the full identity to fail closed on collisions.
        return hashlib.sha256(self.canonical_identity).hexdigest()

Identity resolved by a trusted authorization adapter, never from a request.

seller_id identifies the configured seller, account_id is the seller-resolved account and consumer_id is the canonical authenticated principal. Syntax checks cannot establish provenance: the required authorizer must bind all three to the exact client and its current access. No credential or natural-key account assertion belongs in these fields.

Instance variables

var account_id : str
prop canonical_identity : bytes
Expand source code
@property
def canonical_identity(self) -> bytes:
    return canonical_json_utf8_v1([self.seller_id, self.account_id, self.consumer_id])
var consumer_id : str
var seller_id : str
prop storage_key : str
Expand source code
@property
def storage_key(self) -> str:
    # Index a fixed-size digest, not a possibly 2048-character principal.
    # Stores also compare the full identity to fail closed on collisions.
    return hashlib.sha256(self.canonical_identity).hexdigest()