Module adcp.reporting.receipts

Durable seller receipt ingress. Optional PG driver is imported lazily.

Sub-modules

adcp.reporting.receipts.capture

Immutable receipt dirty inputs; projection activation belongs to B2.4.

adcp.reporting.receipts.errors

Closed, payload-free receipt ingress failures.

adcp.reporting.receipts.handler

One authenticated composition path for direct MCP/A2A and hydrated contexts.

adcp.reporting.receipts.memory

Account-serialized receipt batches on the existing unconditional rollback model.

adcp.reporting.receipts.pg

Durable mixed-batch ordinals on the materializer's account-locked connection.

adcp.reporting.receipts.records

Pure preparation against exact, account-locked retained evidence.

adcp.reporting.receipts.schema

Isolated ingestion readiness; never widens an A/B/C/B1/B2.1 manifest.

adcp.reporting.receipts.store

Optional batch participant; no change to any required legacy store protocol.

adcp.reporting.receipts.transport

Request-local lossless receipt JSON, before transport/model normalization.

adcp.reporting.receipts.wire

Lossless raw batch admission and isolated transport-schema overlays …

Classes

class InMemoryReportingReceiptStore (**kwargs: Any)
Expand source code
class InMemoryReportingReceiptStore(InMemoryReportingMaterializerStore):
    """Reference participant. No production tier or notification readiness claim."""

    # Lazily initialized deliberately: first-use allocations must roll back too.
    _receipt_batches: dict[tuple[str, str, str], _BatchState]
    _receipt_boundaries: list[ReportingReceiptBoundary]
    _receipt_status_heads: dict[ReportingDeliveryPrincipal, int]

    def _receipt_batch(
        self, caller: ReportingDeliveryPrincipal, batch: ReceiptBatch
    ) -> _BatchState:
        if not hasattr(self, "_receipt_batches"):
            self._receipt_batches = {}
        key = caller.account_id, caller.consumer_id, batch.key
        state = self._receipt_batches.get(key)
        if state is None:
            state = _BatchState(batch, len(batch.items), self._clock())
            self._receipt_batches[key] = state
        elif state.batch.canonical_request != batch.canonical_request:
            raise ReportingReceiptError("IDEMPOTENCY_CONFLICT")
        if state.expected_count != len(batch.items) or len(state.results) > state.expected_count:
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        validate_receipt_results(list(state.results), batch)
        return state

    def _append_receipt_result(self, state: _BatchState, result: dict[str, Any]) -> None:
        state.results = (*state.results, deepcopy(result))

    def _save_receipt_response(self, state: _BatchState) -> None:
        response = inject_context(
            state.batch.request, {"status": "completed", "results": deepcopy(list(state.results))}
        )
        validate_receipt_response(response, state.batch)
        state.response = canonical_json_utf8_v1(response)

    def _assemble_receipt_response(self, state: _BatchState) -> dict[str, Any]:
        if state.response is None:
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        response: dict[str, Any] = json.loads(state.response)
        validate_receipt_response(response, state.batch)
        if response["results"] != list(state.results):
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        return response

    async def ingest_receipt_batch(
        self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal
    ) -> dict[str, Any]:
        batch = ReceiptBatch.parse(request)
        while True:
            # Each turn chooses the next *durable* ordinal under the account
            # lock. Concurrent/resumed callers cooperate; no process-local batch
            # lock, pending task, or TTL is a correctness dependency.
            async with self._mutation():
                state = self._receipt_batch(caller, batch)
                if state.response is not None:
                    return self._assemble_receipt_response(state)
                if len(state.results) == state.expected_count:
                    self._save_receipt_response(state)
                    return self._assemble_receipt_response(state)
                kind, item = batch.items[len(state.results)]
                candidate = None
                try:
                    revision_id = item.get(
                        "reporting_revision_id", item.get("adjusts_reporting_revision_id")
                    )
                    revision = (
                        self._revisions.get(revision_id) if isinstance(revision_id, str) else None
                    )
                    obligation_id = item.get("reporting_obligation_id")
                    if (
                        kind == "adjustment_receipt"
                        and revision is not None
                        and revision.account_id == caller.account_id
                    ):
                        obligation_id = revision.reporting_obligation_id
                    obligation = (
                        self._obligations.get(obligation_id)
                        if isinstance(obligation_id, str)
                        else None
                    )
                    candidate = receipt_record(kind, item, caller, obligation)
                    records = tuple(c.record for c in self._caller_changes(caller))
                    stored, added = prepare_receipt(
                        candidate,
                        records,
                        self._delivery_context(candidate),
                        tuple(
                            r
                            for r in self._revisions.values()
                            if obligation is not None
                            and r.account_id == caller.account_id
                            and r.reporting_obligation_id == obligation.reporting_obligation_id
                        ),
                        self._clock(),
                    )
                    result = success_result(stored, added)
                except LedgerConflictError as error:
                    result = failed_result(item, error)
                    added = False
                # Nothing in the semantic-error catch above mutates domain state.
                # Any insertion/capture/result failure below escapes and rolls
                # this entire ordinal back, with notifications both off and on.
                if added:
                    assert candidate is not None
                    stored, inserted = self._commit_record_unlocked(candidate)
                    assert inserted
                    result = success_result(stored, True)
                self._append_receipt_result(state, result)

    def _commit_record_unlocked(
        self, record: RecordT, *, notify: bool = True, dirty: bool = True
    ) -> tuple[RecordT, bool]:
        stored, added = super()._commit_record_unlocked(record, notify=notify, dirty=dirty)
        if added and isinstance(
            stored, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)
        ):
            self._capture_receipt(stored)
        return stored, added

    def _capture_receipt(self, receipt: ReportingReceiptRecord) -> None:
        assert receipt.received_at is not None
        caller = receipt.scope.principal
        if not hasattr(self, "_receipt_status_heads"):
            self._receipt_status_heads = {}
        if not hasattr(self, "_receipt_boundaries"):
            self._receipt_boundaries = []
        sequence = self._receipt_status_heads.get(caller, 0) + 1
        self._receipt_status_heads[caller] = sequence
        # Preserve cross-kind account order using B2.1's original capture clock.
        # No materializer work, retry decision, or readiness event is created.
        account_sequence = self._materializer_account_heads.get(caller.account_id, 0) + 1
        self._materializer_account_heads[caller.account_id] = account_sequence
        core = replace(settle_memory_snapshot(self, caller.account_id), as_of=receipt.received_at)
        self._receipt_boundaries.append(
            ReportingReceiptBoundary(
                caller,
                sequence,
                account_sequence,
                receipt.reporting_receipt_id,
                receipt.received_at,
                private_snapshot(core, caller),
                tuple(c.record for c in self._caller_changes(caller)),
            )
        )

    async def read_receipt_boundaries(
        self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
    ) -> tuple[ReportingReceiptBoundary, ...]:
        if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100:
            raise ReportingReceiptError("INVALID_REQUEST")
        async with self._lock:
            return tuple(
                b
                for b in getattr(self, "_receipt_boundaries", ())
                if b.caller == caller and b.sequence > after
            )[:limit]

