Module adcp.reporting.submissions
Buyer-only durable receipt submission, separate from reconciliation planning.
Public imports are additive and work without the optional PostgreSQL driver. The application supplies a trusted authorizer, a receipt client, and an explicit intent store. Use PgReportingSubmissionIntentStore for restart durability; the memory implementation is a volatile reference for tests. No checkpoint, seller service, consumer-status, or client.reporting facade behavior is changed.
Rollout: use a dedicated buyer schema/pool; explicitly call create_schema before enabling submissions. The migration adds only reporting_buyer_submission_*. Retain pending intents indefinitely and resume them after any uncertain result. Do not drop the tables on rollback or replace pending plans with new receipt IDs. Disable new planning while recovering, then resume the same stored scope with current authorization. Exact final-head review precedes facade integration.
Sub-modules
adcp.reporting.submissions.models-
Immutable buyer intent records and closed diagnostics …
adcp.reporting.submissions.pg-
PostgreSQL buyer intent storage with short transactions and no network locks …
adcp.reporting.submissions.store-
Optional buyer intent protocol, independent of ReportingCheckpointStore.
adcp.reporting.submissions.submit-
Explicit low-level receipt submission with durable uncertainty semantics.
Functions
def prepare_reporting_receipt_submission(scope: ReportingSubmissionScope,
receipts: Sequence[ReportingSubmissionReceipt]) ‑> ReportingReceiptSubmission-
Expand source code
def prepare_reporting_receipt_submission( scope: ReportingSubmissionScope, receipts: Sequence[ReportingSubmissionReceipt], ) -> ReportingReceiptSubmission: """Freeze an already validated receipt plan, with at most 100 items per call. Outbound typed values are normalized once, then retained byte-for-byte. This is not a capture or proof of raw inbound adjustment evidence. Requests contain the authorized resolved account only, with no caller-provided account, idempotency key, context, extensions or credentials. """ return _prepare_for_version(scope, receipts, _VERSION)Freeze an already validated receipt plan, with at most 100 items per call.
Outbound typed values are normalized once, then retained byte-for-byte. This is not a capture or proof of raw inbound adjustment evidence. Requests contain the authorized resolved account only, with no caller-provided account, idempotency key, context, extensions or credentials.
async def submit_reporting_receipts(client: ReportingReceiptSubmissionClient,
*,
authorizer: ReportingSubmissionAuthorizer,
store: ReportingSubmissionIntentStore,
receipts: Sequence[ReportingSubmissionReceipt] | None = None,
timeout_seconds: float = 30.0) ‑> ReportingSubmissionResult-
Expand source code
async def submit_reporting_receipts( client: ReportingReceiptSubmissionClient, *, authorizer: ReportingSubmissionAuthorizer, store: ReportingSubmissionIntentStore, receipts: Sequence[ReportingSubmissionReceipt] | None = None, timeout_seconds: float = 30.0, ) -> ReportingSubmissionResult: """Reserve or resume a receipt plan without replacing uncertain work. Omit ``receipts`` to resume the latest intent for the freshly authorized scope, including its completed outcomes after a lost return. Otherwise the supplied plan is frozen and atomically reserved before ANY network write. An earlier pending scope intent takes priority, with proposal_deferred=True. The caller must reconsider the deferred plan against fresh seller history. A transport failure, timeout, failed task or malformed response leaves the chunk uncertain and returns pending=True. Retry within the live protocol version uses the exact durable body and idempotency key. A pending intent from a retired version remains readable and byte-identical but is refused before transport; it is never rewritten as a new-version request. Cancellation propagates with the intent still retained. Storage failure raises a closed error; resume the stored scope because the last commit may have succeeded. There is no expiry/abandon/replacement path. Every item failure is retained alongside successes. Completion only means every submission outcome is confirmed; it does not mean every receipt was recorded, accepted, readable, or sufficient for definitive reconciliation. This API does not select history, build adjustment evidence, post consumer statuses, or modify the original reconciliation/checkpoint interfaces. """ if ( isinstance(timeout_seconds, bool) or not isinstance(timeout_seconds, (int, float)) or not isfinite(timeout_seconds) or timeout_seconds <= 0 ): raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN) scope = await _authorize(client, authorizer) proposed = ( prepare_reporting_receipt_submission(scope, receipts) if receipts is not None else None ) state = await _storage(store.reserve(proposed) if proposed is not None else store.get(scope)) if state is None: raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND) _check_state(state, scope) deferred = proposed is not None and proposed.submission_id != state.submission_id if state.pending and state.adcp_version != _VERSION: # An unresolved older reservation retains its exact wire bytes, but # the current SDK must not issue a new request under a retired pin. raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN) while state.pending: await _authorize(client, authorizer, scope) chunk = state.confirmed_chunks request = state.request(chunk) response = None try: response = await asyncio.wait_for( client.sync_reporting_receipts(request), timeout=timeout_seconds ) except Exception: response = None if response is None: return ReportingSubmissionResult( state, deferred, ReportingSubmissionCode.TRANSPORT_UNCERTAIN ) if ( response.success is not True or response.status != "completed" or not isinstance(response.data, SyncReportingReceiptsResponse) ): return ReportingSubmissionResult( state, deferred, ReportingSubmissionCode.RESPONSE_UNCONFIRMED ) invalid = False try: updated = await _storage( store.confirm(scope, state.submission_id, chunk, response.data) ) except ReportingSubmissionError as error: if error.code != ReportingSubmissionCode.INVALID_RESPONSE: raise invalid = True if invalid: return ReportingSubmissionResult( state, deferred, ReportingSubmissionCode.INVALID_RESPONSE ) _check_state(updated, scope, state) if updated.confirmed_chunks <= chunk: raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) state = updated await _authorize(client, authorizer, scope) return ReportingSubmissionResult(state, deferred)Reserve or resume a receipt plan without replacing uncertain work.
Omit
receiptsto resume the latest intent for the freshly authorized scope, including its completed outcomes after a lost return. Otherwise the supplied plan is frozen and atomically reserved before ANY network write. An earlier pending scope intent takes priority, with proposal_deferred=True. The caller must reconsider the deferred plan against fresh seller history.A transport failure, timeout, failed task or malformed response leaves the chunk uncertain and returns pending=True. Retry within the live protocol version uses the exact durable body and idempotency key. A pending intent from a retired version remains readable and byte-identical but is refused before transport; it is never rewritten as a new-version request. Cancellation propagates with the intent still retained. Storage failure raises a closed error; resume the stored scope because the last commit may have succeeded. There is no expiry/abandon/replacement path.
Every item failure is retained alongside successes. Completion only means every submission outcome is confirmed; it does not mean every receipt was recorded, accepted, readable, or sufficient for definitive reconciliation. This API does not select history, build adjustment evidence, post consumer statuses, or modify the original reconciliation/checkpoint interfaces.
Classes
class InMemoryReportingSubmissionIntentStore-
Expand source code
class InMemoryReportingSubmissionIntentStore: """Volatile reference implementation for tests; it cannot survive restart. Production callers must explicitly supply a durable implementation such as PgReportingSubmissionIntentStore. There is no implicit in-memory fallback. """ def __init__(self) -> None: self._lock = asyncio.Lock() self._scopes: dict[str, tuple[bytes, str]] = {} self._submissions: dict[tuple[str, str], ReportingReceiptSubmission] = {} def _current(self, scope: ReportingSubmissionScope) -> ReportingReceiptSubmission | None: entry = self._scopes.get(scope.storage_key) if entry is None: return None if entry[0] != scope.canonical_identity: raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) current = self._submissions.get((scope.storage_key, entry[1])) if current is None or current.scope != scope: raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) validate_submission(current) return current async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission: validate_submission(proposed) if proposed.confirmed_chunks: raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN) async with self._lock: current = self._current(proposed.scope) if current is not None and current.pending: return current key = proposed.scope.storage_key, proposed.submission_id prior = self._submissions.get(key) if prior is not None: validate_submission(prior) if prior.scope != proposed.scope or ( prior._plan != proposed._plan and not is_completed_legacy_replay(prior, proposed) ): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) return prior self._submissions[key] = proposed self._scopes[key[0]] = proposed.scope.canonical_identity, proposed.submission_id return proposed async def get( self, scope: ReportingSubmissionScope, submission_id: str | None = None ) -> ReportingReceiptSubmission | None: async with self._lock: current = self._current(scope) if current is None or submission_id is None: return current result = self._submissions.get((scope.storage_key, submission_id)) if result is not None: validate_submission(result) if result.scope != scope: raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) return result async def confirm( self, scope: ReportingSubmissionScope, submission_id: str, chunk: int, response: SyncReportingReceiptsResponse, ) -> ReportingReceiptSubmission: frozen = freeze_response(response) async with self._lock: current = self._current(scope) state = self._submissions.get((scope.storage_key, submission_id)) if current is None or state is None: raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND) validate_submission(state) if state.scope != scope or ( state.pending and state.submission_id != current.submission_id ): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) updated = confirm_submission(state, chunk, frozen) self._submissions[(scope.storage_key, submission_id)] = updated return updatedVolatile reference implementation for tests; it cannot survive restart.
Production callers must explicitly supply a durable implementation such as PgReportingSubmissionIntentStore. There is no implicit in-memory fallback.
Methods
async def confirm(self,
scope: ReportingSubmissionScope,
submission_id: str,
chunk: int,
response: SyncReportingReceiptsResponse) ‑> ReportingReceiptSubmission-
Expand source code
async def confirm( self, scope: ReportingSubmissionScope, submission_id: str, chunk: int, response: SyncReportingReceiptsResponse, ) -> ReportingReceiptSubmission: frozen = freeze_response(response) async with self._lock: current = self._current(scope) state = self._submissions.get((scope.storage_key, submission_id)) if current is None or state is None: raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND) validate_submission(state) if state.scope != scope or ( state.pending and state.submission_id != current.submission_id ): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) updated = confirm_submission(state, chunk, frozen) self._submissions[(scope.storage_key, submission_id)] = updated return updated async def get(self,
scope: ReportingSubmissionScope,
submission_id: str | None = None) ‑> ReportingReceiptSubmission | None-
Expand source code
async def get( self, scope: ReportingSubmissionScope, submission_id: str | None = None ) -> ReportingReceiptSubmission | None: async with self._lock: current = self._current(scope) if current is None or submission_id is None: return current result = self._submissions.get((scope.storage_key, submission_id)) if result is not None: validate_submission(result) if result.scope != scope: raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) return result async def reserve(self,
proposed: ReportingReceiptSubmission) ‑> ReportingReceiptSubmission-
Expand source code
async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission: validate_submission(proposed) if proposed.confirmed_chunks: raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN) async with self._lock: current = self._current(proposed.scope) if current is not None and current.pending: return current key = proposed.scope.storage_key, proposed.submission_id prior = self._submissions.get(key) if prior is not None: validate_submission(prior) if prior.scope != proposed.scope or ( prior._plan != proposed._plan and not is_completed_legacy_replay(prior, proposed) ): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) return prior self._submissions[key] = proposed self._scopes[key[0]] = proposed.scope.canonical_identity, proposed.submission_id return proposed
class PgReportingSubmissionIntentStore (*, pool: AsyncConnectionPool)-
Expand source code
class PgReportingSubmissionIntentStore: """Production persistence for optional buyer receipt submission intents. Pass a trusted application-owned psycopg AsyncConnectionPool, preferably with a dedicated schema/search_path and bounded connection/statement waits. There is no credential or connection string in this object's representation. Run create_schema explicitly during rollout, then validate custom backends with the same memory/PostgreSQL state-machine vectors. Readiness requires all shipped, validated CHECK definitions as well as the keys and reservation index. Snapshot validation happens outside transactions; a mutation compares the complete locked state, retrying at most 16 times. Contention exhaustion raises STORAGE_UNAVAILABLE and leaves the durable intent available to resume. """ def __init__(self, *, pool: AsyncConnectionPool) -> None: self._pool = pool @_closed_storage_errors async def create_schema(self) -> None: available = False try: import psycopg # noqa: F401 available = True except ImportError: # The optional driver is absent; report PG_REQUIRED outside the handler. pass if not available: raise ReportingSubmissionError(ReportingSubmissionCode.PG_REQUIRED) async with self._pool.connection() as connection, connection.transaction(): await connection.execute("SELECT pg_advisory_xact_lock(%s)", (712071172,)) sql = ( files("adcp.reporting.ledger") .joinpath("reporting_buyer_submissions.sql") .read_text() ) await connection.execute(sql) await self._check_schema_on(connection) async def _check_schema_on(self, connection: Any) -> None: # Verify storage bounds as well as concurrency invariants. IF NOT EXISTS # must not silently bless a pre-created table with weaker constraints. rows = await ( await connection.execute( "SELECT c.conrelid='reporting_buyer_submission_scopes'::regclass," " c.conname,c.contype,c.convalidated,c.condeferrable,c.connoinherit," " pg_get_constraintdef(c.oid) FROM pg_constraint c" " WHERE c.conrelid IN ('reporting_buyer_submission_scopes'::regclass," " 'reporting_buyer_submission_intents'::regclass)" ) ).fetchall() constraints = {(row[0], row[1]): tuple(row[2:]) for row in rows} required = { (True, "reporting_buyer_submission_scopes_pkey"): ( "p", True, False, True, "PRIMARY KEY (scope_sha256)", ), (False, "reporting_buyer_submission_intents_pkey"): ( "p", True, False, True, "PRIMARY KEY (scope_sha256, submission_id)", ), (False, "reporting_buyer_submission_scope_fk"): ( "f", True, False, True, "FOREIGN KEY (scope_sha256) REFERENCES " "reporting_buyer_submission_scopes(scope_sha256)", ), } checks = ( (True, "scope_digest", "scope_sha256 ~ '^[a-f0-9]{64}$'::text"), (True, "scope_bound", "octet_length(canonical_identity) <= 32768"), (False, "submission_id", "submission_id ~ '^reporting-submission:[a-f0-9]{64}$'::text"), (False, "plan_digest", "plan_sha256 ~ '^[a-f0-9]{64}$'::text"), (False, "confirmed_digest", "confirmed_sha256 ~ '^[a-f0-9]{64}$'::text"), (False, "plan_bound", "octet_length(canonical_plan) <= 16777216"), (False, "confirmed_bound", "octet_length(confirmed_results) <= 16777216"), ) for scopes_table, suffix, expression in checks: required[(scopes_table, f"reporting_buyer_{suffix}")] = ( "c", True, False, False, f"CHECK (({expression}))", ) if any(constraints.get(name) != expected for name, expected in required.items()): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) index = await ( await connection.execute( "SELECT i.indisunique,i.indisvalid,i.indisready,i.indnkeyatts," " pg_get_indexdef(i.indexrelid,1,true),pg_get_expr(i.indpred,i.indrelid)," " i.indrelid='reporting_buyer_submission_intents'::regclass" " FROM pg_index i WHERE i.indexrelid='reporting_buyer_one_pending_scope'::regclass" ) ).fetchone() if index is None or tuple(index) != (True, True, True, 1, "scope_sha256", "pending", True): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) async def _snapshot_on( self, connection: Any, key: str, selected_id: str | None, *, lock: bool = False ) -> _Snapshot: scope_query = ( "SELECT canonical_identity,current_submission_id" " FROM reporting_buyer_submission_scopes WHERE scope_sha256=%s FOR UPDATE" if lock else "SELECT canonical_identity,current_submission_id" " FROM reporting_buyer_submission_scopes WHERE scope_sha256=%s" ) row = await (await connection.execute(scope_query, (key,))).fetchone() if row is None: return _Snapshot(None, None, None, False) intent_query = ( "SELECT submission_id,canonical_plan,plan_sha256,confirmed_results," " confirmed_sha256,pending FROM reporting_buyer_submission_intents" " WHERE scope_sha256=%s AND (submission_id=%s OR submission_id=%s) FOR UPDATE" if lock else "SELECT submission_id,canonical_plan,plan_sha256,confirmed_results," " confirmed_sha256,pending FROM reporting_buyer_submission_intents" " WHERE scope_sha256=%s AND (submission_id=%s OR submission_id=%s)" ) rows = await (await connection.execute(intent_query, (key, row[1], selected_id))).fetchall() by_id = {value[0]: tuple(value) for value in rows} has_intents = bool(rows) if row[1] is None: has_intents = ( await ( await connection.execute( "SELECT 1 FROM reporting_buyer_submission_intents" " WHERE scope_sha256=%s LIMIT 1", (key,), ) ).fetchone() ) is not None return _Snapshot( tuple(row), by_id.get(row[1]), by_id.get(row[1] if selected_id is None else selected_id), has_intents, ) async def _snapshot(self, key: str, selected_id: str | None) -> _Snapshot: # Release the read-only MVCC snapshot and connection before CPU work. # No scope row lock or caller-owned validation occurs in this phase. async with self._pool.connection() as connection, connection.transaction(): await connection.execute("SET TRANSACTION ISOLATION LEVEL REPEATABLE READ, READ ONLY") return await self._snapshot_on(connection, key, selected_id) def _validate_snapshot( self, scope: ReportingSubmissionScope, snapshot: _Snapshot ) -> tuple[ReportingReceiptSubmission | None, ReportingReceiptSubmission | None]: if snapshot.scope is None: return None, None if snapshot.scope[0] != scope.canonical_identity.decode() or ( snapshot.scope[1] is None and snapshot.has_intents ): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) if snapshot.scope[1] is not None and snapshot.current is None: raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) current = decode_submission(scope, snapshot.current) if snapshot.current else None selected = ( current if snapshot.selected is snapshot.current else decode_submission(scope, snapshot.selected) if snapshot.selected else None ) if selected is not None and selected.pending and selected != current: raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) return current, selected @_closed_storage_errors async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission: validate_submission(proposed) if proposed.confirmed_chunks: raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN) scope = proposed.scope key, identity = scope.storage_key, scope.canonical_identity.decode() plan, plan_digest, confirmed, confirmed_digest = encode_submission(proposed) for _ in range(_MAX_CAS_ATTEMPTS): snapshot = await self._snapshot(key, proposed.submission_id) current, previous = await asyncio.to_thread(self._validate_snapshot, scope, snapshot) if current is not None and current.pending: return current if previous is not None: if previous._plan != proposed._plan and not is_completed_legacy_replay( previous, proposed ): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) return previous async with self._pool.connection() as connection, connection.transaction(): expected = snapshot if snapshot.scope is None: await connection.execute( "INSERT INTO reporting_buyer_submission_scopes" " (scope_sha256,canonical_identity) VALUES (%s,%s)" " ON CONFLICT (scope_sha256) DO NOTHING", (key, identity), ) expected = _Snapshot((identity, None), None, None, False) actual = await self._snapshot_on(connection, key, proposed.submission_id, lock=True) if actual != expected: continue await connection.execute( "INSERT INTO reporting_buyer_submission_intents" " (scope_sha256,submission_id,canonical_plan,plan_sha256," " confirmed_results,confirmed_sha256,pending) VALUES (%s,%s,%s,%s,%s,%s,true)", ( key, proposed.submission_id, plan, plan_digest, confirmed, confirmed_digest, ), ) await connection.execute( "UPDATE reporting_buyer_submission_scopes SET current_submission_id=%s" " WHERE scope_sha256=%s", (proposed.submission_id, key), ) return proposed raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE) @_closed_storage_errors async def get( self, scope: ReportingSubmissionScope, submission_id: str | None = None ) -> ReportingReceiptSubmission | None: snapshot = await self._snapshot(scope.storage_key, submission_id) _, selected = await asyncio.to_thread(self._validate_snapshot, scope, snapshot) return selected @_closed_storage_errors async def confirm( self, scope: ReportingSubmissionScope, submission_id: str, chunk: int, response: SyncReportingReceiptsResponse, ) -> ReportingReceiptSubmission: # Caller-owned Pydantic objects are detached before the first await. # Retries compare the same immutable response, even if its owner edits it. frozen = freeze_response(response) key = scope.storage_key for _ in range(_MAX_CAS_ATTEMPTS): snapshot = await self._snapshot(key, submission_id) _, state = await asyncio.to_thread(self._validate_snapshot, scope, snapshot) if state is None: raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND) updated = await asyncio.to_thread(confirm_submission, state, chunk, frozen) _, _, confirmed, digest = encode_submission(updated) pending = updated.pending async with self._pool.connection() as connection, connection.transaction(): actual = await self._snapshot_on(connection, key, submission_id, lock=True) if actual != snapshot: continue if updated is not state: await connection.execute( "UPDATE reporting_buyer_submission_intents" " SET confirmed_results=%s,confirmed_sha256=%s,pending=%s" " WHERE scope_sha256=%s AND submission_id=%s", (confirmed, digest, pending, key, submission_id), ) return updated # Contention never frees the lane or guesses whether another commit won. raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE)Production persistence for optional buyer receipt submission intents.
Pass a trusted application-owned psycopg AsyncConnectionPool, preferably with a dedicated schema/search_path and bounded connection/statement waits. There is no credential or connection string in this object's representation. Run create_schema explicitly during rollout, then validate custom backends with the same memory/PostgreSQL state-machine vectors. Readiness requires all shipped, validated CHECK definitions as well as the keys and reservation index. Snapshot validation happens outside transactions; a mutation compares the complete locked state, retrying at most 16 times. Contention exhaustion raises STORAGE_UNAVAILABLE and leaves the durable intent available to resume.
Methods
async def confirm(self,
scope: ReportingSubmissionScope,
submission_id: str,
chunk: int,
response: SyncReportingReceiptsResponse) ‑> ReportingReceiptSubmission-
Expand source code
@_closed_storage_errors async def confirm( self, scope: ReportingSubmissionScope, submission_id: str, chunk: int, response: SyncReportingReceiptsResponse, ) -> ReportingReceiptSubmission: # Caller-owned Pydantic objects are detached before the first await. # Retries compare the same immutable response, even if its owner edits it. frozen = freeze_response(response) key = scope.storage_key for _ in range(_MAX_CAS_ATTEMPTS): snapshot = await self._snapshot(key, submission_id) _, state = await asyncio.to_thread(self._validate_snapshot, scope, snapshot) if state is None: raise ReportingSubmissionError(ReportingSubmissionCode.NOT_FOUND) updated = await asyncio.to_thread(confirm_submission, state, chunk, frozen) _, _, confirmed, digest = encode_submission(updated) pending = updated.pending async with self._pool.connection() as connection, connection.transaction(): actual = await self._snapshot_on(connection, key, submission_id, lock=True) if actual != snapshot: continue if updated is not state: await connection.execute( "UPDATE reporting_buyer_submission_intents" " SET confirmed_results=%s,confirmed_sha256=%s,pending=%s" " WHERE scope_sha256=%s AND submission_id=%s", (confirmed, digest, pending, key, submission_id), ) return updated # Contention never frees the lane or guesses whether another commit won. raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE) async def create_schema(self) ‑> None-
Expand source code
@_closed_storage_errors async def create_schema(self) -> None: available = False try: import psycopg # noqa: F401 available = True except ImportError: # The optional driver is absent; report PG_REQUIRED outside the handler. pass if not available: raise ReportingSubmissionError(ReportingSubmissionCode.PG_REQUIRED) async with self._pool.connection() as connection, connection.transaction(): await connection.execute("SELECT pg_advisory_xact_lock(%s)", (712071172,)) sql = ( files("adcp.reporting.ledger") .joinpath("reporting_buyer_submissions.sql") .read_text() ) await connection.execute(sql) await self._check_schema_on(connection) async def get(self,
scope: ReportingSubmissionScope,
submission_id: str | None = None) ‑> ReportingReceiptSubmission | None-
Expand source code
@_closed_storage_errors async def get( self, scope: ReportingSubmissionScope, submission_id: str | None = None ) -> ReportingReceiptSubmission | None: snapshot = await self._snapshot(scope.storage_key, submission_id) _, selected = await asyncio.to_thread(self._validate_snapshot, scope, snapshot) return selected async def reserve(self,
proposed: ReportingReceiptSubmission) ‑> ReportingReceiptSubmission-
Expand source code
@_closed_storage_errors async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission: validate_submission(proposed) if proposed.confirmed_chunks: raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_PLAN) scope = proposed.scope key, identity = scope.storage_key, scope.canonical_identity.decode() plan, plan_digest, confirmed, confirmed_digest = encode_submission(proposed) for _ in range(_MAX_CAS_ATTEMPTS): snapshot = await self._snapshot(key, proposed.submission_id) current, previous = await asyncio.to_thread(self._validate_snapshot, scope, snapshot) if current is not None and current.pending: return current if previous is not None: if previous._plan != proposed._plan and not is_completed_legacy_replay( previous, proposed ): raise ReportingSubmissionError(ReportingSubmissionCode.HISTORY_CORRUPT) return previous async with self._pool.connection() as connection, connection.transaction(): expected = snapshot if snapshot.scope is None: await connection.execute( "INSERT INTO reporting_buyer_submission_scopes" " (scope_sha256,canonical_identity) VALUES (%s,%s)" " ON CONFLICT (scope_sha256) DO NOTHING", (key, identity), ) expected = _Snapshot((identity, None), None, None, False) actual = await self._snapshot_on(connection, key, proposed.submission_id, lock=True) if actual != expected: continue await connection.execute( "INSERT INTO reporting_buyer_submission_intents" " (scope_sha256,submission_id,canonical_plan,plan_sha256," " confirmed_results,confirmed_sha256,pending) VALUES (%s,%s,%s,%s,%s,%s,true)", ( key, proposed.submission_id, plan, plan_digest, confirmed, confirmed_digest, ), ) await connection.execute( "UPDATE reporting_buyer_submission_scopes SET current_submission_id=%s" " WHERE scope_sha256=%s", (proposed.submission_id, key), ) return proposed raise ReportingSubmissionError(ReportingSubmissionCode.STORAGE_UNAVAILABLE)
class ReportingReceiptFailureCode (*args, **kwds)-
Expand source code
class ReportingReceiptFailureCode(str, Enum): """Known item failures; unrecognized seller codes map to UNKNOWN. Seller messages/details are deliberately not persisted or exposed. Each failed item and each error position is retained, including unknown codes. """ UNKNOWN = "UNKNOWN" INVALID_REQUEST = "INVALID_REQUEST" UNAUTHORIZED = "UNAUTHORIZED" RATE_LIMITED = "RATE_LIMITED" NOT_SUPPORTED = "NOT_SUPPORTED" IDEMPOTENCY_CONFLICT = "IDEMPOTENCY_CONFLICT" INVALID_REPORTING_RECORD = "INVALID_REPORTING_RECORD" REPORTING_RECORD_UNAVAILABLE = "REPORTING_RECORD_UNAVAILABLE" REPORTING_HISTORY_CORRUPT = "REPORTING_HISTORY_CORRUPT" REPORTING_IDENTITY_CONFLICT = "REPORTING_IDENTITY_CONFLICT" REPORTING_TIME_INVALID = "REPORTING_TIME_INVALID" RECEIPTS_NOT_ENABLED = "RECEIPTS_NOT_ENABLED" RECEIVED_AT_READ_ONLY = "RECEIVED_AT_READ_ONLY" RECEIPT_PROFILE_MISMATCH = "RECEIPT_PROFILE_MISMATCH" RECEIPT_TOTALS_MISMATCH = "RECEIPT_TOTALS_MISMATCH" RECEIPT_EVIDENCE_MISMATCH = "RECEIPT_EVIDENCE_MISMATCH" MATERIALIZATION_UNREADABLE = "MATERIALIZATION_UNREADABLE" ADJUSTMENT_REQUIRES_OFFICIAL = "ADJUSTMENT_REQUIRES_OFFICIAL" ADJUSTMENT_ORDER_INVALID = "ADJUSTMENT_ORDER_INVALID" ADJUSTMENT_DIGEST_MISMATCH = "ADJUSTMENT_DIGEST_MISMATCH" ACCEPTED_RECEIPT_TERMINAL = "ACCEPTED_RECEIPT_TERMINAL"Known item failures; unrecognized seller codes map to UNKNOWN.
Seller messages/details are deliberately not persisted or exposed. Each failed item and each error position is retained, including unknown codes.
Ancestors
- builtins.str
- enum.Enum
Class variables
var ACCEPTED_RECEIPT_TERMINALvar ADJUSTMENT_DIGEST_MISMATCHvar ADJUSTMENT_ORDER_INVALIDvar ADJUSTMENT_REQUIRES_OFFICIALvar IDEMPOTENCY_CONFLICTvar INVALID_REPORTING_RECORDvar INVALID_REQUESTvar MATERIALIZATION_UNREADABLEvar NOT_SUPPORTEDvar RATE_LIMITEDvar RECEIPTS_NOT_ENABLEDvar RECEIPT_EVIDENCE_MISMATCHvar RECEIPT_PROFILE_MISMATCHvar RECEIPT_TOTALS_MISMATCHvar RECEIVED_AT_READ_ONLYvar REPORTING_HISTORY_CORRUPTvar REPORTING_IDENTITY_CONFLICTvar REPORTING_RECORD_UNAVAILABLEvar REPORTING_TIME_INVALIDvar UNAUTHORIZEDvar UNKNOWN
class ReportingReceiptOutcome (ordinal: int,
kind: ReceiptKind,
result: ReceiptResult,
error_codes: tuple[ReportingReceiptFailureCode, ...] = ())-
Expand source code
@dataclass(frozen=True) class ReportingReceiptOutcome: """One confirmed item, in original input order, with fresh model views.""" ordinal: int kind: ReceiptKind result: ReceiptResult error_codes: tuple[ReportingReceiptFailureCode, ...] = () _submitted: bytes = field(default=b"", repr=False) _stored: bytes | None = field(default=None, repr=False) @property def submitted_receipt(self) -> ReportingSubmissionReceipt: return _receipt(self.kind, json.loads(self._submitted)) @property def receipt(self) -> ReportingSubmissionReceipt | None: """The seller's actual stored receipt on success, including received_at.""" return _receipt(self.kind, json.loads(self._stored)) if self._stored is not None else NoneOne confirmed item, in original input order, with fresh model views.
Instance variables
var error_codes : tuple[ReportingReceiptFailureCode, ...]var kind : Literal['receipt', 'adjustment_receipt']var ordinal : intprop receipt : ReportingSubmissionReceipt | None-
Expand source code
@property def receipt(self) -> ReportingSubmissionReceipt | None: """The seller's actual stored receipt on success, including received_at.""" return _receipt(self.kind, json.loads(self._stored)) if self._stored is not None else NoneThe seller's actual stored receipt on success, including received_at.
var result : Literal['recorded', 'unchanged', 'failed']prop submitted_receipt : ReportingSubmissionReceipt-
Expand source code
@property def submitted_receipt(self) -> ReportingSubmissionReceipt: return _receipt(self.kind, json.loads(self._submitted))
class ReportingReceiptSubmission (scope: ReportingSubmissionScope,
submission_id: str,
_plan: bytes)-
Expand source code
@dataclass(frozen=True) class ReportingReceiptSubmission: """Exact reserved plan and its immutable, fully confirmed chunk prefix. Construct with :func:`prepare_reporting_receipt_submission`. Stores validate persisted records before returning them. Private bytes keep mutable caller models, driver results, and representations out of the retained identity. """ scope: ReportingSubmissionScope = field(repr=False) submission_id: str = field(repr=False) _plan: bytes = field(repr=False) _confirmed: tuple[bytes, ...] = field(default=(), repr=False) @property def chunk_count(self) -> int: return len(_validated_requests(self)) @property def confirmed_chunks(self) -> int: return len(self._confirmed) @property def pending(self) -> bool: return self.confirmed_chunks < self.chunk_count @property def adcp_version(self) -> str: """The immutable request version, including for completed older plans.""" return str(json.loads(_validated_requests(self)[0])["adcp_version"]) def request(self, ordinal: int) -> SyncReportingReceiptsRequest: """Return a fresh typed view of one exact persisted request body.""" body = json.loads(_validated_requests(self)[ordinal]) return SyncReportingReceiptsRequest.model_validate(body) @property def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]: plan = json.loads(self._plan) confirmed = { item["reporting_receipt_id"]: item for chunk in self._confirmed for item in json.loads(chunk) } outcomes = [] for ordinal, item in enumerate(plan["items"]): result = confirmed.get(item["body"]["reporting_receipt_id"]) if result is not None: stored = result.get("receipt") outcomes.append( ReportingReceiptOutcome( ordinal, item["kind"], result["result"], tuple(ReportingReceiptFailureCode(code) for code in result["errors"]), canonical_json_utf8_v1(item["body"]), canonical_json_utf8_v1(stored) if stored is not None else None, ) ) return tuple(outcomes)Exact reserved plan and its immutable, fully confirmed chunk prefix.
Construct with :func:
prepare_reporting_receipt_submission(). Stores validate persisted records before returning them. Private bytes keep mutable caller models, driver results, and representations out of the retained identity.Instance variables
prop adcp_version : str-
Expand source code
@property def adcp_version(self) -> str: """The immutable request version, including for completed older plans.""" return str(json.loads(_validated_requests(self)[0])["adcp_version"])The immutable request version, including for completed older plans.
prop chunk_count : int-
Expand source code
@property def chunk_count(self) -> int: return len(_validated_requests(self)) prop confirmed_chunks : int-
Expand source code
@property def confirmed_chunks(self) -> int: return len(self._confirmed) prop outcomes : tuple[ReportingReceiptOutcome, ...]-
Expand source code
@property def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]: plan = json.loads(self._plan) confirmed = { item["reporting_receipt_id"]: item for chunk in self._confirmed for item in json.loads(chunk) } outcomes = [] for ordinal, item in enumerate(plan["items"]): result = confirmed.get(item["body"]["reporting_receipt_id"]) if result is not None: stored = result.get("receipt") outcomes.append( ReportingReceiptOutcome( ordinal, item["kind"], result["result"], tuple(ReportingReceiptFailureCode(code) for code in result["errors"]), canonical_json_utf8_v1(item["body"]), canonical_json_utf8_v1(stored) if stored is not None else None, ) ) return tuple(outcomes) prop pending : bool-
Expand source code
@property def pending(self) -> bool: return self.confirmed_chunks < self.chunk_count var scope : ReportingSubmissionScopevar submission_id : str
Methods
def request(self, ordinal: int) ‑> SyncReportingReceiptsRequest-
Expand source code
def request(self, ordinal: int) -> SyncReportingReceiptsRequest: """Return a fresh typed view of one exact persisted request body.""" body = json.loads(_validated_requests(self)[ordinal]) return SyncReportingReceiptsRequest.model_validate(body)Return a fresh typed view of one exact persisted request body.
class ReportingReceiptSubmissionClient (*args, **kwargs)-
Expand source code
class ReportingReceiptSubmissionClient(Protocol): async def sync_reporting_receipts( self, request: SyncReportingReceiptsRequest ) -> TaskResult[SyncReportingReceiptsResponse]: ...Base class for protocol classes.
Protocol classes are defined as::
class Proto(Protocol): def meth(self) -> int: ...Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).
For example::
class C: def meth(self) -> int: return 0 def func(x: Proto) -> int: return x.meth() func(C()) # Passes static type 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 sync_reporting_receipts(self, request: SyncReportingReceiptsRequest) ‑> TaskResult[SyncReportingReceiptsResponse]-
Expand source code
async def sync_reporting_receipts( self, request: SyncReportingReceiptsRequest ) -> TaskResult[SyncReportingReceiptsResponse]: ...
class ReportingSubmissionAuthorizer (*args, **kwargs)-
Expand source code
class ReportingSubmissionAuthorizer(Protocol): """Trusted application adapter resolving the exact client's current access. Resolve seller identity from trusted client/registry configuration, account from an authorized seller account lookup, and consumer from authenticated credentials/registry. Never source these identities from receipt payloads, an asserted account reference, or a debug/raw response. Aliases must already be resolved. Raise on revoked access or an identity disagreement. Called before any store read/reservation, before each send, and before returning completed cached outcomes. It must authorize replay afresh too. """ async def __call__( self, client: ReportingReceiptSubmissionClient ) -> ReportingSubmissionScope: ...Trusted application adapter resolving the exact client's current access.
Resolve seller identity from trusted client/registry configuration, account from an authorized seller account lookup, and consumer from authenticated credentials/registry. Never source these identities from receipt payloads, an asserted account reference, or a debug/raw response. Aliases must already be resolved. Raise on revoked access or an identity disagreement.
Called before any store read/reservation, before each send, and before returning completed cached outcomes. It must authorize replay afresh too.
Ancestors
- typing.Protocol
- typing.Generic
class ReportingSubmissionCode (*args, **kwds)-
Expand source code
class ReportingSubmissionCode(str, Enum): INVALID_SCOPE = "INVALID_SUBMISSION_SCOPE" INVALID_PLAN = "INVALID_SUBMISSION_PLAN" UNAUTHORIZED = "SUBMISSION_NOT_AUTHORIZED" NOT_FOUND = "SUBMISSION_NOT_FOUND" HISTORY_CORRUPT = "SUBMISSION_HISTORY_CORRUPT" STORAGE_UNAVAILABLE = "SUBMISSION_STORAGE_UNAVAILABLE" PG_REQUIRED = "SUBMISSION_PG_REQUIRED" INVALID_RESPONSE = "SUBMISSION_RESPONSE_INVALID" TRANSPORT_UNCERTAIN = "SUBMISSION_TRANSPORT_UNCERTAIN" RESPONSE_UNCONFIRMED = "SUBMISSION_RESPONSE_UNCONFIRMED"str(object='') -> str str(bytes_or_buffer[, encoding[, errors]]) -> str
Create a new string object from the given object. If encoding or errors is specified, then the object must expose a data buffer that will be decoded using the given encoding and error handler. Otherwise, returns the result of object.str() (if defined) or repr(object). encoding defaults to sys.getdefaultencoding(). errors defaults to 'strict'.
Ancestors
- builtins.str
- enum.Enum
Class variables
var HISTORY_CORRUPTvar INVALID_PLANvar INVALID_RESPONSEvar INVALID_SCOPEvar NOT_FOUNDvar PG_REQUIREDvar RESPONSE_UNCONFIRMEDvar STORAGE_UNAVAILABLEvar TRANSPORT_UNCERTAINvar UNAUTHORIZED
class ReportingSubmissionError (code: ReportingSubmissionCode)-
Expand source code
class ReportingSubmissionError(RuntimeError): """Closed code only: no provider, authentication, SQL or wire body details.""" def __init__(self, code: ReportingSubmissionCode) -> None: self.code = ( code if isinstance(code, ReportingSubmissionCode) else ReportingSubmissionCode.INVALID_PLAN ) super().__init__(self.code.value)Closed code only: no provider, authentication, SQL or wire body details.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class ReportingSubmissionIntentStore (*args, **kwargs)-
Expand source code
class ReportingSubmissionIntentStore(Protocol): """Atomic exact-request reservation and monotonic confirmed chunk storage. Adopters implementing this protocol must run the shared store vectors against their durable backend. A scope has one unresolved reservation, never a TTL, lease expiry, or automatic abandonment. Each mutation commits before return. No transport call may run while holding a storage transaction or lock. The supplied scope is already resolved/authorized by a trusted adapter. The store is an internal persistence boundary, not a request authentication API. Old ReportingCheckpointStore implementations need no new methods. """ async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission: """Atomically reserve or return the scope's earlier unresolved intent. A different proposal must not replace an uncertain request. An already completed proposal for the same frozen receipt items returns its retained outcomes across the rc.6/rc.7 to 3.2 upgrade. Preserve all old request bytes, chunk keys/order and confirmed outcomes across restart. """ ... async def get( self, scope: ReportingSubmissionScope, submission_id: str | None = None ) -> ReportingReceiptSubmission | None: """Read a validated reservation, including completed work after a lost return. With no ID return the latest reserved intent; do not erase it when it completes. Never resolve an ID outside the exact trusted scope. """ ... async def confirm( self, scope: ReportingSubmissionScope, submission_id: str, chunk: int, response: SyncReportingReceiptsResponse, ) -> ReportingReceiptSubmission: """Atomically acknowledge full unique ID/body coverage for the next chunk. Retain item successes AND failures in the original plan order. Confirmed chunks never change. A concurrent duplicate acknowledgement returns the retained state; it cannot replace an earlier confirmed outcome. """ ...Atomic exact-request reservation and monotonic confirmed chunk storage.
Adopters implementing this protocol must run the shared store vectors against their durable backend. A scope has one unresolved reservation, never a TTL, lease expiry, or automatic abandonment. Each mutation commits before return. No transport call may run while holding a storage transaction or lock.
The supplied scope is already resolved/authorized by a trusted adapter. The store is an internal persistence boundary, not a request authentication API. Old ReportingCheckpointStore implementations need no new methods.
Ancestors
- typing.Protocol
- typing.Generic
Methods
async def confirm(self,
scope: ReportingSubmissionScope,
submission_id: str,
chunk: int,
response: SyncReportingReceiptsResponse) ‑> ReportingReceiptSubmission-
Expand source code
async def confirm( self, scope: ReportingSubmissionScope, submission_id: str, chunk: int, response: SyncReportingReceiptsResponse, ) -> ReportingReceiptSubmission: """Atomically acknowledge full unique ID/body coverage for the next chunk. Retain item successes AND failures in the original plan order. Confirmed chunks never change. A concurrent duplicate acknowledgement returns the retained state; it cannot replace an earlier confirmed outcome. """ ...Atomically acknowledge full unique ID/body coverage for the next chunk.
Retain item successes AND failures in the original plan order. Confirmed chunks never change. A concurrent duplicate acknowledgement returns the retained state; it cannot replace an earlier confirmed outcome.
async def get(self,
scope: ReportingSubmissionScope,
submission_id: str | None = None) ‑> ReportingReceiptSubmission | None-
Expand source code
async def get( self, scope: ReportingSubmissionScope, submission_id: str | None = None ) -> ReportingReceiptSubmission | None: """Read a validated reservation, including completed work after a lost return. With no ID return the latest reserved intent; do not erase it when it completes. Never resolve an ID outside the exact trusted scope. """ ...Read a validated reservation, including completed work after a lost return.
With no ID return the latest reserved intent; do not erase it when it completes. Never resolve an ID outside the exact trusted scope.
async def reserve(self,
proposed: ReportingReceiptSubmission) ‑> ReportingReceiptSubmission-
Expand source code
async def reserve(self, proposed: ReportingReceiptSubmission) -> ReportingReceiptSubmission: """Atomically reserve or return the scope's earlier unresolved intent. A different proposal must not replace an uncertain request. An already completed proposal for the same frozen receipt items returns its retained outcomes across the rc.6/rc.7 to 3.2 upgrade. Preserve all old request bytes, chunk keys/order and confirmed outcomes across restart. """ ...Atomically reserve or return the scope's earlier unresolved intent.
A different proposal must not replace an uncertain request. An already completed proposal for the same frozen receipt items returns its retained outcomes across the rc.6/rc.7 to 3.2 upgrade. Preserve all old request bytes, chunk keys/order and confirmed outcomes across restart.
class ReportingSubmissionResult (submission: ReportingReceiptSubmission,
proposal_deferred: bool = False,
diagnostic: ReportingSubmissionCode | None = None)-
Expand source code
@dataclass(frozen=True) class ReportingSubmissionResult: """Submission outcomes only; completion does not establish reconciliation. ``proposal_deferred`` means an earlier pending scope reservation was resumed instead. Inspect its outcomes before planning any subsequent submission. """ submission: ReportingReceiptSubmission proposal_deferred: bool = False diagnostic: ReportingSubmissionCode | None = None @property def pending(self) -> bool: return self.submission.pending @property def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]: return self.submission.outcomes @property def submitted_receipts(self) -> tuple[ReportingSubmissionReceipt, ...]: return tuple( receipt for outcome in self.outcomes if (receipt := outcome.receipt) is not None )Submission outcomes only; completion does not establish reconciliation.
proposal_deferredmeans an earlier pending scope reservation was resumed instead. Inspect its outcomes before planning any subsequent submission.Instance variables
var diagnostic : ReportingSubmissionCode | Noneprop outcomes : tuple[ReportingReceiptOutcome, ...]-
Expand source code
@property def outcomes(self) -> tuple[ReportingReceiptOutcome, ...]: return self.submission.outcomes prop pending : bool-
Expand source code
@property def pending(self) -> bool: return self.submission.pending var proposal_deferred : boolvar submission : ReportingReceiptSubmissionprop submitted_receipts : tuple[ReportingSubmissionReceipt, ...]-
Expand source code
@property def submitted_receipts(self) -> tuple[ReportingSubmissionReceipt, ...]: return tuple( receipt for outcome in self.outcomes if (receipt := outcome.receipt) is not None )
class ReportingSubmissionScope (seller_id: str, account_id: str, consumer_id: str)-
Expand source code
@dataclass(frozen=True) class ReportingSubmissionScope: """Identity resolved by a trusted authorization adapter, never from a request. ``seller_id`` identifies the configured seller, ``account_id`` is the seller-resolved account and ``consumer_id`` is the canonical authenticated principal. Syntax checks cannot establish provenance: the required authorizer must bind all three to the exact client and its current access. No credential or natural-key account assertion belongs in these fields. """ seller_id: str = field(repr=False) account_id: str = field(repr=False) consumer_id: str = field(repr=False) def __post_init__(self) -> None: valid = False try: consumer_reference(self.seller_id) principal_reference(self.account_id) canonical_consumer(self.consumer_id) valid = all(value == value.strip() for value in (self.seller_id, self.account_id)) except (ValueError, TypeError, RuntimeError): # Map provider validation details to the closed INVALID_SCOPE error below. pass if not valid: raise ReportingSubmissionError(ReportingSubmissionCode.INVALID_SCOPE) @property def canonical_identity(self) -> bytes: return canonical_json_utf8_v1([self.seller_id, self.account_id, self.consumer_id]) @property def storage_key(self) -> str: # Index a fixed-size digest, not a possibly 2048-character principal. # Stores also compare the full identity to fail closed on collisions. return hashlib.sha256(self.canonical_identity).hexdigest()Identity resolved by a trusted authorization adapter, never from a request.
seller_ididentifies the configured seller,account_idis the seller-resolved account andconsumer_idis the canonical authenticated principal. Syntax checks cannot establish provenance: the required authorizer must bind all three to the exact client and its current access. No credential or natural-key account assertion belongs in these fields.Instance variables
var account_id : strprop canonical_identity : bytes-
Expand source code
@property def canonical_identity(self) -> bytes: return canonical_json_utf8_v1([self.seller_id, self.account_id, self.consumer_id]) var consumer_id : strvar seller_id : strprop storage_key : str-
Expand source code
@property def storage_key(self) -> str: # Index a fixed-size digest, not a possibly 2048-character principal. # Stores also compare the full identity to fail closed on collisions. return hashlib.sha256(self.canonical_identity).hexdigest()