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
- InMemoryReportingMaterializerStore
- InMemoryReportingReconciliationStore
- InMemoryReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
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
- PgReportingMaterializerStore
- PgReportingReconciliationStore
- PgReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
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 : intvar as_of : datetime.datetimevar caller : ReportingDeliveryPrincipalvar core : ReportingStatusSnapshotvar reconciliation : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...]var reporting_receipt_id : strvar 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 checkSee PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::
class GenProto(Protocol[T]): def meth(self) -> T: ...Ancestors
- typing.Protocol
- typing.Generic
Methods
async def 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_accountis an application ACL, called even for a completed batch.buyer_agentsoptionally 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
- ADCPHandler
- abc.ABC
- typing.Generic
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_versionSelect 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
ADCPHandler:accept_proposalacquire_rightsactivate_signalbuild_creative_legacybuy_productscalibrate_contentcheck_governancecomply_test_controllercontext_matchcontrol_media_buycreate_collection_listcreate_content_standardscreate_media_buycreate_property_listdecline_proposalsdelete_collection_listdelete_property_listget_account_financialsget_adcp_capabilitiesget_brand_identityget_collection_listget_content_standardsget_creative_deliveryget_creative_featuresget_media_buy_artifactsget_media_buy_deliveryget_media_buysget_plan_audit_logsget_principalget_productsget_property_listget_rightsget_signalsget_task_statusidentity_matchlist_account_changeslist_accountslist_collection_listslist_content_standardslist_creative_formats_legacylist_creativeslist_productslist_property_listslist_taskslist_transformerslog_eventpreview_creative_legacyprovide_performance_feedbackrefine_proposalsreport_plan_adjustmentreport_plan_outcomereport_usagerequest_proposalssi_get_offeringsi_initiate_sessionsi_send_messagesi_terminate_sessionsync_accountssync_agent_notification_configssync_audiencessync_catalogssync_creativessync_event_sourcessync_governancesync_planssync_principalsync_reporting_receiptssync_reporting_statusupdate_collection_listupdate_content_standardsupdate_media_buyupdate_property_listupdate_rightsvalidate_content_deliveryvalidate_inputverify_brand_claimverify_brand_claims