Reference participant. No production tier or notification readiness claim.

Ancestors

Subclasses

Methods

async def ingest_receipt_batch(self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal) ‑> dict[str, typing.Any]
Expand source code
async def ingest_receipt_batch(
    self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal
) -> dict[str, Any]:
    batch = ReceiptBatch.parse(request)
    while True:
        # Each turn chooses the next *durable* ordinal under the account
        # lock. Concurrent/resumed callers cooperate; no process-local batch
        # lock, pending task, or TTL is a correctness dependency.
        async with self._mutation():
            state = self._receipt_batch(caller, batch)
            if state.response is not None:
                return self._assemble_receipt_response(state)
            if len(state.results) == state.expected_count:
                self._save_receipt_response(state)
                return self._assemble_receipt_response(state)
            kind, item = batch.items[len(state.results)]
            candidate = None
            try:
                revision_id = item.get(
                    "reporting_revision_id", item.get("adjusts_reporting_revision_id")
                )
                revision = (
                    self._revisions.get(revision_id) if isinstance(revision_id, str) else None
                )
                obligation_id = item.get("reporting_obligation_id")
                if (
                    kind == "adjustment_receipt"
                    and revision is not None
                    and revision.account_id == caller.account_id
                ):
                    obligation_id = revision.reporting_obligation_id
                obligation = (
                    self._obligations.get(obligation_id)
                    if isinstance(obligation_id, str)
                    else None
                )
                candidate = receipt_record(kind, item, caller, obligation)
                records = tuple(c.record for c in self._caller_changes(caller))
                stored, added = prepare_receipt(
                    candidate,
                    records,
                    self._delivery_context(candidate),
                    tuple(
                        r
                        for r in self._revisions.values()
                        if obligation is not None
                        and r.account_id == caller.account_id
                        and r.reporting_obligation_id == obligation.reporting_obligation_id
                    ),
                    self._clock(),
                )
                result = success_result(stored, added)
            except LedgerConflictError as error:
                result = failed_result(item, error)
                added = False
            # Nothing in the semantic-error catch above mutates domain state.
            # Any insertion/capture/result failure below escapes and rolls
            # this entire ordinal back, with notifications both off and on.
            if added:
                assert candidate is not None
                stored, inserted = self._commit_record_unlocked(candidate)
                assert inserted
                result = success_result(stored, True)
            self._append_receipt_result(state, result)
async def read_receipt_boundaries(self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100) ‑> tuple[ReportingReceiptBoundary, ...]
Expand source code
async def read_receipt_boundaries(
    self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
) -> tuple[ReportingReceiptBoundary, ...]:
    if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100:
        raise ReportingReceiptError("INVALID_REQUEST")
    async with self._lock:
        return tuple(
            b
            for b in getattr(self, "_receipt_boundaries", ())
            if b.caller == caller and b.sequence > after
        )[:limit]

Inherited members

class PgReportingReceiptStore (*,
pool: AsyncConnectionPool,
clock: Callable[[], datetime] | None = None,
notifications: bool = False)
Expand source code
class PgReportingReceiptStore(PgReportingMaterializerStore):
    """One composition: old public stores plus optional ingestion and capture.

    No network I/O, session lock, or second pool participates in an ordinal.
    All actual receipt timestamps come from PostgreSQL, including when an old
    public conformance clock was supplied to the inherited constructor.
    """

    @_storage_errors("store.create_schema")
    async def create_schema(self) -> None:
        async with self._connection() as connection, connection.transaction():
            await self._create_schema_on(connection)
            root = files("adcp.reporting.ledger")
            await connection.execute(root.joinpath("reporting_materializer.sql").read_text())
            await connection.execute(root.joinpath("reporting_receipt_ingestion.sql").read_text())

    @_storage_errors("store.receipt_ingestion_ready")
    async def receipt_ingestion_ready(self) -> bool:
        async with self._connection() as connection:
            await validate_receipt_schema(connection, notifications=self._notifications_enabled)
        return True

    async def _delivery_time_on(self, connection: Any, record: ReportingDeliveryRecord) -> datetime:
        if isinstance(record, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)):
            return await _now(connection)
        return await super()._delivery_time_on(connection, record)

    async def _commit_record_on(
        self, connection: Any, record: RecordT, *, notify: bool = True, dirty: bool = True
    ) -> tuple[RecordT, bool]:
        stored, added = await super()._commit_record_on(
            connection, record, notify=notify, dirty=dirty
        )
        if added and isinstance(
            stored, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)
        ):
            await self._capture_receipt_on(connection, stored)
        return stored, added

    async def _receipt_batch_on(
        self, connection: Any, caller: ReportingDeliveryPrincipal, batch: ReceiptBatch
    ) -> tuple[Any, ...]:
        key = caller.account_id, caller.consumer_id, batch.key
        row = await (
            await connection.execute(
                "SELECT canonical_request,request_sha256,expected_count,final_response,final_sha256"
                " FROM reporting_receipt_ingestion_batches"
                " WHERE account_id=%s AND consumer_id=%s AND idempotency_key=%s FOR UPDATE",
                key,
            )
        ).fetchone()
        if row is None:
            await connection.execute(
                "INSERT INTO reporting_receipt_ingestion_batches"
                " (account_id,consumer_id,idempotency_key,canonical_request,"
                " request_sha256,expected_count)"
                " VALUES (%s,%s,%s,%s,%s,%s)",
                (*key, batch.canonical_request.decode(), batch.digest, len(batch.items)),
            )
            return batch.canonical_request.decode(), batch.digest, len(batch.items), None, None
        if row[0] != batch.canonical_request.decode() or row[1] != batch.digest:
            raise ReportingReceiptError("IDEMPOTENCY_CONFLICT")
        if row[2] != len(batch.items):
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        return tuple(row)

    async def _receipt_results_on(
        self, connection: Any, caller: ReportingDeliveryPrincipal, batch: ReceiptBatch
    ) -> list[dict[str, Any]]:
        rows = await (
            await connection.execute(
                "SELECT ordinal,receipt_kind,reporting_receipt_id,result,content_sha256"
                " FROM reporting_receipt_ingestion_results"
                " WHERE account_id=%s AND consumer_id=%s AND idempotency_key=%s ORDER BY ordinal",
                (caller.account_id, caller.consumer_id, batch.key),
            )
        ).fetchall()
        items = batch.items
        if len(rows) > len(items) or any(
            r[0] != i
            or r[1] != items[i][0]
            or r[2] != items[i][1]["reporting_receipt_id"]
            or canonical_json_sha256_v1(r[3]) != r[4]
            for i, r in enumerate(rows)
        ):
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        results = [r[3] for r in rows]
        validate_receipt_results(results, batch)
        return results

    async def _prepare_receipt_on(
        self,
        connection: Any,
        caller: ReportingDeliveryPrincipal,
        kind: ReceiptKind,
        item: dict[str, Any],
    ) -> tuple[ReportingReceiptRecord, dict[str, Any], bool]:
        revision_id = item.get("reporting_revision_id", item.get("adjusts_reporting_revision_id"))
        revision = await (
            await connection.execute(
                # SDK-owned column list; every request value is bound below.
                f"SELECT {_REVISION_COLUMNS} FROM reporting_revisions"  # nosec B608
                " WHERE account_id=%s AND reporting_revision_id=%s",
                (caller.account_id, revision_id),
            )
        ).fetchone()
        obligation_id = item.get("reporting_obligation_id")
        if kind == "adjustment_receipt" and revision is not None:
            obligation_id = _revision_from_row(revision).reporting_obligation_id
        row = await (
            await connection.execute(
                # SDK-owned column list; every request value is bound below.
                f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations"  # nosec B608
                " WHERE account_id=%s AND reporting_obligation_id=%s",
                (caller.account_id, obligation_id),
            )
        ).fetchone()
        obligation = _obligation_from_row(row) if row else None
        candidate = receipt_record(kind, item, caller, obligation)
        records = await self._records(connection, caller)
        revisions = await (
            await connection.execute(
                # SDK-owned column list; every request value is bound below.
                f"SELECT {_REVISION_COLUMNS} FROM reporting_revisions"  # nosec B608
                " WHERE account_id=%s AND reporting_obligation_id=%s",
                (caller.account_id, obligation_id),
            )
        ).fetchall()
        stored, added = prepare_receipt(
            candidate,
            records,
            await self._delivery_context(connection, candidate),
            tuple(_revision_from_row(r) for r in revisions),
            await _now(connection),
        )
        return candidate, success_result(stored, added), added

    async def _insert_receipt_result_on(
        self,
        connection: Any,
        caller: ReportingDeliveryPrincipal,
        batch: ReceiptBatch,
        ordinal: int,
        result: dict[str, Any],
    ) -> None:
        kind, item = batch.items[ordinal]
        await connection.execute(
            "INSERT INTO reporting_receipt_ingestion_results"
            " (account_id,consumer_id,idempotency_key,ordinal,receipt_kind,reporting_receipt_id,"
            " receipt_record_id,result,content_sha256) VALUES (%s,%s,%s,%s,%s,%s,%s,%s::jsonb,%s)",
            (
                caller.account_id,
                caller.consumer_id,
                batch.key,
                ordinal,
                kind,
                item["reporting_receipt_id"],
                item["reporting_receipt_id"] if result["result"] != "failed" else None,
                _json(result),
                canonical_json_sha256_v1(result),
            ),
        )

    async def _save_receipt_response_on(
        self,
        connection: Any,
        caller: ReportingDeliveryPrincipal,
        batch: ReceiptBatch,
        results: list[dict[str, Any]],
    ) -> dict[str, Any]:
        response = inject_context(batch.request, {"status": "completed", "results": results})
        validate_receipt_response(response, batch)
        await connection.execute(
            "UPDATE reporting_receipt_ingestion_batches"
            " SET final_response=%s::jsonb,final_sha256=%s,finalized_at=clock_timestamp()"
            " WHERE account_id=%s AND consumer_id=%s AND idempotency_key=%s",
            (
                _json(response),
                canonical_json_sha256_v1(response),
                caller.account_id,
                caller.consumer_id,
                batch.key,
            ),
        )
        return response

    def _assemble_receipt_response(
        self,
        response: dict[str, Any],
        digest: str,
        results: list[dict[str, Any]],
        batch: ReceiptBatch,
    ) -> dict[str, Any]:
        validate_receipt_response(response, batch)
        if canonical_json_sha256_v1(response) != digest or response["results"] != results:
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        return cast(dict[str, Any], json.loads(_json(response)))

    @_storage_errors("store.ingest_receipt_batch")
    async def ingest_receipt_batch(
        self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal
    ) -> dict[str, Any]:
        batch = ReceiptBatch.parse(request)
        checked = False
        while True:
            async with self._connection() as connection, connection.transaction():
                # Account advisory lock precedes the batch row; no lock pool or
                # second connection can deadlock a size-one application pool.
                await self._lock_account(connection, caller.account_id)
                if not checked:
                    await validate_receipt_schema(
                        connection, notifications=self._notifications_enabled
                    )
                    checked = True
                row = await self._receipt_batch_on(connection, caller, batch)
                results = await self._receipt_results_on(connection, caller, batch)
                if row[3] is not None:
                    return self._assemble_receipt_response(row[3], row[4], results, batch)
                if len(results) == len(batch.items):
                    response = await self._save_receipt_response_on(
                        connection, caller, batch, results
                    )
                    return self._assemble_receipt_response(
                        response, canonical_json_sha256_v1(response), results, batch
                    )
                kind, item = batch.items[len(results)]
                candidate = None
                try:
                    candidate, result, added = await self._prepare_receipt_on(
                        connection, caller, kind, item
                    )
                except LedgerConflictError as error:
                    result, added = failed_result(item, error), False
                # Only pure semantic preparation is caught. SQL/capture/result
                # failures roll back this ordinal instead of recording a failure
                # beside a partially inserted receipt.
                if added:
                    assert candidate is not None
                    stored, inserted = await self._commit_record_on(connection, candidate)
                    assert inserted
                    result = success_result(stored, True)
                await self._insert_receipt_result_on(
                    connection, caller, batch, len(results), result
                )

    async def _capture_receipt_on(self, connection: Any, receipt: ReportingReceiptRecord) -> None:
        assert receipt.received_at is not None
        caller = receipt.scope.principal
        core = await settle_snapshot_on(
            self, connection, account_id=caller.account_id, as_of=receipt.received_at
        )
        row = await (
            await connection.execute(
                "INSERT INTO reporting_receipt_ingestion_heads"
                " (account_id,consumer_id,max_sequence)"
                " VALUES (%s,%s,1) ON CONFLICT (account_id,consumer_id) DO UPDATE"
                " SET max_sequence=reporting_receipt_ingestion_heads.max_sequence+1"
                " RETURNING max_sequence",
                (caller.account_id, caller.consumer_id),
            )
        ).fetchone()
        account = await (
            await connection.execute(
                "UPDATE reporting_materializer_accounts SET captured_sequence=captured_sequence+1"
                " WHERE account_id=%s RETURNING captured_sequence",
                (caller.account_id,),
            )
        ).fetchone()
        if row is None or account is None:
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        value = ReportingReceiptBoundary(
            caller,
            row[0],
            account[0],
            receipt.reporting_receipt_id,
            receipt.received_at,
            private_snapshot(core, caller),
            await self._records(connection, caller),
        ).to_storage()
        await connection.execute(
            "INSERT INTO reporting_receipt_ingestion_boundaries"
            " (account_id,consumer_id,sequence,account_sequence,reporting_receipt_id,"
            " as_of,input,content_sha256)"
            " VALUES (%s,%s,%s,%s,%s,%s,%s::jsonb,%s)",
            (
                caller.account_id,
                caller.consumer_id,
                row[0],
                account[0],
                receipt.reporting_receipt_id,
                receipt.received_at,
                _json(value),
                canonical_json_sha256_v1(value),
            ),
        )

    @_storage_errors("store.read_receipt_boundaries")
    async def read_receipt_boundaries(
        self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
    ) -> tuple[ReportingReceiptBoundary, ...]:
        if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100:
            raise ReportingReceiptError("INVALID_REQUEST")
        async with self._connection() as connection, connection.transaction():
            await self._lock_account(connection, caller.account_id)
            await validate_receipt_schema(connection, notifications=self._notifications_enabled)
            rows = await (
                await connection.execute(
                    "SELECT input,content_sha256=reporting_receipt_ingestion_sha256(input)"
                    " FROM reporting_receipt_ingestion_boundaries"
                    " WHERE account_id=%s AND consumer_id=%s AND sequence>%s"
                    " ORDER BY sequence LIMIT %s",
                    (caller.account_id, caller.consumer_id, after, limit),
                )
            ).fetchall()
            if any(not r[1] for r in rows):
                raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
            return tuple(decode_receipt_boundary(r[0]) for r in rows)

One composition: old public stores plus optional ingestion and capture.

No network I/O, session lock, or second pool participates in an ordinal. All actual receipt timestamps come from PostgreSQL, including when an old public conformance clock was supplied to the inherited constructor.

Ancestors

Subclasses

Methods

async def ingest_receipt_batch(self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal) ‑> dict[str, typing.Any]
Expand source code
@_storage_errors("store.ingest_receipt_batch")
async def ingest_receipt_batch(
    self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal
) -> dict[str, Any]:
    batch = ReceiptBatch.parse(request)
    checked = False
    while True:
        async with self._connection() as connection, connection.transaction():
            # Account advisory lock precedes the batch row; no lock pool or
            # second connection can deadlock a size-one application pool.
            await self._lock_account(connection, caller.account_id)
            if not checked:
                await validate_receipt_schema(
                    connection, notifications=self._notifications_enabled
                )
                checked = True
            row = await self._receipt_batch_on(connection, caller, batch)
            results = await self._receipt_results_on(connection, caller, batch)
            if row[3] is not None:
                return self._assemble_receipt_response(row[3], row[4], results, batch)
            if len(results) == len(batch.items):
                response = await self._save_receipt_response_on(
                    connection, caller, batch, results
                )
                return self._assemble_receipt_response(
                    response, canonical_json_sha256_v1(response), results, batch
                )
            kind, item = batch.items[len(results)]
            candidate = None
            try:
                candidate, result, added = await self._prepare_receipt_on(
                    connection, caller, kind, item
                )
            except LedgerConflictError as error:
                result, added = failed_result(item, error), False
            # Only pure semantic preparation is caught. SQL/capture/result
            # failures roll back this ordinal instead of recording a failure
            # beside a partially inserted receipt.
            if added:
                assert candidate is not None
                stored, inserted = await self._commit_record_on(connection, candidate)
                assert inserted
                result = success_result(stored, True)
            await self._insert_receipt_result_on(
                connection, caller, batch, len(results), result
            )
async def read_receipt_boundaries(self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100) ‑> tuple[ReportingReceiptBoundary, ...]
Expand source code
@_storage_errors("store.read_receipt_boundaries")
async def read_receipt_boundaries(
    self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
) -> tuple[ReportingReceiptBoundary, ...]:
    if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100:
        raise ReportingReceiptError("INVALID_REQUEST")
    async with self._connection() as connection, connection.transaction():
        await self._lock_account(connection, caller.account_id)
        await validate_receipt_schema(connection, notifications=self._notifications_enabled)
        rows = await (
            await connection.execute(
                "SELECT input,content_sha256=reporting_receipt_ingestion_sha256(input)"
                " FROM reporting_receipt_ingestion_boundaries"
                " WHERE account_id=%s AND consumer_id=%s AND sequence>%s"
                " ORDER BY sequence LIMIT %s",
                (caller.account_id, caller.consumer_id, after, limit),
            )
        ).fetchall()
        if any(not r[1] for r in rows):
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        return tuple(decode_receipt_boundary(r[0]) for r in rows)
async def receipt_ingestion_ready(self) ‑> bool
Expand source code
@_storage_errors("store.receipt_ingestion_ready")
async def receipt_ingestion_ready(self) -> bool:
    async with self._connection() as connection:
        await validate_receipt_schema(connection, notifications=self._notifications_enabled)
    return True

Inherited members

class ReportingReceiptBatchStore (*args, **kwargs)
Expand source code
@runtime_checkable
class ReportingReceiptBatchStore(Protocol):
    """Trusted callers only. Each ordinal co-commits with all of its domain effects."""

    async def ingest_receipt_batch(
        self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal
    ) -> dict[str, Any]: ...

Trusted callers only. Each ordinal co-commits with all of its domain effects.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def ingest_receipt_batch(self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal) ‑> dict[str, typing.Any]
Expand source code
async def ingest_receipt_batch(
    self, request: dict[str, Any], *, caller: ReportingDeliveryPrincipal
) -> dict[str, Any]: ...
class ReportingReceiptBoundary (caller: ReportingDeliveryPrincipal,
sequence: int,
account_sequence: int,
reporting_receipt_id: str,
as_of: datetime,
core: ReportingStatusSnapshot,
reconciliation: tuple[ReportingDeliveryRecord, ...])
Expand source code
@dataclass(frozen=True)
class ReportingReceiptBoundary:
    caller: ReportingDeliveryPrincipal
    sequence: int
    account_sequence: int
    reporting_receipt_id: str
    as_of: datetime
    core: ReportingStatusSnapshot = field(repr=False)
    reconciliation: tuple[ReportingDeliveryRecord, ...] = field(repr=False)

    def __post_init__(self) -> None:
        receipts = tuple(
            r
            for r in self.reconciliation
            if isinstance(r, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord))
            and r.reporting_receipt_id == self.reporting_receipt_id
        )
        if (
            type(self.sequence) is not int
            or self.sequence < 1
            or type(self.account_sequence) is not int
            or self.account_sequence < self.sequence
            or self.core.account_id != self.caller.account_id
            or self.core.as_of != self.as_of
            or self.core.consumer_ids != (self.caller.consumer_id,)
            or any(s.consumer_id != self.caller.consumer_id for s in self.core.statuses)
            or any(
                i.consumer_id not in {None, self.caller.consumer_id} for i in self.core.lifecycles
            )
            or any(principal(r) != self.caller for r in self.reconciliation)
            or len(receipts) != 1
            or receipts[0].received_at != self.as_of
        ):
            raise ReportingReceiptError("RECEIPT_HISTORY_CORRUPT")
        object.__setattr__(self, "as_of", aware_utc(self.as_of))

    def to_storage(self) -> dict[str, Any]:
        return {
            "version": 1,
            "admission_epoch": 0,
            "account_id": self.caller.account_id,
            "consumer_id": self.caller.consumer_id,
            "sequence": self.sequence,
            "account_sequence": self.account_sequence,
            "reporting_receipt_id": self.reporting_receipt_id,
            "as_of": self.as_of.isoformat(),
            "core": _core_adapter().dump_python(self.core, mode="json"),
            "reconciliation": [payload(r) for r in self.reconciliation],
        }

ReportingReceiptBoundary(caller: 'ReportingDeliveryPrincipal', sequence: 'int', account_sequence: 'int', reporting_receipt_id: 'str', as_of: 'datetime', core: 'ReportingStatusSnapshot', reconciliation: 'tuple[ReportingDeliveryRecord, …]')

Instance variables

var account_sequence : int
var as_of : datetime.datetime
var caller : ReportingDeliveryPrincipal
var core : ReportingStatusSnapshot
var reconciliation : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...]
var reporting_receipt_id : str
var sequence : int

Methods

def to_storage(self) ‑> dict[str, typing.Any]
Expand source code
def to_storage(self) -> dict[str, Any]:
    return {
        "version": 1,
        "admission_epoch": 0,
        "account_id": self.caller.account_id,
        "consumer_id": self.caller.consumer_id,
        "sequence": self.sequence,
        "account_sequence": self.account_sequence,
        "reporting_receipt_id": self.reporting_receipt_id,
        "as_of": self.as_of.isoformat(),
        "core": _core_adapter().dump_python(self.core, mode="json"),
        "reconciliation": [payload(r) for r in self.reconciliation],
    }
class ReportingReceiptCaptureStore (*args, **kwargs)
Expand source code
@runtime_checkable
class ReportingReceiptCaptureStore(Protocol):
    async def read_receipt_boundaries(
        self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
    ) -> tuple[ReportingReceiptBoundary, ...]: ...

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 read_receipt_boundaries(self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100) ‑> tuple[ReportingReceiptBoundary, ...]
Expand source code
async def read_receipt_boundaries(
    self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
) -> tuple[ReportingReceiptBoundary, ...]: ...
class ReportingReceiptError (code: ReceiptErrorCode)
Expand source code
class ReportingReceiptError(RuntimeError):
    """An actionable classification; never includes request, driver or auth details."""

    def __init__(self, code: ReceiptErrorCode) -> None:
        self.code = code
        super().__init__(_MESSAGES[code])

An actionable classification; never includes request, driver or auth details.

Ancestors

  • builtins.RuntimeError
  • builtins.Exception
  • builtins.BaseException
class ReportingReceiptHandler (store: ReportingReceiptBatchStore,
*,
resolve_account: ReceiptAccountResolver,
buyer_agents: BuyerAgentRegistry | None = None,
consumer_status_enabled: bool = False,
adcp_version: str | None = None)
Expand source code
class ReportingReceiptHandler(ADCPHandler[ToolContext]):
    """Mount receipts and the optional frozen feed; tier activation is separately gated.

    ``resolve_account`` is an application ACL, called even for a completed
    batch. ``buyer_agents`` optionally re-resolves API/OAuth/signed commercial
    identity on every call. Neither transport tenancy nor body fields supply
    the consumer. Authentication middleware must populate trusted context.
    """

    advertised_tools = {TASK, "get_reporting_status"}

    def __init__(
        self,
        store: ReportingReceiptBatchStore,
        *,
        resolve_account: ReceiptAccountResolver,
        buyer_agents: BuyerAgentRegistry | None = None,
        consumer_status_enabled: bool = False,
        adcp_version: str | None = None,
    ) -> None:
        super().__init__()
        if not isinstance(store, ReportingReceiptBatchStore):
            raise TypeError("receipt ingress requires an atomic ReportingReceiptBatchStore")
        self.receipt_store = store
        self._receipt_account_resolver = resolve_account
        self._receipt_registry = buyer_agents
        from adcp.reporting.feed.store import ReportingFeedStore

        self.reporting_feed_store: ReportingFeedStore | None = (
            store if isinstance(store, ReportingFeedStore) else None
        )
        self._feed_consumer_status_enabled = consumer_status_enabled
        self._adcp_version = resolve_adcp_version(adcp_version)
        if not is_adcp_version_at_least(self._adcp_version, "3.2-rc.6"):
            raise ConfigurationError(
                "reporting mounts require a supported AdCP 3.2 reporting contract; "
                "use 3.2 or omit the pin for the packaged default"
            )

    def get_adcp_version(self) -> str:
        """Select rendering and advertised MCP/A2A schemas for this mount."""
        return self._adcp_version

    def advertised_tools_for_instance(self) -> set[str]:
        return {TASK, "get_reporting_status"} if self.reporting_feed_store is not None else {TASK}

    async def get_reporting_status(
        self,
        params: GetReportingStatusRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any] | NotImplementedResponse:
        """One authenticated mount; the optional store freezes periods walks."""
        from adcp.reporting.feed.errors import ReportingFeedError
        from adcp.reporting.feed.request import FeedRequest
        from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
        from adcp.reporting.ledger.store import LedgerConflictError, ReportingLedgerStore

        if self.reporting_feed_store is None:
            return self._not_supported("get_reporting_status")
        request = (
            dict(params)
            if isinstance(params, dict)
            else params.model_dump(mode="json", exclude_unset=True)
        )
        if context is not None and context.resolved_adcp_version is not None:
            request["adcp_version"] = context.resolved_adcp_version
        else:
            request.setdefault("adcp_version", self.get_adcp_version())
        try:
            if request.get("view") == "periods":
                FeedRequest.parse(request)

            async def authorize() -> ReportingDeliveryPrincipal:
                if context is None:
                    raise ReportingFeedError("UNAUTHORIZED")
                try:
                    consumer = await _consumer(context, self._receipt_registry)
                    account = await self._receipt_account_resolver(
                        dict(request["account"]), context, consumer
                    )
                    if isinstance(context, RequestContext) and context.account.id != account:
                        raise ReportingFeedError("UNAUTHORIZED")
                    return ReportingDeliveryPrincipal(account, consumer)
                except (ReportingReceiptError, ReportingNotificationError, ValueError, TypeError):
                    raise ReportingFeedError("UNAUTHORIZED") from None

            caller = await authorize()

            async def reauthorize() -> None:
                # A still-authorized alias must not change which account or
                # canonical consumer owns the already captured boundary.
                if await authorize() != caller:
                    raise ReportingFeedError("UNAUTHORIZED")

            if request.get("view") == "periods":
                response = await self.reporting_feed_store.read_reporting_feed(
                    request,
                    caller=caller,
                    consumer_status_enabled=self._feed_consumer_status_enabled,
                    reauthorize=reauthorize,
                )
                await reauthorize()
                return response
            pagination = request.get("pagination")
            positions = (
                request.get("changes_after"),
                pagination.get("cursor") if isinstance(pagination, dict) else None,
            )
            if any(isinstance(p, str) and p.startswith("rpf1.") for p in positions):
                # A view change cannot send a frozen position to the mutable
                # legacy projector. Leave ordinary Core requests compatible.
                raise ReportingFeedError("INVALID_CHECKPOINT")
            if isinstance(self.receipt_store, ReportingLedgerStore):
                try:
                    from adcp.reporting.projection.wire import ReportingTierStatusStore

                    if isinstance(self.receipt_store, ReportingTierStatusStore):
                        response = await self.receipt_store.read_tier_status(
                            request,
                            caller=caller,
                            consumer_status_enabled=self._feed_consumer_status_enabled,
                        )
                    else:
                        response = await ReportingStatusHandler(
                            self.receipt_store,
                            consumer_status_enabled=self._feed_consumer_status_enabled,
                        ).handle(
                            request,
                            caller=ReportingStatusCaller(caller.account_id, caller.consumer_id),
                        )
                    await reauthorize()
                    return response
                except LedgerConflictError as error:
                    # Only the legacy Core projector exposes its established
                    # domain errors. ACL/provider/feed failures stay redacted.
                    code, message = error.code, str(error)
            else:
                return self._not_supported("get_reporting_status")
        except ReportingFeedError as error:
            code, message = error.code, str(error)
        except Exception:
            unavailable = ReportingFeedError("REPORTING_FEED_STORAGE_UNAVAILABLE")
            code, message = unavailable.code, str(unavailable)
        raise ADCPTaskError(
            operation="get_reporting_status", errors=[Error(code=code, message=message)]
        )

    async def sync_reporting_receipts(
        self,
        params: SyncReportingReceiptsRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any]:
        request = (
            params
            if isinstance(params, dict)
            else params.model_dump(mode="json", exclude_unset=True)
        )
        try:
            validate_receipt_request(request)
            if context is None:
                raise ReportingReceiptError("UNAUTHORIZED")
            try:
                consumer = await _consumer(context, self._receipt_registry)
            except ReportingNotificationError:
                raise ReportingReceiptError("UNAUTHORIZED") from None
            account = await self._receipt_account_resolver(
                dict(request["account"]), context, consumer
            )
            if isinstance(context, RequestContext) and context.account.id != account:
                raise ReportingReceiptError("UNAUTHORIZED")
            try:
                caller = ReportingDeliveryPrincipal(account, consumer)
            except (ValueError, TypeError):
                raise ReportingReceiptError("UNAUTHORIZED") from None
            return await self.receipt_store.ingest_receipt_batch(request, caller=caller)
        except ReportingReceiptError as error:
            code, message = error.code, str(error)
        except Exception as error:
            _storage_failure(error, boundary="handler")
            unavailable = ReportingReceiptError("RECEIPT_STORAGE_UNAVAILABLE")
            code, message = unavailable.code, str(unavailable)
        # Leave the exception scope before translating. Credential/ACL adapters
        # may raise provider errors; only safe origin coordinates were logged.
        raise ADCPTaskError(operation=TASK, errors=[Error(code=code, message=message)])

Mount receipts and the optional frozen feed; tier activation is separately gated.

resolve_account is an application ACL, called even for a completed batch. buyer_agents optionally re-resolves API/OAuth/signed commercial identity on every call. Neither transport tenancy nor body fields supply the consumer. Authentication middleware must populate trusted context.

Ancestors

Subclasses

Class variables

var advertised_tools

Methods

def advertised_tools_for_instance(self) ‑> set[str]
Expand source code
def advertised_tools_for_instance(self) -> set[str]:
    return {TASK, "get_reporting_status"} if self.reporting_feed_store is not None else {TASK}
def get_adcp_version(self) ‑> str
Expand source code
def get_adcp_version(self) -> str:
    """Select rendering and advertised MCP/A2A schemas for this mount."""
    return self._adcp_version

Select rendering and advertised MCP/A2A schemas for this mount.

async def get_reporting_status(self,
params: GetReportingStatusRequest | dict[str, Any],
context: ToolContext | None = None) ‑> dict[str, typing.Any] | NotImplementedResponse
Expand source code
async def get_reporting_status(
    self,
    params: GetReportingStatusRequest | dict[str, Any],
    context: ToolContext | None = None,
) -> dict[str, Any] | NotImplementedResponse:
    """One authenticated mount; the optional store freezes periods walks."""
    from adcp.reporting.feed.errors import ReportingFeedError
    from adcp.reporting.feed.request import FeedRequest
    from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
    from adcp.reporting.ledger.store import LedgerConflictError, ReportingLedgerStore

    if self.reporting_feed_store is None:
        return self._not_supported("get_reporting_status")
    request = (
        dict(params)
        if isinstance(params, dict)
        else params.model_dump(mode="json", exclude_unset=True)
    )
    if context is not None and context.resolved_adcp_version is not None:
        request["adcp_version"] = context.resolved_adcp_version
    else:
        request.setdefault("adcp_version", self.get_adcp_version())
    try:
        if request.get("view") == "periods":
            FeedRequest.parse(request)

        async def authorize() -> ReportingDeliveryPrincipal:
            if context is None:
                raise ReportingFeedError("UNAUTHORIZED")
            try:
                consumer = await _consumer(context, self._receipt_registry)
                account = await self._receipt_account_resolver(
                    dict(request["account"]), context, consumer
                )
                if isinstance(context, RequestContext) and context.account.id != account:
                    raise ReportingFeedError("UNAUTHORIZED")
                return ReportingDeliveryPrincipal(account, consumer)
            except (ReportingReceiptError, ReportingNotificationError, ValueError, TypeError):
                raise ReportingFeedError("UNAUTHORIZED") from None

        caller = await authorize()

        async def reauthorize() -> None:
            # A still-authorized alias must not change which account or
            # canonical consumer owns the already captured boundary.
            if await authorize() != caller:
                raise ReportingFeedError("UNAUTHORIZED")

        if request.get("view") == "periods":
            response = await self.reporting_feed_store.read_reporting_feed(
                request,
                caller=caller,
                consumer_status_enabled=self._feed_consumer_status_enabled,
                reauthorize=reauthorize,
            )
            await reauthorize()
            return response
        pagination = request.get("pagination")
        positions = (
            request.get("changes_after"),
            pagination.get("cursor") if isinstance(pagination, dict) else None,
        )
        if any(isinstance(p, str) and p.startswith("rpf1.") for p in positions):
            # A view change cannot send a frozen position to the mutable
            # legacy projector. Leave ordinary Core requests compatible.
            raise ReportingFeedError("INVALID_CHECKPOINT")
        if isinstance(self.receipt_store, ReportingLedgerStore):
            try:
                from adcp.reporting.projection.wire import ReportingTierStatusStore

                if isinstance(self.receipt_store, ReportingTierStatusStore):
                    response = await self.receipt_store.read_tier_status(
                        request,
                        caller=caller,
                        consumer_status_enabled=self._feed_consumer_status_enabled,
                    )
                else:
                    response = await ReportingStatusHandler(
                        self.receipt_store,
                        consumer_status_enabled=self._feed_consumer_status_enabled,
                    ).handle(
                        request,
                        caller=ReportingStatusCaller(caller.account_id, caller.consumer_id),
                    )
                await reauthorize()
                return response
            except LedgerConflictError as error:
                # Only the legacy Core projector exposes its established
                # domain errors. ACL/provider/feed failures stay redacted.
                code, message = error.code, str(error)
        else:
            return self._not_supported("get_reporting_status")
    except ReportingFeedError as error:
        code, message = error.code, str(error)
    except Exception:
        unavailable = ReportingFeedError("REPORTING_FEED_STORAGE_UNAVAILABLE")
        code, message = unavailable.code, str(unavailable)
    raise ADCPTaskError(
        operation="get_reporting_status", errors=[Error(code=code, message=message)]
    )

One authenticated mount; the optional store freezes periods walks.

Inherited members