Module adcp.reporting.materializer
Destination contracts, verification and optional durable materialization.
Use the immutable registry to prepare all frozen source rows, then invoke write and verify explicitly with separate authorization sessions. The reference writer is exclusively for tests/development. The durable service owns fencing, retry allocation, target reselection and atomic finish. Production tier and notification activation additionally require the complete seller projection.
Sub-modules
adcp.reporting.materializer.capture-
Immutable finish boundaries and a real, isolated notification enqueue path …
adcp.reporting.materializer.contracts-
B1 destination I/O contracts. No scheduling, persistence or capability claims …
adcp.reporting.materializer.memory-
Deterministic in-memory model of the durable materializer state machine …
adcp.reporting.materializer.pg-
PostgreSQL account-first reservation, fenced resume and verified atomic finish.
adcp.reporting.materializer.publication-
SDK-owned canonical evidence before an immutable source revision is committed.
adcp.reporting.materializer.reference-
Deterministic memory destination for tests/development. NEVER production support.
adcp.reporting.materializer.schema-
Materializer readiness is independent of every prior mandatory manifest.
adcp.reporting.materializer.service-
Service-owned heartbeat around destination I/O outside database transactions.
adcp.reporting.materializer.verification-
SDK-owned verification of frozen source content and actual destination reads …
adcp.reporting.materializer.work-
Optional durable work contracts, independent of the foundation store protocols.
Functions
def parse_reporting_json(payload: bytes,
limits: ReportingVerificationLimits = ReportingVerificationLimits(max_depth=32, max_value_bytes=1048576, max_total_bytes=67108864, max_items=1000000, max_rows=100000, max_pages=10000, max_objects=10000, max_chunks=1000000)) ‑> bool | int | str | list[bool | int | str | list[JsonValue] | dict[str, bool | int | str | list[JsonValue] | dict[str, JsonValue] | None] | None] | dict[str, bool | int | str | list[JsonValue] | dict[str, JsonValue] | None] | None-
Expand source code
def parse_reporting_json( payload: bytes, limits: ReportingVerificationLimits = DEFAULT_LIMITS ) -> JsonValue: """Reject ambiguous JSON before it can be normalized into a Python mapping.""" if type(payload) is not bytes: raise failure("SOURCE_INVALID") if len(payload) > limits.max_value_bytes: raise failure("LIMIT_EXCEEDED") # Bound parser recursion before json.loads allocates a nested tree. depth, quoted, escaped = 0, False, False for byte in payload: if quoted: if escaped: escaped = False elif byte == 92: escaped = True elif byte == 34: quoted = False elif byte == 34: quoted = True elif byte in (91, 123): depth += 1 if depth > limits.max_depth: raise failure("LIMIT_EXCEEDED") elif byte in (93, 125): depth -= 1 def pairs(values: list[tuple[str, JsonValue]]) -> dict[str, JsonValue]: result: dict[str, JsonValue] = {} for key, value in values: if key in result: raise ValueError result[key] = value return result def reject(value: str) -> JsonValue: raise ValueError invalid = False try: value = json.loads( payload.decode("utf-8"), object_pairs_hook=pairs, parse_constant=reject, parse_float=reject, ) except (ValueError, UnicodeError, RecursionError): invalid = True if invalid: raise failure("SOURCE_INVALID") strict_reporting_json(value, limits) return cast(JsonValue, value)Reject ambiguous JSON before it can be normalized into a Python mapping.
def reference_digest(verifier: ReportingRevisionVerifier,
rows: list[object]) ‑> ReportingCanonicalDigest-
Expand source code
def reference_digest( verifier: ReportingRevisionVerifier, rows: list[object] ) -> ReportingCanonicalDigest: encoded, _ = verifier.canonicalize(rows) contract = verifier.key.canonicalization return ReportingCanonicalDigest( hashlib.sha256(b"[" + b",".join(encoded) + b"]").hexdigest(), contract.canonicalization_id, contract.canonicalization_uri, contract.canonicalization_sha256, ) def reference_verifier(capability: ReportingWriterCapability | None = None) ‑> ReportingRevisionVerifier-
Expand source code
def reference_verifier( capability: ReportingWriterCapability | None = None, ) -> ReportingRevisionVerifier: """One installed example definition/canonicalization; no network resolution.""" root = files("adcp.reporting.materializer").joinpath("assets") definition = root.joinpath("reference-definition.json").read_bytes() schema = root.joinpath("reference-row-schema.json").read_bytes() contract = root.joinpath("reference-canonicalization.json").read_bytes() capability = capability or ReportingWriterCapability( "file_transfer", "reference-memory", "jsonl", "canonical_digest", "producer", "immutable_location", "sha256", "conditional_create", ) return ReportingRevisionVerifier( ReportingVerificationKey( "reference-report-v1", "paid_media_delivery", ReportingDefinitionBinding( "https://contracts.example.test/reference-definition.json", hashlib.sha256(definition).hexdigest(), "1.0.0", "https://contracts.example.test/reference-row-schema.json", hashlib.sha256(schema).hexdigest(), monetary_metric_units=(("spend", "USD"),), monetary_control_total_units=(("spend", "USD"),), ), ReportingCanonicalization( "reference-jcs-rows-v1", "https://contracts.example.test/reference-canonicalization.json", hashlib.sha256(contract).hexdigest(), ), capability, ), definition, schema, contract, )One installed example definition/canonicalization; no network resolution.
def strict_reporting_json(value: object,
limits: ReportingVerificationLimits = ReportingVerificationLimits(max_depth=32, max_value_bytes=1048576, max_total_bytes=67108864, max_items=1000000, max_rows=100000, max_pages=10000, max_objects=10000, max_chunks=1000000)) ‑> bytes-
Expand source code
def strict_reporting_json( value: object, limits: ReportingVerificationLimits = DEFAULT_LIMITS ) -> bytes: """Encode exact JSON types, without coercion, subclass hooks or Unicode normalization. Integers must be JS-safe; exact decimals are strings. Floats, tuples, subclasses and lone surrogates are deliberately outside this SDK profile. """ stack = [(value, 0)] count, encoded_bytes = 0, 0 while stack: item, depth = stack.pop() count += 1 if depth > limits.max_depth or count > limits.max_items: raise failure("LIMIT_EXCEEDED") kind = type(item) if kind is str: text = cast(str, item) if len(text) > limits.max_value_bytes: raise failure("LIMIT_EXCEEDED") encoded_bytes += 2 # Opening/closing quotes, including escaped UTF-8. for char in text: code = ord(char) if 0xD800 <= code <= 0xDFFF: raise failure("SOURCE_INVALID") if char in '\\"\b\t\n\f\r': encoded_bytes += 2 elif code < 0x20: encoded_bytes += 6 else: encoded_bytes += ( 1 if code < 0x80 else 2 if code < 0x800 else 3 if code < 0x10000 else 4 ) if encoded_bytes > limits.max_value_bytes: raise failure("LIMIT_EXCEEDED") elif kind is int: if abs(cast(int, item)) > MAX_SAFE_INTEGER: raise failure("SOURCE_INVALID") encoded_bytes += len(str(item)) elif item is None or kind is bool: encoded_bytes += 5 if item is False else 4 elif kind is list: values = cast(list[object], item) if len(values) + count + len(stack) > limits.max_items: raise failure("LIMIT_EXCEEDED") encoded_bytes += 2 + max(0, len(values) - 1) stack.extend((v, depth + 1) for v in values) elif kind is dict: mapping = cast(dict[object, object], item) if len(mapping) * 2 + count + len(stack) > limits.max_items: raise failure("LIMIT_EXCEEDED") encoded_bytes += 2 + max(0, 2 * len(mapping) - 1) for key, child in mapping.items(): if type(key) is not str: raise failure("SOURCE_INVALID") stack.extend(((key, depth + 1), (child, depth + 1))) else: raise failure("SOURCE_INVALID") if encoded_bytes > limits.max_value_bytes: raise failure("LIMIT_EXCEEDED") result = canonical_json_utf8_v1(value) if len(result) > limits.max_value_bytes: raise failure("LIMIT_EXCEEDED") return resultEncode exact JSON types, without coercion, subclass hooks or Unicode normalization.
Integers must be JS-safe; exact decimals are strings. Floats, tuples, subclasses and lone surrogates are deliberately outside this SDK profile.
def validate_materialization_target(prepared: ReportingPreparedRevision,
*,
binding: ReportingDestinationBinding,
revisions: Sequence[ReportingRevisionRecord]) ‑> None-
Expand source code
def validate_materialization_target( prepared: ReportingPreparedRevision, *, binding: ReportingDestinationBinding, revisions: Sequence[ReportingRevisionRecord], ) -> None: """Pure B2 finish seam. Caller must supply its locked, complete current history.""" selection = select_reporting_revision( revisions, account_id=prepared.request.principal.account_id, reporting_obligation_id=prepared.request.reporting_obligation_id, required_finality=prepared.obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if ( selection.kind != "selected" or not selection.revision.readable or selection.revision != prepared.revision or binding_fingerprint(binding) != prepared.request.binding_fingerprint ): raise failure("CURRENT_REVISION_CHANGED")Pure B2 finish seam. Caller must supply its locked, complete current history.
Classes
class InMemoryReportingMaterializerStore (**kwargs: Any)-
Expand source code
class InMemoryReportingMaterializerStore(InMemoryReportingReconciliationStore): _materializer_writer_epoch = 0 def __init__(self, **kwargs: Any) -> None: super().__init__(**kwargs) self._materializer_candidates: dict[ReportingDeliveryScope, _Candidate] = {} self._materializer_work: dict[tuple[str, str, str], _Work] = {} self._materializer_account_turns: dict[str, int] = {} self._materializer_turn = 0 self._materializer_outbox = ( MaterializerNotificationState() if self._notification_state is not None else None ) self._materializer_boundaries: list[ReportingMaterializerBoundary] = [] self._materializer_status_heads: dict[ReportingDeliveryPrincipal, int] = {} self._materializer_account_heads: dict[str, int] = {} def _wake(self, scope: ReportingDeliveryScope) -> None: candidate = self._materializer_candidates.get(scope) if candidate is None: candidate = _Candidate(scope) self._materializer_candidates[scope] = candidate else: candidate.generation += 1 candidate.due_at, candidate.reason = self._clock(), "ready" def _wake_obligation(self, account_id: str, obligation_id: str) -> None: obligation = self._obligations.get(obligation_id) if obligation is None or obligation.account_id != account_id: return for _, _, record in self._retained_delivery_records(): if isinstance(record, ReportingDestinationBinding) and ( record.generation_key == obligation.generation_key ): self._wake( ReportingDeliveryScope( obligation.generation_key, record.consumer_id, obligation_id ) ) def _append(self, account_id: str, kind: Any, record_id: str, *, consumer_id: str) -> None: super()._append(account_id, kind, record_id, consumer_id=consumer_id) if kind == "obligation": self._wake_obligation(account_id, record_id) elif kind == "revision": self._wake_obligation(account_id, self._revisions[record_id].reporting_obligation_id) async def put_configuration(self, configuration: ReportingConfiguration) -> None: async with self._mutation(): before = self._configurations.get(configuration.generation_key) await super().put_configuration(configuration) if before != configuration: for scope in self._materializer_candidates: if scope.generation_key == configuration.generation_key: self._wake(scope) async def set_revision_readable( self, *, account_id: str, reporting_revision_id: str, readable: bool, ) -> None: async with self._mutation(): before = self._revisions.get(reporting_revision_id) await super().set_revision_readable( account_id=account_id, reporting_revision_id=reporting_revision_id, readable=readable, ) revision = self._revisions.get(reporting_revision_id) if revision is not None and revision != before: self._wake_obligation(account_id, revision.reporting_obligation_id) def _commit_record_unlocked( self, record: RecordT, *, notify: bool = True, dirty: bool = True ) -> tuple[RecordT, bool]: # Only the fenced verified finish may enqueue a readiness intent. An # ordinary public outcome keeps the projection dirty work it has always # produced; suppressing that would silently stall the status projector. stored, added = super()._commit_record_unlocked( record, notify=notify and not isinstance(record, ReportingMaterializationRecord), dirty=dirty, ) if added: if isinstance(record, ReportingDestinationBinding): for obligation in self._obligations.values(): if obligation.generation_key == record.generation_key: self._wake( ReportingDeliveryScope( obligation.generation_key, record.consumer_id, obligation.reporting_obligation_id, ) ) elif isinstance(record, ReportingMaterializationCheck): self._wake(record.scope) return stored, added def _context(self, scope: ReportingDeliveryScope) -> MaterializerContext: records = tuple(c.record for c in self._caller_changes(scope.principal)) binding = next( ( r for r in records if isinstance(r, ReportingDestinationBinding) and r.generation_key == scope.generation_key ), None, ) obligation = self._obligations.get(scope.reporting_obligation_id) configuration = self._configurations.get(scope.generation_key) if ( binding is None or obligation is None or configuration is None or configuration.quarantined or obligation.generation_key != scope.generation_key ): raise failure("BINDING_MISMATCH") delivery = next( ( r for r in records if isinstance(r, ReportingObligationDeliveryRecord) and r.scope == scope ), None, ) return MaterializerContext( configuration, obligation, binding, delivery, tuple( r for r in self._revisions.values() if r.account_id == scope.principal.account_id and r.reporting_obligation_id == scope.reporting_obligation_id ), records, ) def _work_key(self, attempt: ReportingMaterializationAttempt) -> tuple[str, str, str]: return ( attempt.scope.principal.account_id, attempt.scope.consumer_id, attempt.reporting_materialization_id, ) def _park( self, candidate: _Candidate, reason: MaterializerReason, due: datetime | None = None, ) -> ReportingMaterializerTurn: candidate.reason, candidate.due_at = reason, due return ReportingMaterializerTurn("parked", reason) def _materializer_candidate_enabled(self, account_id: str) -> bool: return True @materializer_errors async def claim_materialization( self, *, keys: tuple[ReportingVerificationKey, ...], lease_seconds: int = 30, ) -> ReportingMaterializerLease | ReportingMaterializerTurn: validate_lease_seconds(lease_seconds) async with self._mutation(): now = self._clock() pending = {w.scope for w in self._materializer_work.values() if not w.acked} works = [ w for w in self._materializer_work.values() if not w.acked and w.due_at is not None and w.due_at <= now and (w.lease_until is None or w.lease_until <= now) ] candidates = [ c for c in self._materializer_candidates.values() if c.due_at is not None and c.due_at <= now and c.scope not in pending and self._materializer_candidate_enabled(c.scope.principal.account_id) ] accounts = {w.scope.principal.account_id for w in works} | { c.scope.principal.account_id for c in candidates } if not accounts: return ReportingMaterializerTurn("idle") account_id = min( accounts, key=lambda a: (self._materializer_account_turns.get(a, 0), a) ) self._materializer_turn += 1 self._materializer_account_turns[account_id] = self._materializer_turn account_work = [w for w in works if w.scope.principal.account_id == account_id] if account_work: work = min(account_work, key=lambda w: (w.due_at or now, w.request.external_id)) return self._lease(work, keys, lease_seconds) candidate = min( (c for c in candidates if c.scope.principal.account_id == account_id), key=lambda c: (c.turn, c.scope.consumer_id, c.scope.reporting_obligation_id), ) candidate.turn = self._materializer_turn try: context = self._context(candidate.scope) key = key_for(context.binding, context.obligation, keys) except LedgerConflictError: return self._park(candidate, "history_corrupt") except ReportingWriterError: return self._park(candidate, "component_unavailable") revision, reason = context.selection(now) if revision is None or reason != "ready": at = context.configuration.activated_at return self._park(candidate, reason, at if at is not None and at > now else None) attempts = context.attempts(revision.reporting_revision_id) if tuple(a.attempt for a in attempts) != tuple(range(1, len(attempts) + 1)): return self._park(candidate, "history_corrupt") if context.pending_attempts(): return self._park(candidate, "legacy_pending") if attempts: outcome = context.outcome(attempts[-1]) assert outcome is not None if outcome.status != "failed": reason, expires = context.retained_success(outcome, now) return self._park(candidate, reason, expires) owned = self._materializer_work.get(self._work_key(attempts[-1])) if owned is None or not owned.retry_allowed: return self._park( candidate, "legacy_terminal" if owned is None else "operator_required" ) try: admission_epoch = self._new_admission_epoch(context, key) except ReportingWriterError: return self._park(candidate, "component_unavailable") if context.delivery is None: if context.obligation.currency is None: return self._park(candidate, "operator_required") delivery = ReportingObligationDeliveryRecord( candidate.scope, context.obligation.currency, max(now, context.obligation.period.end) + timedelta(days=context.binding.resource_retention_days), now, ) self._commit_record_unlocked(delivery, notify=False, dirty=False) attempt = ReportingMaterializationAttempt( candidate.scope, revision.reporting_revision_id, "rpm_" + uuid4().hex, len(attempts) + 1, now, ) self._commit_record_unlocked(attempt, notify=False, dirty=False) request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) work = _Work( candidate.scope, candidate.generation, attempt, request, now, self._notification_state is not None, admission_epoch=admission_epoch, ) self._materializer_work[self._work_key(attempt)] = work self._park(candidate, "ready") return self._lease(work, keys, lease_seconds) def _lease( self, work: _Work, keys: tuple[ReportingVerificationKey, ...], seconds: int, ) -> ReportingMaterializerLease | ReportingMaterializerTurn: if work.admission_epoch > self._materializer_writer_epoch: raise failure("UNSUPPORTED_VERIFICATION") if work.notifications_enabled != (self._notification_state is not None): return self._park_work(work, "component_unavailable") try: context = self._context(work.scope) if not context.resumable(work.attempt): return self._park_work(work, "history_corrupt") key = key_for( context.binding, context.obligation, keys, required=verification_key_id(work.request.verification_key), ) if ( ReportingDestinationRequest.from_binding(context.binding, work.attempt, key) != work.request ): raise failure("BINDING_MISMATCH") except (ReportingWriterError, LedgerConflictError): return self._park_work(work, "component_unavailable") work.token = str(uuid4()) work.lease_until = self._clock() + timedelta(seconds=seconds) work.due_at = work.lease_until return ReportingMaterializerLease( work.scope, work.generation, work.token, work.lease_until, work.attempt, work.request, context, work.notifications_enabled, work.admission_epoch, ) def _park_work(self, work: _Work, reason: MaterializerReason) -> ReportingMaterializerTurn: work.due_at, work.reason = None, reason work.token, work.lease_until = None, None return self._park(self._materializer_candidates[work.scope], reason) def _held(self, lease: ReportingMaterializerLease) -> _Work | None: work = self._materializer_work.get(self._work_key(lease.attempt)) if ( work is None or work.acked or work.token != lease.token or work.lease_until is None or work.lease_until <= self._clock() ): return None if ( work.request != lease.request or work.generation != lease.generation or work.notifications_enabled != lease.notifications_enabled or work.notifications_enabled != (self._notification_state is not None) or work.admission_epoch != lease.admission_epoch or work.admission_epoch > self._materializer_writer_epoch ): raise failure("BINDING_MISMATCH") return work def _target( self, lease: ReportingMaterializerLease ) -> tuple[MaterializerContext, MaterializerReason]: context = self._context(lease.scope) selected, reason = context.selection(self._clock()) if reason != "ready": return context, reason candidate = self._materializer_candidates[lease.scope] if ( candidate.generation != lease.generation or selected is None or selected.reporting_revision_id != lease.attempt.reporting_revision_id ): return context, "target_changed" if binding_fingerprint(context.binding) != lease.request.binding_fingerprint: return context, "binding_changed" return context, "ready" @materializer_errors async def renew_materialization( self, lease: ReportingMaterializerLease, *, lease_seconds: int ) -> bool: validate_lease_seconds(lease_seconds) async with self._mutation(): work = self._held(lease) if work is None: return False work.lease_until = self._clock() + timedelta(seconds=lease_seconds) work.due_at = work.lease_until return True @materializer_errors async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None: async with self._mutation(): if self._held(lease) is None: raise failure("LEASE_LOST") _, reason = self._target(lease) if reason != "ready": raise failure( "HISTORY_CORRUPT" if reason == "history_corrupt" else "CURRENT_REVISION_CHANGED" ) @materializer_errors async def finish_materialization( self, lease: ReportingMaterializerLease, *, prepared: ReportingPreparedRevision | None = None, verified: ReportingVerifiedDestination | None = None, error: ReportingWriterFailure | None = None, ) -> ReportingMaterializerTurn: if verified is not None: validate_verified_destination(verified, lease.request, token=lease.token) if prepared is None or prepared.request != lease.request or error is not None: raise failure("BINDING_MISMATCH") elif error is None: raise failure("BINDING_MISMATCH") async with self._mutation(): work = self._held(lease) if work is None: previous = self._materializer_work.get(self._work_key(lease.attempt)) if ( previous is not None and previous.acked and previous.completion_token == lease.token and previous.request == lease.request ): return ReportingMaterializerTurn( "verified" if previous.reason == "verified" else "failed", previous.reason, lease.attempt.reporting_materialization_id, ) return ReportingMaterializerTurn( "pending", "effect_unknown", lease.attempt.reporting_materialization_id ) context, reason = self._target(lease) if reason != "ready": error = ReportingWriterFailure("CURRENT_REVISION_CHANGED", "new_attempt", "applied") verified = None elif verified is not None and prepared is not None: validate_materialization_target( prepared, binding=context.binding, revisions=context.revisions ) if error is not None and ( error.effect == "unknown" or error.retry == "same_identity" or error.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"} ): work.token, work.lease_until, work.reason = None, None, "effect_unknown" work.due_at = self._clock() + timedelta( seconds=max(1, min(error.retry_after_seconds or 5, 300)) ) return ReportingMaterializerTurn( "pending", work.reason, work.attempt.reporting_materialization_id ) now = self._clock() if verified is not None and ( context.delivery is None or verified.resource.expires_at < max( context.delivery.resource_retained_until, now + timedelta(days=context.binding.resource_retention_days), ) ): error = ReportingWriterFailure("RESOURCE_UNAVAILABLE", "new_attempt", "applied") verified = None if verified is not None: # Validate again under the finish lock, immediately before # copying observations into the immutable terminal record. validate_verified_destination(verified, lease.request, token=lease.token) outcome = ReportingMaterializationRecord( lease.scope, lease.attempt.reporting_revision_id, lease.attempt.reporting_materialization_id, context.binding.success_status if verified is not None else "failed", now, verified.resource if verified is not None else None, replace(verified.verification, verified_at=now) if verified is not None else None, public_failure(error) if error is not None else None, ) stored, inserted = self._commit_record_unlocked(outcome, notify=False, dirty=False) self._materializer_dirty(stored, context) if inserted and verified is not None and self._notification_state is not None: revision = next( r for r in context.revisions if r.reporting_revision_id == stored.reporting_revision_id ) event = materialization_event( stored, context.records, context.obligation, revision, context.configuration, now, ) if event is None: raise failure("BINDING_MISMATCH") self._enqueue_materializer(event, lease) if self._held(lease) is None: raise failure("LEASE_LOST") work.completion_token = work.token work.acked, work.token, work.lease_until = True, None, None work.retry_allowed = error is not None and error.retry == "new_attempt" if reason == "ready": reason = ( "verified" if verified is not None else ("retry" if work.retry_allowed else "operator_required") ) work.reason = reason due = ( now + timedelta(seconds=max(1, error.retry_after_seconds or 5)) if work.retry_allowed and error else None ) if reason == "target_changed": due = now elif reason not in {"retry", "verified"}: due = None if ( reason == "inactive" and context.configuration.activated_at is not None and context.configuration.activated_at > now ): due = context.configuration.activated_at if verified is not None: due = verified.resource.expires_at self._park(self._materializer_candidates[lease.scope], reason, due) return ReportingMaterializerTurn( "verified" if verified is not None else "failed", reason, stored.reporting_materialization_id, ) def _new_admission_epoch( self, context: MaterializerContext, key: ReportingVerificationKey ) -> int: return 0 def _enqueue_materializer( self, event: ReportingDomainEvent, lease: ReportingMaterializerLease ) -> None: if lease.admission_epoch != 0: raise failure("UNSUPPORTED_VERIFICATION") assert self._materializer_outbox is not None self._materializer_outbox.enqueue(event) if ( self._materializer_outbox.events.get( (event.account_id, event.consumer_namespace, event.notification_id) ) != event ): raise failure("BINDING_MISMATCH") def _materializer_dirty( self, outcome: ReportingMaterializationRecord, context: MaterializerContext ) -> None: scope, reason, evidence = delivery_dirty(outcome, context.obligation) self._dirty_status(scope, reason, after=evidence) from adcp.reporting.ledger.status_snapshot import settle_memory_snapshot core = settle_memory_snapshot(self, scope.account_id) core = replace(core, as_of=outcome.completed_at) heads, boundaries = self._materializer_capture_collections(outcome) sequence = heads.get(outcome.scope.principal, 0) + 1 heads[outcome.scope.principal] = sequence account_sequence = self._materializer_account_heads.get(scope.account_id, 0) + 1 self._materializer_account_heads[scope.account_id] = account_sequence boundaries.append( ReportingMaterializerBoundary( outcome.scope.principal, sequence, account_sequence, outcome.reporting_materialization_id, outcome.completed_at, private_snapshot(core, outcome.scope.principal), tuple(c.record for c in self._caller_changes(outcome.scope.principal)), ) ) def _materializer_capture_collections( self, outcome: ReportingMaterializationRecord ) -> tuple[dict[ReportingDeliveryPrincipal, int], list[ReportingMaterializerBoundary]]: return self._materializer_status_heads, self._materializer_boundaries @materializer_errors async def read_materializer_boundaries( self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100 ) -> tuple[ReportingMaterializerBoundary, ...]: if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100: raise ReportingMaterializerUsageError( "materializer boundary reads require bounded positions" ) async with self._lock: return tuple( b for b in self._materializer_boundaries if b.caller == caller and b.sequence > after )[:limit] @materializer_errors async def import_pending_materialization( self, *, scope: ReportingDeliveryScope, reporting_materialization_id: str, original_external_id: str, keys: tuple[ReportingVerificationKey, ...], ) -> None: async with self._mutation(): context = self._context(scope) attempt = next( ( r for r in context.records if isinstance(r, ReportingMaterializationAttempt) and r.reporting_materialization_id == reporting_materialization_id and r.scope == scope ), None, ) if attempt is None or context.outcome(attempt) is not None: raise failure("BINDING_MISMATCH") if not context.resumable(attempt): raise failure("HISTORY_CORRUPT") key = key_for(context.binding, context.obligation, keys) request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) if original_external_id != request.external_id: raise failure("BINDING_MISMATCH") if scope not in self._materializer_candidates: self._wake(scope) existing = self._materializer_work.get(self._work_key(attempt)) if existing is not None: if existing.token is None and not existing.acked: existing.due_at = self._clock() return if any(w.scope == scope and not w.acked for w in self._materializer_work.values()): raise failure("BINDING_MISMATCH") self._materializer_work[self._work_key(attempt)] = _Work( scope, self._materializer_candidates[scope].generation, attempt, request, self._clock(), self._notification_state is not None, )Optional reference extension. No destination/receipt services at construction.
Ancestors
- InMemoryReportingReconciliationStore
- InMemoryReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
Subclasses
Methods
-
Expand source code
@materializer_errors async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None: async with self._mutation(): if self._held(lease) is None: raise failure("LEASE_LOST") _, reason = self._target(lease) if reason != "ready": raise failure( "HISTORY_CORRUPT" if reason == "history_corrupt" else "CURRENT_REVISION_CHANGED" ) async def claim_materialization(self,
*,
keys: tuple[ReportingVerificationKey, ...],
lease_seconds: int = 30) ‑> ReportingMaterializerLease | ReportingMaterializerTurn-
Expand source code
@materializer_errors async def claim_materialization( self, *, keys: tuple[ReportingVerificationKey, ...], lease_seconds: int = 30, ) -> ReportingMaterializerLease | ReportingMaterializerTurn: validate_lease_seconds(lease_seconds) async with self._mutation(): now = self._clock() pending = {w.scope for w in self._materializer_work.values() if not w.acked} works = [ w for w in self._materializer_work.values() if not w.acked and w.due_at is not None and w.due_at <= now and (w.lease_until is None or w.lease_until <= now) ] candidates = [ c for c in self._materializer_candidates.values() if c.due_at is not None and c.due_at <= now and c.scope not in pending and self._materializer_candidate_enabled(c.scope.principal.account_id) ] accounts = {w.scope.principal.account_id for w in works} | { c.scope.principal.account_id for c in candidates } if not accounts: return ReportingMaterializerTurn("idle") account_id = min( accounts, key=lambda a: (self._materializer_account_turns.get(a, 0), a) ) self._materializer_turn += 1 self._materializer_account_turns[account_id] = self._materializer_turn account_work = [w for w in works if w.scope.principal.account_id == account_id] if account_work: work = min(account_work, key=lambda w: (w.due_at or now, w.request.external_id)) return self._lease(work, keys, lease_seconds) candidate = min( (c for c in candidates if c.scope.principal.account_id == account_id), key=lambda c: (c.turn, c.scope.consumer_id, c.scope.reporting_obligation_id), ) candidate.turn = self._materializer_turn try: context = self._context(candidate.scope) key = key_for(context.binding, context.obligation, keys) except LedgerConflictError: return self._park(candidate, "history_corrupt") except ReportingWriterError: return self._park(candidate, "component_unavailable") revision, reason = context.selection(now) if revision is None or reason != "ready": at = context.configuration.activated_at return self._park(candidate, reason, at if at is not None and at > now else None) attempts = context.attempts(revision.reporting_revision_id) if tuple(a.attempt for a in attempts) != tuple(range(1, len(attempts) + 1)): return self._park(candidate, "history_corrupt") if context.pending_attempts(): return self._park(candidate, "legacy_pending") if attempts: outcome = context.outcome(attempts[-1]) assert outcome is not None if outcome.status != "failed": reason, expires = context.retained_success(outcome, now) return self._park(candidate, reason, expires) owned = self._materializer_work.get(self._work_key(attempts[-1])) if owned is None or not owned.retry_allowed: return self._park( candidate, "legacy_terminal" if owned is None else "operator_required" ) try: admission_epoch = self._new_admission_epoch(context, key) except ReportingWriterError: return self._park(candidate, "component_unavailable") if context.delivery is None: if context.obligation.currency is None: return self._park(candidate, "operator_required") delivery = ReportingObligationDeliveryRecord( candidate.scope, context.obligation.currency, max(now, context.obligation.period.end) + timedelta(days=context.binding.resource_retention_days), now, ) self._commit_record_unlocked(delivery, notify=False, dirty=False) attempt = ReportingMaterializationAttempt( candidate.scope, revision.reporting_revision_id, "rpm_" + uuid4().hex, len(attempts) + 1, now, ) self._commit_record_unlocked(attempt, notify=False, dirty=False) request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) work = _Work( candidate.scope, candidate.generation, attempt, request, now, self._notification_state is not None, admission_epoch=admission_epoch, ) self._materializer_work[self._work_key(attempt)] = work self._park(candidate, "ready") return self._lease(work, keys, lease_seconds) async def finish_materialization(self,
lease: ReportingMaterializerLease,
*,
prepared: ReportingPreparedRevision | None = None,
verified: ReportingVerifiedDestination | None = None,
error: ReportingWriterFailure | None = None) ‑> ReportingMaterializerTurn-
Expand source code
@materializer_errors async def finish_materialization( self, lease: ReportingMaterializerLease, *, prepared: ReportingPreparedRevision | None = None, verified: ReportingVerifiedDestination | None = None, error: ReportingWriterFailure | None = None, ) -> ReportingMaterializerTurn: if verified is not None: validate_verified_destination(verified, lease.request, token=lease.token) if prepared is None or prepared.request != lease.request or error is not None: raise failure("BINDING_MISMATCH") elif error is None: raise failure("BINDING_MISMATCH") async with self._mutation(): work = self._held(lease) if work is None: previous = self._materializer_work.get(self._work_key(lease.attempt)) if ( previous is not None and previous.acked and previous.completion_token == lease.token and previous.request == lease.request ): return ReportingMaterializerTurn( "verified" if previous.reason == "verified" else "failed", previous.reason, lease.attempt.reporting_materialization_id, ) return ReportingMaterializerTurn( "pending", "effect_unknown", lease.attempt.reporting_materialization_id ) context, reason = self._target(lease) if reason != "ready": error = ReportingWriterFailure("CURRENT_REVISION_CHANGED", "new_attempt", "applied") verified = None elif verified is not None and prepared is not None: validate_materialization_target( prepared, binding=context.binding, revisions=context.revisions ) if error is not None and ( error.effect == "unknown" or error.retry == "same_identity" or error.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"} ): work.token, work.lease_until, work.reason = None, None, "effect_unknown" work.due_at = self._clock() + timedelta( seconds=max(1, min(error.retry_after_seconds or 5, 300)) ) return ReportingMaterializerTurn( "pending", work.reason, work.attempt.reporting_materialization_id ) now = self._clock() if verified is not None and ( context.delivery is None or verified.resource.expires_at < max( context.delivery.resource_retained_until, now + timedelta(days=context.binding.resource_retention_days), ) ): error = ReportingWriterFailure("RESOURCE_UNAVAILABLE", "new_attempt", "applied") verified = None if verified is not None: # Validate again under the finish lock, immediately before # copying observations into the immutable terminal record. validate_verified_destination(verified, lease.request, token=lease.token) outcome = ReportingMaterializationRecord( lease.scope, lease.attempt.reporting_revision_id, lease.attempt.reporting_materialization_id, context.binding.success_status if verified is not None else "failed", now, verified.resource if verified is not None else None, replace(verified.verification, verified_at=now) if verified is not None else None, public_failure(error) if error is not None else None, ) stored, inserted = self._commit_record_unlocked(outcome, notify=False, dirty=False) self._materializer_dirty(stored, context) if inserted and verified is not None and self._notification_state is not None: revision = next( r for r in context.revisions if r.reporting_revision_id == stored.reporting_revision_id ) event = materialization_event( stored, context.records, context.obligation, revision, context.configuration, now, ) if event is None: raise failure("BINDING_MISMATCH") self._enqueue_materializer(event, lease) if self._held(lease) is None: raise failure("LEASE_LOST") work.completion_token = work.token work.acked, work.token, work.lease_until = True, None, None work.retry_allowed = error is not None and error.retry == "new_attempt" if reason == "ready": reason = ( "verified" if verified is not None else ("retry" if work.retry_allowed else "operator_required") ) work.reason = reason due = ( now + timedelta(seconds=max(1, error.retry_after_seconds or 5)) if work.retry_allowed and error else None ) if reason == "target_changed": due = now elif reason not in {"retry", "verified"}: due = None if ( reason == "inactive" and context.configuration.activated_at is not None and context.configuration.activated_at > now ): due = context.configuration.activated_at if verified is not None: due = verified.resource.expires_at self._park(self._materializer_candidates[lease.scope], reason, due) return ReportingMaterializerTurn( "verified" if verified is not None else "failed", reason, stored.reporting_materialization_id, ) async def import_pending_materialization(self,
*,
scope: ReportingDeliveryScope,
reporting_materialization_id: str,
original_external_id: str,
keys: tuple[ReportingVerificationKey, ...]) ‑> None-
Expand source code
@materializer_errors async def import_pending_materialization( self, *, scope: ReportingDeliveryScope, reporting_materialization_id: str, original_external_id: str, keys: tuple[ReportingVerificationKey, ...], ) -> None: async with self._mutation(): context = self._context(scope) attempt = next( ( r for r in context.records if isinstance(r, ReportingMaterializationAttempt) and r.reporting_materialization_id == reporting_materialization_id and r.scope == scope ), None, ) if attempt is None or context.outcome(attempt) is not None: raise failure("BINDING_MISMATCH") if not context.resumable(attempt): raise failure("HISTORY_CORRUPT") key = key_for(context.binding, context.obligation, keys) request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) if original_external_id != request.external_id: raise failure("BINDING_MISMATCH") if scope not in self._materializer_candidates: self._wake(scope) existing = self._materializer_work.get(self._work_key(attempt)) if existing is not None: if existing.token is None and not existing.acked: existing.due_at = self._clock() return if any(w.scope == scope and not w.acked for w in self._materializer_work.values()): raise failure("BINDING_MISMATCH") self._materializer_work[self._work_key(attempt)] = _Work( scope, self._materializer_candidates[scope].generation, attempt, request, self._clock(), self._notification_state is not None, ) async def put_configuration(self, configuration: ReportingConfiguration) ‑> None-
Expand source code
async def put_configuration(self, configuration: ReportingConfiguration) -> None: async with self._mutation(): before = self._configurations.get(configuration.generation_key) await super().put_configuration(configuration) if before != configuration: for scope in self._materializer_candidates: if scope.generation_key == configuration.generation_key: self._wake(scope) async def read_materializer_boundaries(self,
*,
caller: ReportingDeliveryPrincipal,
after: int = 0,
limit: int = 100) ‑> tuple[ReportingMaterializerBoundary, ...]-
Expand source code
@materializer_errors async def read_materializer_boundaries( self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100 ) -> tuple[ReportingMaterializerBoundary, ...]: if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100: raise ReportingMaterializerUsageError( "materializer boundary reads require bounded positions" ) async with self._lock: return tuple( b for b in self._materializer_boundaries if b.caller == caller and b.sequence > after )[:limit] async def renew_materialization(self,
lease: ReportingMaterializerLease,
*,
lease_seconds: int) ‑> bool-
Expand source code
@materializer_errors async def renew_materialization( self, lease: ReportingMaterializerLease, *, lease_seconds: int ) -> bool: validate_lease_seconds(lease_seconds) async with self._mutation(): work = self._held(lease) if work is None: return False work.lease_until = self._clock() + timedelta(seconds=lease_seconds) work.due_at = work.lease_until return True async def set_revision_readable(self, *, account_id: str, reporting_revision_id: str, readable: bool) ‑> None-
Expand source code
async def set_revision_readable( self, *, account_id: str, reporting_revision_id: str, readable: bool, ) -> None: async with self._mutation(): before = self._revisions.get(reporting_revision_id) await super().set_revision_readable( account_id=account_id, reporting_revision_id=reporting_revision_id, readable=readable, ) revision = self._revisions.get(reporting_revision_id) if revision is not None and revision != before: self._wake_obligation(account_id, revision.reporting_obligation_id)
Inherited members
class PgReportingMaterializerStore (*,
pool: AsyncConnectionPool,
clock: Callable[[], datetime] | None = None,
notifications: bool = False)-
Expand source code
class PgReportingMaterializerStore(PgReportingReconciliationStore): """Additive durable extension. No account enumeration or destination I/O. Each turn samples at most sixteen due accounts without row locks, tries the Core advisory lock first, then claims within that account. Long work returns its connection before any source/destination call; even a size-one pool can renew its lease. Clock overrides affect old conformance APIs only: work time is always PostgreSQL time. """ # Only a read-only sampling position, never a lease or durable correctness # boundary. Walk past busy accounts in bounded pages, wrapping after the end. # Restart may repeat a page; it cannot lose or acknowledge work. _materializer_sample_after: tuple[str, str] | None = None async def _commit_record_on( self, connection: Any, record: RecordT, *, notify: bool = True, dirty: bool = True ) -> tuple[RecordT, bool]: # Only the fenced verified finish may enqueue a readiness intent. An # ordinary public outcome keeps the projection dirty work it has always # produced; suppressing that would silently stall the status projector. return await super()._commit_record_on( connection, record, notify=notify and not isinstance(record, ReportingMaterializationRecord), dirty=dirty, ) @materializer_errors async def create_schema(self) -> None: async with self._connection() as connection, connection.transaction(): await self._create_schema_on(connection) await connection.execute( files("adcp.reporting.ledger").joinpath("reporting_materializer.sql").read_text() ) @materializer_errors async def materializer_ready(self) -> bool: async with self._connection() as connection: await validate_materializer_schema( connection, notifications=self._notifications_enabled ) return True async def _materializer_context_on( self, connection: Any, scope: ReportingDeliveryScope ) -> MaterializerContext: records = await self._records(connection, scope.principal) binding = next( ( r for r in records if isinstance(r, ReportingDestinationBinding) and r.generation_key == scope.generation_key ), None, ) obligation_row = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # nosec B608 " WHERE account_id=%s AND reporting_obligation_id=%s", (scope.principal.account_id, scope.reporting_obligation_id), ) ).fetchone() config_row = await ( await connection.execute( "SELECT delivery_config_id, delivery_config_version, account_id," " report_definition_id, reporting_profile, feed_purpose, required_finality," " account_timezone, schedule, media_buy_ids, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, definition," " authoritative_party, consumer_id, quarantined" " FROM reporting_configurations WHERE account_id=%s" " AND consumer_id=%s AND delivery_config_id=%s AND delivery_config_version=%s", ( scope.principal.account_id, scope.consumer_id, scope.generation_key.delivery_config_id, scope.generation_key.delivery_config_version, ), ) ).fetchone() revision_rows = await ( await connection.execute( f"SELECT {_REVISION_COLUMNS} FROM reporting_revisions" # nosec B608 " WHERE account_id=%s AND reporting_obligation_id=%s", (scope.principal.account_id, scope.reporting_obligation_id), ) ).fetchall() if binding is None or obligation_row is None or config_row is None: raise failure("BINDING_MISMATCH") obligation = _obligation_from_row(obligation_row) configuration = _configuration_from_row(config_row) if configuration.quarantined or obligation.generation_key != scope.generation_key: raise failure("BINDING_MISMATCH") delivery = next( ( r for r in records if isinstance(r, ReportingObligationDeliveryRecord) and r.scope == scope ), None, ) return MaterializerContext( configuration, obligation, binding, delivery, tuple(_revision_from_row(r) for r in revision_rows), records, ) async def _schedule_account_on(self, connection: Any, account_id: str) -> None: await connection.execute( "UPDATE reporting_materializer_accounts SET due_at=(SELECT min(due) FROM (" " SELECT due_at AS due FROM reporting_materializer_work" " WHERE account_id=%s AND state='pending'" " UNION ALL SELECT c.due_at FROM reporting_materializer_candidates c" " WHERE c.account_id=%s AND c.due_at IS NOT NULL AND NOT EXISTS (" " SELECT 1 FROM reporting_materializer_work w WHERE w.account_id=c.account_id" " AND w.consumer_id=c.consumer_id AND w.delivery_config_id=c.delivery_config_id" " AND w.delivery_config_version=c.delivery_config_version" " AND w.reporting_obligation_id=c.reporting_obligation_id AND w.state='pending')" " UNION ALL SELECT clock_timestamp() FROM reporting_materializer_discovery" " WHERE account_id=%s AND NOT complete) ready) WHERE account_id=%s", (account_id,) * 4, ) async def _discover_on(self, connection: Any, account_id: str) -> bool: row = await ( await connection.execute( "SELECT consumer_id, delivery_config_id, delivery_config_version," " after_obligation_id" " FROM reporting_materializer_discovery WHERE account_id=%s AND NOT complete" " ORDER BY consumer_id, delivery_config_id, delivery_config_version" " LIMIT 1 FOR UPDATE", (account_id,), ) ).fetchone() if row is None: return False consumer, config, version, after = row obligations = await ( await connection.execute( "SELECT reporting_obligation_id FROM reporting_obligations WHERE account_id=%s" " AND consumer_id=%s AND delivery_config_id=%s AND delivery_config_version=%s" " AND reporting_obligation_id>%s ORDER BY reporting_obligation_id LIMIT 32", (account_id, consumer, config, version, after), ) ).fetchall() for (obligation_id,) in obligations: # Only missing candidates: a later backfill batch never invalidates # a worker already reserved through the publication trigger. await connection.execute( "INSERT INTO reporting_materializer_candidates" " (account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id,due_at) VALUES (%s,%s,%s,%s,%s,clock_timestamp())" " ON CONFLICT DO NOTHING", (account_id, consumer, config, version, obligation_id), ) await connection.execute( "UPDATE reporting_materializer_discovery SET after_obligation_id=%s,complete=%s" " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s" " AND delivery_config_version=%s", ( obligations[-1][0] if obligations else after, len(obligations) < 32, account_id, consumer, config, version, ), ) return True async def _park_on( self, connection: Any, scope: ReportingDeliveryScope, reason: MaterializerReason, *, due_at: datetime | None = None, ) -> ReportingMaterializerTurn: await connection.execute( "UPDATE reporting_materializer_candidates SET reason=%s,due_at=%s," f" served_at=clock_timestamp() WHERE {_SCOPE_WHERE}", # nosec B608 (reason, due_at, *_scope_args(scope)), ) return ReportingMaterializerTurn("parked", reason) @materializer_errors async def claim_materialization( self, *, keys: tuple[ReportingVerificationKey, ...], lease_seconds: int = 30 ) -> ReportingMaterializerLease | ReportingMaterializerTurn: validate_lease_seconds(lease_seconds) await self.materializer_ready() async with self._connection() as connection: # Read-only sampling. Never lock global candidate rows before the # account advisory lock, including when a competing worker is busy. for _ in range(2): # Keep the cursor serialized, but order native timestamps: the # served_at::text output alias otherwise makes fairness locale-dependent. after = self._materializer_sample_after if after: query = ( "SELECT account_id,served_at::text FROM reporting_materializer_accounts" " WHERE due_at<=%s AND (served_at,account_id)>(%s::timestamptz,%s)" " ORDER BY reporting_materializer_accounts.served_at,account_id LIMIT 16" ) else: query = ( "SELECT account_id,served_at::text FROM reporting_materializer_accounts" " WHERE due_at<=%s" " ORDER BY reporting_materializer_accounts.served_at,account_id LIMIT 16" ) accounts = await ( await connection.execute( query, (await _now(connection), *(after or ())), ) ).fetchall() if accounts: break self._materializer_sample_after = None if after is None: break for account_id, served_at in accounts: self._materializer_sample_after = (served_at, account_id) async with self._connection() as connection, connection.transaction(): held = await ( await connection.execute( "SELECT pg_try_advisory_xact_lock(hashtext(%s))", (f"adcp.reporting:{account_id}",), ) ).fetchone() if held is None or not held[0]: continue await self._lock_account(connection, account_id) await validate_materializer_schema( connection, notifications=self._notifications_enabled ) await connection.execute( "UPDATE reporting_materializer_accounts SET served_at=clock_timestamp()" " WHERE account_id=%s", (account_id,), ) result = await self._claim_account_on(connection, account_id, keys, lease_seconds) await self._schedule_account_on(connection, account_id) # A committed account turn moves its durable served_at rank. Start # there again so continuously due peers cannot hide newly enrolled # accounts behind this pre-update cursor. Lock misses above retain # the continuation, allowing the next bounded page past a busy prefix. self._materializer_sample_after = None return result return ReportingMaterializerTurn("idle") async def _claim_account_on( self, connection: Any, account_id: str, keys: tuple[ReportingVerificationKey, ...], lease_seconds: int, ) -> ReportingMaterializerLease | ReportingMaterializerTurn: # A bound DB instant is an indexable range predicate. clock_timestamp() # is volatile and cannot be used as a PostgreSQL index scan boundary. now = await _now(connection) pending = await ( await connection.execute( "SELECT to_jsonb(w) FROM reporting_materializer_work w WHERE account_id=%s" " AND state='pending' AND due_at<=%s" " AND (lease_until IS NULL OR lease_until<=%s)" " ORDER BY due_at,reporting_materialization_id LIMIT 1 FOR UPDATE", (account_id, now, now), ) ).fetchone() if pending is not None: return await self._lease_on(connection, pending[0], keys, lease_seconds) discovered = await self._discover_on(connection, account_id) row = await ( await connection.execute( "SELECT to_jsonb(c) FROM reporting_materializer_candidates c WHERE account_id=%s" " AND due_at<=%s AND NOT EXISTS (" " SELECT 1 FROM reporting_materializer_work w WHERE w.account_id=c.account_id" " AND w.consumer_id=c.consumer_id AND w.delivery_config_id=c.delivery_config_id" " AND w.delivery_config_version=c.delivery_config_version" " AND w.reporting_obligation_id=c.reporting_obligation_id AND w.state='pending')" " ORDER BY served_at,consumer_id,reporting_obligation_id LIMIT 1 FOR UPDATE", (account_id, await _now(connection)), ) ).fetchone() if row is None: return ReportingMaterializerTurn("discovered" if discovered else "idle") scope = _scope(row[0]) try: context = await self._materializer_context_on(connection, scope) key = key_for(context.binding, context.obligation, keys) except LedgerConflictError: return await self._park_on(connection, scope, "history_corrupt") except ReportingWriterError: return await self._park_on(connection, scope, "component_unavailable") now = await _now(connection) revision, reason = context.selection(now) if reason != "ready" or revision is None: due = context.configuration.activated_at return await self._park_on( connection, scope, reason, due_at=due if due is not None and due > now else None ) attempts = context.attempts(revision.reporting_revision_id) if tuple(a.attempt for a in attempts) != tuple(range(1, len(attempts) + 1)): return await self._park_on(connection, scope, "history_corrupt") if context.pending_attempts(): # Includes legacy pending effects on an older selected revision. # Publication cannot make their unknown external history disappear. return await self._park_on(connection, scope, "legacy_pending") if attempts: outcome = context.outcome(attempts[-1]) assert outcome is not None if outcome.status != "failed": reason, expires = context.retained_success(outcome, now) return await self._park_on(connection, scope, reason, due_at=expires) owned = await ( await connection.execute( "SELECT retry_allowed FROM reporting_materializer_work WHERE account_id=%s" " AND consumer_id=%s AND reporting_materialization_id=%s AND state='acked'", (account_id, scope.consumer_id, attempts[-1].reporting_materialization_id), ) ).fetchone() if owned is None or not owned[0]: return await self._park_on( connection, scope, "legacy_terminal" if owned is None else "operator_required" ) delivery = context.delivery if delivery is None: if context.obligation.currency is None: return await self._park_on(connection, scope, "operator_required") delivery = ReportingObligationDeliveryRecord( scope, context.obligation.currency, max(now, context.obligation.period.end) + timedelta(days=context.binding.resource_retention_days), now, ) await self._commit_record_on(connection, delivery, notify=False, dirty=False) attempt = ReportingMaterializationAttempt( scope, revision.reporting_revision_id, "rpm_" + uuid4().hex, len(attempts) + 1, now ) await self._commit_record_on(connection, attempt, notify=False, dirty=False) request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) inserted = await ( await connection.execute( "INSERT INTO reporting_materializer_work" " (account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id,reporting_revision_id,reporting_materialization_id," " generation,binding_sha256,verification_key_sha256,external_id," " notifications_enabled)" " VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)" " RETURNING to_jsonb(reporting_materializer_work)", ( *_scope_args(scope), attempt.reporting_revision_id, attempt.reporting_materialization_id, row[0]["generation"], request.binding_fingerprint, verification_key_id(key), request.external_id, self._notifications_enabled, ), ) ).fetchone() assert inserted is not None await self._park_on(connection, scope, "ready") return await self._lease_on(connection, inserted[0], keys, lease_seconds) async def _lease_on( self, connection: Any, work: dict[str, Any], keys: tuple[ReportingVerificationKey, ...], lease_seconds: int, ) -> ReportingMaterializerLease | ReportingMaterializerTurn: scope = _scope(work) try: if work["notifications_enabled"] != self._notifications_enabled: raise failure("UNSUPPORTED_VERIFICATION") context = await self._materializer_context_on(connection, scope) key = key_for( context.binding, context.obligation, keys, required=work["verification_key_sha256"] ) attempt = next( r for r in context.records if isinstance(r, ReportingMaterializationAttempt) and r.reporting_materialization_id == work["reporting_materialization_id"] ) if not context.resumable(attempt): return await self._park_work_on(connection, work, "history_corrupt") request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) if ( request.external_id != work["external_id"] or request.binding_fingerprint != work["binding_sha256"] ): raise failure("BINDING_MISMATCH") except (LedgerConflictError, ReportingWriterError, StopIteration): # Park pending unknown effects until an operator restores the exact # installed contract. No hot loop and no replacement identity. return await self._park_work_on(connection, work, "component_unavailable") token = uuid4() updated = await ( await connection.execute( "UPDATE reporting_materializer_work SET lease_token=%s," " lease_until=clock_timestamp()+make_interval(secs=>%s)," " due_at=clock_timestamp()+make_interval(secs=>%s)" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s" " AND state='pending' AND (lease_until IS NULL OR lease_until<=clock_timestamp())" " RETURNING lease_until", ( token, lease_seconds, lease_seconds, scope.principal.account_id, scope.consumer_id, work["reporting_materialization_id"], ), ) ).fetchone() if updated is None: return ReportingMaterializerTurn("idle") return ReportingMaterializerLease( scope, work["generation"], str(token), updated[0], attempt, request, context, work["notifications_enabled"], work["admission_epoch"], ) async def _park_work_on( self, connection: Any, work: dict[str, Any], reason: MaterializerReason ) -> ReportingMaterializerTurn: scope = _scope(work) await connection.execute( "UPDATE reporting_materializer_work SET due_at='infinity'," " lease_token=NULL,lease_until=NULL,reason=%s" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s", ( reason, scope.principal.account_id, scope.consumer_id, work["reporting_materialization_id"], ), ) return await self._park_on(connection, scope, reason) async def _held_on( self, connection: Any, lease: ReportingMaterializerLease ) -> dict[str, Any] | None: row = await ( await connection.execute( "SELECT to_jsonb(w) FROM reporting_materializer_work w WHERE account_id=%s" " AND consumer_id=%s AND reporting_materialization_id=%s AND state='pending'" " AND lease_token=%s::uuid AND lease_until>clock_timestamp() FOR UPDATE", ( lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, lease.token, ), ) ).fetchone() if row is None: return None work = cast(dict[str, Any], row[0]) if ( work["external_id"] != lease.request.external_id or work["generation"] != lease.generation or work["binding_sha256"] != lease.request.binding_fingerprint or work["notifications_enabled"] != lease.notifications_enabled or work["notifications_enabled"] != self._notifications_enabled or work["admission_epoch"] != lease.admission_epoch ): raise failure("BINDING_MISMATCH") return work @materializer_errors async def renew_materialization( self, lease: ReportingMaterializerLease, *, lease_seconds: int ) -> bool: validate_lease_seconds(lease_seconds) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, lease.scope.principal.account_id) if await self._held_on(connection, lease) is None: return False await connection.execute( "UPDATE reporting_materializer_work SET lease_until=clock_timestamp()" " +make_interval(secs=>%s),due_at=clock_timestamp()+make_interval(secs=>%s)" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s", ( lease_seconds, lease_seconds, lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, ), ) await self._schedule_account_on(connection, lease.scope.principal.account_id) return True async def _target_on( self, connection: Any, lease: ReportingMaterializerLease, ) -> tuple[MaterializerContext, MaterializerReason]: context = await self._materializer_context_on(connection, lease.scope) selected, reason = context.selection(await _now(connection)) if reason != "ready": return context, reason row = await ( await connection.execute( f"SELECT generation FROM reporting_materializer_candidates WHERE {_SCOPE_WHERE}", # nosec B608 _scope_args(lease.scope), ) ).fetchone() if ( row is None or row[0] != lease.generation or selected is None or selected.reporting_revision_id != lease.attempt.reporting_revision_id ): return context, "target_changed" if binding_fingerprint(context.binding) != lease.request.binding_fingerprint: return context, "binding_changed" return context, "ready" @materializer_errors async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None: async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, lease.scope.principal.account_id) await validate_materializer_schema( connection, notifications=self._notifications_enabled ) if await self._held_on(connection, lease) is None: raise failure("LEASE_LOST") _, reason = await self._target_on(connection, lease) if reason != "ready": raise failure( "HISTORY_CORRUPT" if reason == "history_corrupt" else "CURRENT_REVISION_CHANGED" ) @materializer_errors async def finish_materialization( self, lease: ReportingMaterializerLease, *, prepared: ReportingPreparedRevision | None = None, verified: ReportingVerifiedDestination | None = None, error: ReportingWriterFailure | None = None, ) -> ReportingMaterializerTurn: if verified is not None: validate_verified_destination(verified, lease.request, token=lease.token) if prepared is None or prepared.request != lease.request or error is not None: raise failure("BINDING_MISMATCH") elif error is None: raise failure("BINDING_MISMATCH") async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, lease.scope.principal.account_id) await validate_materializer_schema( connection, notifications=self._notifications_enabled ) held = await self._held_on(connection, lease) if held is None: replay = await ( await connection.execute( "SELECT reason FROM reporting_materializer_work WHERE account_id=%s" " AND consumer_id=%s AND reporting_materialization_id=%s AND state='acked'" " AND completion_token=%s::uuid AND external_id=%s", ( lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, lease.token, lease.request.external_id, ), ) ).fetchone() if replay is not None: return ReportingMaterializerTurn( "verified" if replay[0] == "verified" else "failed", replay[0], lease.attempt.reporting_materialization_id, ) return ReportingMaterializerTurn( "pending", "effect_unknown", lease.attempt.reporting_materialization_id ) context, reason = await self._target_on(connection, lease) if reason != "ready": error = ReportingWriterFailure("CURRENT_REVISION_CHANGED", "new_attempt", "applied") verified = None elif verified is not None and prepared is not None: validate_materialization_target( prepared, binding=context.binding, revisions=context.revisions ) if error is not None and ( error.effect == "unknown" or error.retry == "same_identity" or error.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"} ): delay = max(1, min(error.retry_after_seconds or 5, 300)) await connection.execute( "UPDATE reporting_materializer_work SET lease_token=NULL,lease_until=NULL," " due_at=clock_timestamp()+make_interval(secs=>%s),reason='effect_unknown'" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s", ( delay, lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, ), ) await self._schedule_account_on(connection, lease.scope.principal.account_id) return ReportingMaterializerTurn( "pending", "effect_unknown", lease.attempt.reporting_materialization_id ) now = await _now(connection) if verified is not None and ( context.delivery is None or verified.resource.expires_at < max( context.delivery.resource_retained_until, now + timedelta(days=context.binding.resource_retention_days), ) ): error = ReportingWriterFailure("RESOURCE_UNAVAILABLE", "new_attempt", "applied") verified = None if verified is not None: # Same-connection finish validation after all lock acquisition # and DB-time reads; no I/O separates this from the copy below. validate_verified_destination(verified, lease.request, token=lease.token) outcome = ReportingMaterializationRecord( lease.scope, lease.attempt.reporting_revision_id, lease.attempt.reporting_materialization_id, context.binding.success_status if verified is not None else "failed", now, verified.resource if verified is not None else None, replace(verified.verification, verified_at=now) if verified is not None else None, public_failure(error) if error is not None else None, ) stored, inserted = await self._commit_record_on( connection, outcome, notify=False, dirty=False ) await self._materializer_dirty_on(connection, stored, context) if inserted and verified is not None and self._notifications_enabled: revision = next( r for r in context.revisions if r.reporting_revision_id == stored.reporting_revision_id ) event = materialization_event( stored, context.records, context.obligation, revision, context.configuration, now, ) if event is None: raise failure("BINDING_MISMATCH") await enqueue_materializer_event_on(connection, event) retry = error is not None and error.retry == "new_attempt" if reason == "ready": reason = ( "verified" if verified is not None else ("retry" if retry else "operator_required") ) ack = await ( await connection.execute( "UPDATE reporting_materializer_work SET state='acked',acknowledged_at=%s," " completion_token=lease_token,lease_token=NULL,lease_until=NULL," " retry_allowed=%s,reason=%s" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s" " AND lease_token=%s::uuid AND lease_until>clock_timestamp() RETURNING 1", ( now, retry, reason, lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, lease.token, ), ) ).fetchone() if ack is None: raise failure("LEASE_LOST") # Rolls back evidence, dirty, ACK and outbox together. due = ( now + timedelta(seconds=max(1, error.retry_after_seconds or 5)) if retry and error else None ) if reason == "target_changed": due = now elif reason not in {"retry", "verified"}: due = None if ( reason == "inactive" and context.configuration.activated_at is not None and context.configuration.activated_at > now ): due = context.configuration.activated_at if verified is not None: due = verified.resource.expires_at await self._park_on(connection, lease.scope, reason, due_at=due) await self._schedule_account_on(connection, lease.scope.principal.account_id) return ReportingMaterializerTurn( "verified" if verified is not None else "failed", reason, stored.reporting_materialization_id, ) async def _materializer_dirty_on( self, connection: Any, outcome: ReportingMaterializationRecord, context: MaterializerContext ) -> None: from adcp.reporting.ledger.notification_events import delivery_dirty scope, reason, evidence = delivery_dirty(outcome, context.obligation) await self._dirty_status(connection, scope, reason, after=evidence) from adcp.reporting.canonical_json import canonical_json_sha256_v1 from adcp.reporting.ledger.pg import _json from adcp.reporting.ledger.status_snapshot import settle_snapshot_on core = await settle_snapshot_on( self, connection, account_id=scope.account_id, as_of=outcome.completed_at ) row = await ( await connection.execute( "INSERT INTO reporting_materializer_status_heads" " (account_id,consumer_id,max_sequence)" " VALUES (%s,%s,1) ON CONFLICT (account_id,consumer_id) DO UPDATE" " SET max_sequence=reporting_materializer_status_heads.max_sequence+1" " RETURNING max_sequence", (scope.account_id, outcome.scope.consumer_id), ) ).fetchone() assert row is not None account = await ( await connection.execute( "UPDATE reporting_materializer_accounts SET captured_sequence=captured_sequence+1" " WHERE account_id=%s RETURNING captured_sequence", (scope.account_id,), ) ).fetchone() assert account is not None boundary = ReportingMaterializerBoundary( outcome.scope.principal, row[0], account[0], outcome.reporting_materialization_id, outcome.completed_at, private_snapshot(core, outcome.scope.principal), await self._records(connection, outcome.scope.principal), ).to_storage() await connection.execute( "INSERT INTO reporting_materializer_status_boundaries" " (account_id,consumer_id,sequence,account_sequence,reporting_materialization_id," " as_of,input,content_sha256)" " VALUES (%s,%s,%s,%s,%s,%s,%s::jsonb,%s)", ( scope.account_id, outcome.scope.consumer_id, row[0], account[0], outcome.reporting_materialization_id, outcome.completed_at, _json(boundary), canonical_json_sha256_v1(boundary), ), ) @materializer_errors async def read_materializer_boundaries( self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100 ) -> tuple[ReportingMaterializerBoundary, ...]: if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100: raise ReportingMaterializerUsageError( "materializer boundary reads require bounded positions" ) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, caller.account_id) rows = await ( await connection.execute( "SELECT input, content_sha256=reporting_payload_sha256(input)" " FROM reporting_materializer_status_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 row[1] for row in rows): raise failure("HISTORY_CORRUPT") return tuple(decode_materializer_boundary(row[0]) for row in rows) @materializer_errors async def import_pending_materialization( self, *, scope: ReportingDeliveryScope, reporting_materialization_id: str, original_external_id: str, keys: tuple[ReportingVerificationKey, ...], ) -> None: """Explicit operator recovery after verifying the original external identity. Never infer legacy effect history. An unknown/non-SDK external identity requires operator cleanup and a public terminal outcome before recovery. """ await self.materializer_ready() async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, scope.principal.account_id) context = await self._materializer_context_on(connection, scope) attempt = next( ( r for r in context.records if isinstance(r, ReportingMaterializationAttempt) and r.reporting_materialization_id == reporting_materialization_id and r.scope == scope ), None, ) if attempt is None or context.outcome(attempt) is not None: raise failure("BINDING_MISMATCH") if not context.resumable(attempt): raise failure("HISTORY_CORRUPT") key = key_for(context.binding, context.obligation, keys) request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) if original_external_id != request.external_id: raise failure("BINDING_MISMATCH") await connection.execute( "SELECT reporting_materializer_wake(%s)", (scope.principal.account_id,) ) await connection.execute( "INSERT INTO reporting_materializer_candidates" " (account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id)" " VALUES (%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING", _scope_args(scope), ) await connection.execute( "INSERT INTO reporting_materializer_work" " (account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id,reporting_revision_id,reporting_materialization_id," " generation,binding_sha256,verification_key_sha256,external_id,imported," " notifications_enabled)" " SELECT account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id,%s,%s,generation,%s,%s,%s,TRUE,%s" f" FROM reporting_materializer_candidates WHERE {_SCOPE_WHERE}" # nosec B608 " ON CONFLICT (account_id,consumer_id,reporting_materialization_id) DO NOTHING", ( attempt.reporting_revision_id, reporting_materialization_id, request.binding_fingerprint, verification_key_id(key), request.external_id, self._notifications_enabled, *_scope_args(scope), ), ) await connection.execute( "UPDATE reporting_materializer_work SET due_at=clock_timestamp()" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s" " AND state='pending' AND lease_token IS NULL", (scope.principal.account_id, scope.consumer_id, reporting_materialization_id), )Additive durable extension. No account enumeration or destination I/O.
Each turn samples at most sixteen due accounts without row locks, tries the Core advisory lock first, then claims within that account. Long work returns its connection before any source/destination call; even a size-one pool can renew its lease. Clock overrides affect old conformance APIs only: work time is always PostgreSQL time.
Ancestors
- PgReportingReconciliationStore
- PgReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
Subclasses
Methods
-
Expand source code
@materializer_errors async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None: async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, lease.scope.principal.account_id) await validate_materializer_schema( connection, notifications=self._notifications_enabled ) if await self._held_on(connection, lease) is None: raise failure("LEASE_LOST") _, reason = await self._target_on(connection, lease) if reason != "ready": raise failure( "HISTORY_CORRUPT" if reason == "history_corrupt" else "CURRENT_REVISION_CHANGED" ) async def claim_materialization(self,
*,
keys: tuple[ReportingVerificationKey, ...],
lease_seconds: int = 30) ‑> ReportingMaterializerLease | ReportingMaterializerTurn-
Expand source code
@materializer_errors async def claim_materialization( self, *, keys: tuple[ReportingVerificationKey, ...], lease_seconds: int = 30 ) -> ReportingMaterializerLease | ReportingMaterializerTurn: validate_lease_seconds(lease_seconds) await self.materializer_ready() async with self._connection() as connection: # Read-only sampling. Never lock global candidate rows before the # account advisory lock, including when a competing worker is busy. for _ in range(2): # Keep the cursor serialized, but order native timestamps: the # served_at::text output alias otherwise makes fairness locale-dependent. after = self._materializer_sample_after if after: query = ( "SELECT account_id,served_at::text FROM reporting_materializer_accounts" " WHERE due_at<=%s AND (served_at,account_id)>(%s::timestamptz,%s)" " ORDER BY reporting_materializer_accounts.served_at,account_id LIMIT 16" ) else: query = ( "SELECT account_id,served_at::text FROM reporting_materializer_accounts" " WHERE due_at<=%s" " ORDER BY reporting_materializer_accounts.served_at,account_id LIMIT 16" ) accounts = await ( await connection.execute( query, (await _now(connection), *(after or ())), ) ).fetchall() if accounts: break self._materializer_sample_after = None if after is None: break for account_id, served_at in accounts: self._materializer_sample_after = (served_at, account_id) async with self._connection() as connection, connection.transaction(): held = await ( await connection.execute( "SELECT pg_try_advisory_xact_lock(hashtext(%s))", (f"adcp.reporting:{account_id}",), ) ).fetchone() if held is None or not held[0]: continue await self._lock_account(connection, account_id) await validate_materializer_schema( connection, notifications=self._notifications_enabled ) await connection.execute( "UPDATE reporting_materializer_accounts SET served_at=clock_timestamp()" " WHERE account_id=%s", (account_id,), ) result = await self._claim_account_on(connection, account_id, keys, lease_seconds) await self._schedule_account_on(connection, account_id) # A committed account turn moves its durable served_at rank. Start # there again so continuously due peers cannot hide newly enrolled # accounts behind this pre-update cursor. Lock misses above retain # the continuation, allowing the next bounded page past a busy prefix. self._materializer_sample_after = None return result return ReportingMaterializerTurn("idle") async def finish_materialization(self,
lease: ReportingMaterializerLease,
*,
prepared: ReportingPreparedRevision | None = None,
verified: ReportingVerifiedDestination | None = None,
error: ReportingWriterFailure | None = None) ‑> ReportingMaterializerTurn-
Expand source code
@materializer_errors async def finish_materialization( self, lease: ReportingMaterializerLease, *, prepared: ReportingPreparedRevision | None = None, verified: ReportingVerifiedDestination | None = None, error: ReportingWriterFailure | None = None, ) -> ReportingMaterializerTurn: if verified is not None: validate_verified_destination(verified, lease.request, token=lease.token) if prepared is None or prepared.request != lease.request or error is not None: raise failure("BINDING_MISMATCH") elif error is None: raise failure("BINDING_MISMATCH") async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, lease.scope.principal.account_id) await validate_materializer_schema( connection, notifications=self._notifications_enabled ) held = await self._held_on(connection, lease) if held is None: replay = await ( await connection.execute( "SELECT reason FROM reporting_materializer_work WHERE account_id=%s" " AND consumer_id=%s AND reporting_materialization_id=%s AND state='acked'" " AND completion_token=%s::uuid AND external_id=%s", ( lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, lease.token, lease.request.external_id, ), ) ).fetchone() if replay is not None: return ReportingMaterializerTurn( "verified" if replay[0] == "verified" else "failed", replay[0], lease.attempt.reporting_materialization_id, ) return ReportingMaterializerTurn( "pending", "effect_unknown", lease.attempt.reporting_materialization_id ) context, reason = await self._target_on(connection, lease) if reason != "ready": error = ReportingWriterFailure("CURRENT_REVISION_CHANGED", "new_attempt", "applied") verified = None elif verified is not None and prepared is not None: validate_materialization_target( prepared, binding=context.binding, revisions=context.revisions ) if error is not None and ( error.effect == "unknown" or error.retry == "same_identity" or error.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"} ): delay = max(1, min(error.retry_after_seconds or 5, 300)) await connection.execute( "UPDATE reporting_materializer_work SET lease_token=NULL,lease_until=NULL," " due_at=clock_timestamp()+make_interval(secs=>%s),reason='effect_unknown'" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s", ( delay, lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, ), ) await self._schedule_account_on(connection, lease.scope.principal.account_id) return ReportingMaterializerTurn( "pending", "effect_unknown", lease.attempt.reporting_materialization_id ) now = await _now(connection) if verified is not None and ( context.delivery is None or verified.resource.expires_at < max( context.delivery.resource_retained_until, now + timedelta(days=context.binding.resource_retention_days), ) ): error = ReportingWriterFailure("RESOURCE_UNAVAILABLE", "new_attempt", "applied") verified = None if verified is not None: # Same-connection finish validation after all lock acquisition # and DB-time reads; no I/O separates this from the copy below. validate_verified_destination(verified, lease.request, token=lease.token) outcome = ReportingMaterializationRecord( lease.scope, lease.attempt.reporting_revision_id, lease.attempt.reporting_materialization_id, context.binding.success_status if verified is not None else "failed", now, verified.resource if verified is not None else None, replace(verified.verification, verified_at=now) if verified is not None else None, public_failure(error) if error is not None else None, ) stored, inserted = await self._commit_record_on( connection, outcome, notify=False, dirty=False ) await self._materializer_dirty_on(connection, stored, context) if inserted and verified is not None and self._notifications_enabled: revision = next( r for r in context.revisions if r.reporting_revision_id == stored.reporting_revision_id ) event = materialization_event( stored, context.records, context.obligation, revision, context.configuration, now, ) if event is None: raise failure("BINDING_MISMATCH") await enqueue_materializer_event_on(connection, event) retry = error is not None and error.retry == "new_attempt" if reason == "ready": reason = ( "verified" if verified is not None else ("retry" if retry else "operator_required") ) ack = await ( await connection.execute( "UPDATE reporting_materializer_work SET state='acked',acknowledged_at=%s," " completion_token=lease_token,lease_token=NULL,lease_until=NULL," " retry_allowed=%s,reason=%s" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s" " AND lease_token=%s::uuid AND lease_until>clock_timestamp() RETURNING 1", ( now, retry, reason, lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, lease.token, ), ) ).fetchone() if ack is None: raise failure("LEASE_LOST") # Rolls back evidence, dirty, ACK and outbox together. due = ( now + timedelta(seconds=max(1, error.retry_after_seconds or 5)) if retry and error else None ) if reason == "target_changed": due = now elif reason not in {"retry", "verified"}: due = None if ( reason == "inactive" and context.configuration.activated_at is not None and context.configuration.activated_at > now ): due = context.configuration.activated_at if verified is not None: due = verified.resource.expires_at await self._park_on(connection, lease.scope, reason, due_at=due) await self._schedule_account_on(connection, lease.scope.principal.account_id) return ReportingMaterializerTurn( "verified" if verified is not None else "failed", reason, stored.reporting_materialization_id, ) async def import_pending_materialization(self,
*,
scope: ReportingDeliveryScope,
reporting_materialization_id: str,
original_external_id: str,
keys: tuple[ReportingVerificationKey, ...]) ‑> None-
Expand source code
@materializer_errors async def import_pending_materialization( self, *, scope: ReportingDeliveryScope, reporting_materialization_id: str, original_external_id: str, keys: tuple[ReportingVerificationKey, ...], ) -> None: """Explicit operator recovery after verifying the original external identity. Never infer legacy effect history. An unknown/non-SDK external identity requires operator cleanup and a public terminal outcome before recovery. """ await self.materializer_ready() async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, scope.principal.account_id) context = await self._materializer_context_on(connection, scope) attempt = next( ( r for r in context.records if isinstance(r, ReportingMaterializationAttempt) and r.reporting_materialization_id == reporting_materialization_id and r.scope == scope ), None, ) if attempt is None or context.outcome(attempt) is not None: raise failure("BINDING_MISMATCH") if not context.resumable(attempt): raise failure("HISTORY_CORRUPT") key = key_for(context.binding, context.obligation, keys) request = ReportingDestinationRequest.from_binding(context.binding, attempt, key) if original_external_id != request.external_id: raise failure("BINDING_MISMATCH") await connection.execute( "SELECT reporting_materializer_wake(%s)", (scope.principal.account_id,) ) await connection.execute( "INSERT INTO reporting_materializer_candidates" " (account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id)" " VALUES (%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING", _scope_args(scope), ) await connection.execute( "INSERT INTO reporting_materializer_work" " (account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id,reporting_revision_id,reporting_materialization_id," " generation,binding_sha256,verification_key_sha256,external_id,imported," " notifications_enabled)" " SELECT account_id,consumer_id,delivery_config_id,delivery_config_version," " reporting_obligation_id,%s,%s,generation,%s,%s,%s,TRUE,%s" f" FROM reporting_materializer_candidates WHERE {_SCOPE_WHERE}" # nosec B608 " ON CONFLICT (account_id,consumer_id,reporting_materialization_id) DO NOTHING", ( attempt.reporting_revision_id, reporting_materialization_id, request.binding_fingerprint, verification_key_id(key), request.external_id, self._notifications_enabled, *_scope_args(scope), ), ) await connection.execute( "UPDATE reporting_materializer_work SET due_at=clock_timestamp()" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s" " AND state='pending' AND lease_token IS NULL", (scope.principal.account_id, scope.consumer_id, reporting_materialization_id), )Explicit operator recovery after verifying the original external identity.
Never infer legacy effect history. An unknown/non-SDK external identity requires operator cleanup and a public terminal outcome before recovery.
async def materializer_ready(self) ‑> bool-
Expand source code
@materializer_errors async def materializer_ready(self) -> bool: async with self._connection() as connection: await validate_materializer_schema( connection, notifications=self._notifications_enabled ) return True async def read_materializer_boundaries(self,
*,
caller: ReportingDeliveryPrincipal,
after: int = 0,
limit: int = 100) ‑> tuple[ReportingMaterializerBoundary, ...]-
Expand source code
@materializer_errors async def read_materializer_boundaries( self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100 ) -> tuple[ReportingMaterializerBoundary, ...]: if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100: raise ReportingMaterializerUsageError( "materializer boundary reads require bounded positions" ) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, caller.account_id) rows = await ( await connection.execute( "SELECT input, content_sha256=reporting_payload_sha256(input)" " FROM reporting_materializer_status_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 row[1] for row in rows): raise failure("HISTORY_CORRUPT") return tuple(decode_materializer_boundary(row[0]) for row in rows) async def renew_materialization(self,
lease: ReportingMaterializerLease,
*,
lease_seconds: int) ‑> bool-
Expand source code
@materializer_errors async def renew_materialization( self, lease: ReportingMaterializerLease, *, lease_seconds: int ) -> bool: validate_lease_seconds(lease_seconds) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, lease.scope.principal.account_id) if await self._held_on(connection, lease) is None: return False await connection.execute( "UPDATE reporting_materializer_work SET lease_until=clock_timestamp()" " +make_interval(secs=>%s),due_at=clock_timestamp()+make_interval(secs=>%s)" " WHERE account_id=%s AND consumer_id=%s AND reporting_materialization_id=%s", ( lease_seconds, lease_seconds, lease.scope.principal.account_id, lease.scope.consumer_id, lease.attempt.reporting_materialization_id, ), ) await self._schedule_account_on(connection, lease.scope.principal.account_id) return True
Inherited members
class ReferenceReportingDestinationWriter (capabilities: tuple[ReportingWriterCapability, ...])-
Expand source code
@final class ReferenceReportingDestinationWriter: """Conditional, deterministic external identities in shared process memory. Losing this object loses its artifacts. No durability or Managed capability follows from this type or its descriptors. The production flag is a constant property, with no constructor/config override. """ def __init_subclass__(cls, **kwargs: Any) -> None: raise TypeError("the non-production reference writer cannot be promoted by subclassing") def __init__(self, capabilities: tuple[ReportingWriterCapability, ...]) -> None: if type(capabilities) is not tuple or any( type(c) is not ReportingWriterCapability for c in capabilities ): raise ValueError("reference capabilities require an immutable typed tuple") self._capabilities = capabilities self._artifacts: dict[str, _Artifact] = {} self.write_effects = 0 self.open_count = 0 self.close_count = 0 @property def capabilities(self) -> tuple[ReportingWriterCapability, ...]: return self._capabilities @property def production_eligible(self) -> Literal[False]: return False def __repr__(self) -> str: return "<ReferenceReportingDestinationWriter non-production>"Conditional, deterministic external identities in shared process memory.
Losing this object loses its artifacts. No durability or Managed capability follows from this type or its descriptors. The production flag is a constant property, with no constructor/config override.
Instance variables
prop capabilities : tuple[ReportingWriterCapability, ...]-
Expand source code
@property def capabilities(self) -> tuple[ReportingWriterCapability, ...]: return self._capabilities prop production_eligible : Literal[False]-
Expand source code
@property def production_eligible(self) -> Literal[False]: return False
class ReferenceReportingResolver (writer: ReferenceReportingDestinationWriter,
registry: ReportingRevisionVerifierRegistry,
bindings: tuple[ReportingDestinationBinding, ...])-
Expand source code
@dataclass(frozen=True) class ReferenceReportingResolver: writer: ReferenceReportingDestinationWriter registry: ReportingRevisionVerifierRegistry bindings: tuple[ReportingDestinationBinding, ...] = field(repr=False) _revoked: set[ReportingDeliveryPrincipal] = field( default_factory=set, init=False, repr=False, compare=False ) _rotations: list[int] = field( default_factory=lambda: [0], init=False, repr=False, compare=False ) def __post_init__(self) -> None: if type(self.bindings) is not tuple or any( type(b) is not ReportingDestinationBinding for b in self.bindings ): raise ValueError("reference resolver requires frozen bindings") def revoke(self, principal: ReportingDeliveryPrincipal) -> None: self._revoked.add(principal) def rotate(self) -> None: self._rotations[0] += 1 def resolve( self, request: ReportingDestinationRequest, *, phase: ReportingIOPhase, context: ReportingIOContext, ) -> ReportingDestinationSession: return _Session(self, request, phase, context)ReferenceReportingResolver(writer: 'ReferenceReportingDestinationWriter', registry: 'ReportingRevisionVerifierRegistry', bindings: 'tuple[ReportingDestinationBinding, …]')
Instance variables
var bindings : tuple[ReportingDestinationBinding, ...]var registry : ReportingRevisionVerifierRegistryvar writer : ReferenceReportingDestinationWriter
Methods
def resolve(self,
request: ReportingDestinationRequest,
*,
phase: ReportingIOPhase,
context: ReportingIOContext) ‑> ReportingDestinationSession-
Expand source code
def resolve( self, request: ReportingDestinationRequest, *, phase: ReportingIOPhase, context: ReportingIOContext, ) -> ReportingDestinationSession: return _Session(self, request, phase, context) def revoke(self,
principal: ReportingDeliveryPrincipal) ‑> None-
Expand source code
def revoke(self, principal: ReportingDeliveryPrincipal) -> None: self._revoked.add(principal) def rotate(self) ‑> None-
Expand source code
def rotate(self) -> None: self._rotations[0] += 1
class ReportingCanonicalization (canonicalization_id: str,
canonicalization_uri: str,
canonicalization_sha256: str)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingCanonicalization(_ClosedValue): canonicalization_id: str canonicalization_uri: str canonicalization_sha256: str def __post_init__(self) -> None: _freeze_fields(self) ReportingCanonicalDigest( "0" * 64, self.canonicalization_id, self.canonicalization_uri, self.canonicalization_sha256, ) object.__setattr__(self, "canonicalization_sha256", self.canonicalization_sha256.lower())ReportingCanonicalization(canonicalization_id: 'str', canonicalization_uri: 'str', canonicalization_sha256: 'str')
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var canonicalization_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingCanonicalization(_ClosedValue): canonicalization_id: str canonicalization_uri: str canonicalization_sha256: str def __post_init__(self) -> None: _freeze_fields(self) ReportingCanonicalDigest( "0" * 64, self.canonicalization_id, self.canonicalization_uri, self.canonicalization_sha256, ) object.__setattr__(self, "canonicalization_sha256", self.canonicalization_sha256.lower()) var canonicalization_sha256 : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingCanonicalization(_ClosedValue): canonicalization_id: str canonicalization_uri: str canonicalization_sha256: str def __post_init__(self) -> None: _freeze_fields(self) ReportingCanonicalDigest( "0" * 64, self.canonicalization_id, self.canonicalization_uri, self.canonicalization_sha256, ) object.__setattr__(self, "canonicalization_sha256", self.canonicalization_sha256.lower()) var canonicalization_uri : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingCanonicalization(_ClosedValue): canonicalization_id: str canonicalization_uri: str canonicalization_sha256: str def __post_init__(self) -> None: _freeze_fields(self) ReportingCanonicalDigest( "0" * 64, self.canonicalization_id, self.canonicalization_uri, self.canonicalization_sha256, ) object.__setattr__(self, "canonicalization_sha256", self.canonicalization_sha256.lower())
class ReportingDeliveryPrincipal (account_id: str, consumer_id: str)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDeliveryPrincipal(_ClosedValue): """Trusted account and consumer identity, resolved from authenticated transport.""" account_id: str consumer_id: str def __post_init__(self) -> None: _freeze_fields(self) principal_reference(self.account_id) from adcp.reporting.evidence import consumer_reference consumer_reference(self.consumer_id)Trusted account and consumer identity, resolved from authenticated transport.
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var account_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDeliveryPrincipal(_ClosedValue): """Trusted account and consumer identity, resolved from authenticated transport.""" account_id: str consumer_id: str def __post_init__(self) -> None: _freeze_fields(self) principal_reference(self.account_id) from adcp.reporting.evidence import consumer_reference consumer_reference(self.consumer_id) var consumer_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDeliveryPrincipal(_ClosedValue): """Trusted account and consumer identity, resolved from authenticated transport.""" account_id: str consumer_id: str def __post_init__(self) -> None: _freeze_fields(self) principal_reference(self.account_id) from adcp.reporting.evidence import consumer_reference consumer_reference(self.consumer_id)
class ReportingDestinationBinding (generation_key: ReportingConfigurationGenerationKey,
consumer_id: str,
destination_ref: str,
trusted_binding_ref: str,
method: DeliveryMethod,
transport: str,
verification_profile: VerificationProfile,
reconciliation_mode: "Literal['delivery_only', 'consumer_receipt']",
feed_purpose: "Literal['pacing', 'analytics', 'billing']",
resource_retention_days: int,
created_at: datetime,
format: ReportingFormat | None = None,
reader_compatibility: tuple[str, ...] = (),
success_status: "Literal['available', 'delivered']" = 'available',
*,
kind: "Literal['destination_binding']" = 'destination_binding')-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id)Immutable public portion of one trusted principal/configuration binding.
trusted_binding_refidentifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability.Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var consumer_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var created_at : datetime.datetime-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var destination_ref : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var feed_purpose : Literal['pacing', 'analytics', 'billing']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var format : Literal['jsonl', 'csv', 'parquet', 'avro', 'orc'] | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var generation_key : ReportingConfigurationGenerationKey-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var kind : Literal['destination_binding']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var method : Literal['file_transfer', 'dataset_share', 'warehouse_materialization']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) prop principal : ReportingDeliveryPrincipal-
Expand source code
@property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var reader_compatibility : tuple[str, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var reconciliation_mode : Literal['delivery_only', 'consumer_receipt']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var resource_retention_days : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var success_status : Literal['available', 'delivered']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var transport : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var trusted_binding_ref : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id) var verification_profile : Literal['canonical_digest', 'manifest_checksums', 'native_commit']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationBinding(_ClosedValue): """Immutable public portion of one trusted principal/configuration binding. ``trusted_binding_ref`` identifies an immutable trusted configuration, not a mutable 'latest' alias. Credentials are resolved later behind that reference. Storing this record neither configures a writer nor advertises a capability. """ generation_key: ReportingConfigurationGenerationKey consumer_id: str destination_ref: str trusted_binding_ref: str = field(repr=False) method: DeliveryMethod transport: str verification_profile: VerificationProfile reconciliation_mode: Literal["delivery_only", "consumer_receipt"] feed_purpose: Literal["pacing", "analytics", "billing"] resource_retention_days: int created_at: datetime format: ReportingFormat | None = None reader_compatibility: tuple[str, ...] = () success_status: Literal["available", "delivered"] = "available" kind: Literal["destination_binding"] = field(default="destination_binding", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) ReportingDeliveryScope(self.generation_key, self.consumer_id, "binding") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("reporting transport requires a public protocol label") reporting_identifier(self.transport, maximum=64) _positive(self.resource_retention_days) object.__setattr__(self, "created_at", aware_utc(self.created_at)) object.__setattr__( self, "reader_compatibility", tuple(reader_feature_reference(value) for value in self.reader_compatibility), ) if len(set(self.reader_compatibility)) != len(self.reader_compatibility): raise ValueError("reader compatibility requirements must be unique") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a declared format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest checksum verification requires file transfer") if self.method == "warehouse_materialization" and self.success_status != "delivered": raise ValueError("warehouse materialization requires destination delivery") if self.method == "dataset_share" and self.success_status != "available": raise ValueError("dataset share requires representative-consumer availability") if self.feed_purpose == "billing" and self.verification_profile != "canonical_digest": raise ValueError("billing receipts require canonical digest verification") @property def principal(self) -> ReportingDeliveryPrincipal: return ReportingDeliveryPrincipal(self.generation_key.account_id, self.consumer_id)
class ReportingDestinationIO (registry: ReportingRevisionVerifierRegistry,
resolver: ReportingDestinationResolver)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationIO: """Explicit one-shot write/readback operations; no coordinator or lease loop.""" registry: ReportingRevisionVerifierRegistry resolver: ReportingDestinationResolver = field(repr=False) @asynccontextmanager async def _session( self, prepared: ReportingPreparedRevision, phase: ReportingIOPhase, context: ReportingIOContext, ) -> AsyncIterator[ReportingDestinationSession]: self.registry.require(prepared.request.verification_key) problem = False canceled = False try: session = self.resolver.resolve(prepared.request, phase=phase, context=context) except asyncio.CancelledError: canceled = True except Exception: problem = True if canceled: raise asyncio.CancelledError if problem: raise ReportingWriterError( ReportingWriterFailure("RESOURCE_UNAVAILABLE", "same_identity") ) if not isinstance(session, ReportingDestinationSession): raise failure("BINDING_MISMATCH") try: _session_binding(session, prepared.request, phase, context) except ReportingWriterError: problem = True if problem: await session.aclose() raise failure("BINDING_MISMATCH") failure_record: ReportingWriterFailure | None = None canceled = False try: async with session: _session_binding(session, prepared.request, phase, context) yield session except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: failure_record = exc.failure except Exception: failure_record = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", "unknown" ) if canceled: raise asyncio.CancelledError if failure_record is not None: raise ReportingWriterError(failure_record) async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem) def _validate_prepared(self, prepared: ReportingPreparedRevision) -> ReportingRevisionVerifier: verifier = self.registry.require(prepared.request.verification_key) request, obligation, binding = prepared.request, prepared.obligation, prepared.binding cap = request.verification_key.capability if ( binding_fingerprint(binding) != request.binding_fingerprint or request.principal != binding.principal or request.generation != binding.generation_key or request.destination_ref != binding.destination_ref or request.trusted_binding_ref != binding.trusted_binding_ref or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != prepared.revision.reporting_revision_id or prepared.revision.account_id != obligation.account_id or prepared.revision.reporting_obligation_id != obligation.reporting_obligation_id or not prepared.revision.readable or prepared.delivery.scope.principal != request.principal or prepared.delivery.scope.generation_key != request.generation or prepared.delivery.scope.reporting_obligation_id != request.reporting_obligation_id or prepared.delivery.currency != obligation.currency or not _same_definition(verifier.key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (verifier.key.report_definition_id, verifier.key.reporting_profile) or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") rows, totals = verifier.canonicalize( [parse_reporting_json(r, verifier.limits) for r in prepared.rows] ) if rows != prepared.rows: raise failure("SOURCE_INVALID") _verify_content(prepared.revision, verifier.key, rows, totals) return verifier async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem)Explicit one-shot write/readback operations; no coordinator or lease loop.
Instance variables
var registry : ReportingRevisionVerifierRegistry-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationIO: """Explicit one-shot write/readback operations; no coordinator or lease loop.""" registry: ReportingRevisionVerifierRegistry resolver: ReportingDestinationResolver = field(repr=False) @asynccontextmanager async def _session( self, prepared: ReportingPreparedRevision, phase: ReportingIOPhase, context: ReportingIOContext, ) -> AsyncIterator[ReportingDestinationSession]: self.registry.require(prepared.request.verification_key) problem = False canceled = False try: session = self.resolver.resolve(prepared.request, phase=phase, context=context) except asyncio.CancelledError: canceled = True except Exception: problem = True if canceled: raise asyncio.CancelledError if problem: raise ReportingWriterError( ReportingWriterFailure("RESOURCE_UNAVAILABLE", "same_identity") ) if not isinstance(session, ReportingDestinationSession): raise failure("BINDING_MISMATCH") try: _session_binding(session, prepared.request, phase, context) except ReportingWriterError: problem = True if problem: await session.aclose() raise failure("BINDING_MISMATCH") failure_record: ReportingWriterFailure | None = None canceled = False try: async with session: _session_binding(session, prepared.request, phase, context) yield session except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: failure_record = exc.failure except Exception: failure_record = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", "unknown" ) if canceled: raise asyncio.CancelledError if failure_record is not None: raise ReportingWriterError(failure_record) async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem) def _validate_prepared(self, prepared: ReportingPreparedRevision) -> ReportingRevisionVerifier: verifier = self.registry.require(prepared.request.verification_key) request, obligation, binding = prepared.request, prepared.obligation, prepared.binding cap = request.verification_key.capability if ( binding_fingerprint(binding) != request.binding_fingerprint or request.principal != binding.principal or request.generation != binding.generation_key or request.destination_ref != binding.destination_ref or request.trusted_binding_ref != binding.trusted_binding_ref or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != prepared.revision.reporting_revision_id or prepared.revision.account_id != obligation.account_id or prepared.revision.reporting_obligation_id != obligation.reporting_obligation_id or not prepared.revision.readable or prepared.delivery.scope.principal != request.principal or prepared.delivery.scope.generation_key != request.generation or prepared.delivery.scope.reporting_obligation_id != request.reporting_obligation_id or prepared.delivery.currency != obligation.currency or not _same_definition(verifier.key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (verifier.key.report_definition_id, verifier.key.reporting_profile) or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") rows, totals = verifier.canonicalize( [parse_reporting_json(r, verifier.limits) for r in prepared.rows] ) if rows != prepared.rows: raise failure("SOURCE_INVALID") _verify_content(prepared.revision, verifier.key, rows, totals) return verifier async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem) var resolver : ReportingDestinationResolver-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationIO: """Explicit one-shot write/readback operations; no coordinator or lease loop.""" registry: ReportingRevisionVerifierRegistry resolver: ReportingDestinationResolver = field(repr=False) @asynccontextmanager async def _session( self, prepared: ReportingPreparedRevision, phase: ReportingIOPhase, context: ReportingIOContext, ) -> AsyncIterator[ReportingDestinationSession]: self.registry.require(prepared.request.verification_key) problem = False canceled = False try: session = self.resolver.resolve(prepared.request, phase=phase, context=context) except asyncio.CancelledError: canceled = True except Exception: problem = True if canceled: raise asyncio.CancelledError if problem: raise ReportingWriterError( ReportingWriterFailure("RESOURCE_UNAVAILABLE", "same_identity") ) if not isinstance(session, ReportingDestinationSession): raise failure("BINDING_MISMATCH") try: _session_binding(session, prepared.request, phase, context) except ReportingWriterError: problem = True if problem: await session.aclose() raise failure("BINDING_MISMATCH") failure_record: ReportingWriterFailure | None = None canceled = False try: async with session: _session_binding(session, prepared.request, phase, context) yield session except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: failure_record = exc.failure except Exception: failure_record = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", "unknown" ) if canceled: raise asyncio.CancelledError if failure_record is not None: raise ReportingWriterError(failure_record) async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem) def _validate_prepared(self, prepared: ReportingPreparedRevision) -> ReportingRevisionVerifier: verifier = self.registry.require(prepared.request.verification_key) request, obligation, binding = prepared.request, prepared.obligation, prepared.binding cap = request.verification_key.capability if ( binding_fingerprint(binding) != request.binding_fingerprint or request.principal != binding.principal or request.generation != binding.generation_key or request.destination_ref != binding.destination_ref or request.trusted_binding_ref != binding.trusted_binding_ref or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != prepared.revision.reporting_revision_id or prepared.revision.account_id != obligation.account_id or prepared.revision.reporting_obligation_id != obligation.reporting_obligation_id or not prepared.revision.readable or prepared.delivery.scope.principal != request.principal or prepared.delivery.scope.generation_key != request.generation or prepared.delivery.scope.reporting_obligation_id != request.reporting_obligation_id or prepared.delivery.currency != obligation.currency or not _same_definition(verifier.key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (verifier.key.report_definition_id, verifier.key.reporting_profile) or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") rows, totals = verifier.canonicalize( [parse_reporting_json(r, verifier.limits) for r in prepared.rows] ) if rows != prepared.rows: raise failure("SOURCE_INVALID") _verify_content(prepared.revision, verifier.key, rows, totals) return verifier async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem)
Methods
async def verify(self,
prepared: ReportingPreparedRevision,
locator: ReportingDestinationLocator,
*,
context: ReportingIOContext) ‑> ReportingVerifiedDestination-
Expand source code
async def verify( self, prepared: ReportingPreparedRevision, locator: ReportingDestinationLocator, *, context: ReportingIOContext, ) -> ReportingVerifiedDestination: verifier = self._validate_prepared(prepared) _locator_binding(locator, prepared.request) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "readback", context) as session: _session_binding(session, prepared.request, "readback", context) return await _verify_destination(verifier, prepared, locator, session, context) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None if problem.code == "SOURCE_INVALID": problem = ReportingWriterFailure("DESTINATION_CORRUPT", "new_attempt", "applied") raise ReportingWriterError(problem) async def write(self,
prepared: ReportingPreparedRevision,
*,
context: ReportingIOContext) ‑> ReportingDestinationLocator-
Expand source code
async def write( self, prepared: ReportingPreparedRevision, *, context: ReportingIOContext ) -> ReportingDestinationLocator: self._validate_prepared(prepared) problem: ReportingWriterFailure | None = None canceled = False try: async with self._session(prepared, "write", context) as session: _session_binding(session, prepared.request, "write", context) locator = await context.run(lambda: session.write(prepared), effect="unknown") invalid_locator = False try: _locator_binding(locator, prepared.request) except ReportingWriterError: invalid_locator = True if invalid_locator: raise ReportingWriterError( ReportingWriterFailure("BINDING_MISMATCH", "same_identity", "unknown") ) return locator except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled: raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem)
class ReportingDestinationLocator (external_id: str, binding_fingerprint: str, resource: ReportingResourceRecord)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationLocator(_ClosedValue): """Writer claims identify what to read. They are never verification proof.""" external_id: str binding_fingerprint: str resource: ReportingResourceRecord def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.external_id) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower())Writer claims identify what to read. They are never verification proof.
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var binding_fingerprint : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationLocator(_ClosedValue): """Writer claims identify what to read. They are never verification proof.""" external_id: str binding_fingerprint: str resource: ReportingResourceRecord def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.external_id) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) var external_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationLocator(_ClosedValue): """Writer claims identify what to read. They are never verification proof.""" external_id: str binding_fingerprint: str resource: ReportingResourceRecord def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.external_id) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) var resource : ReportingResourceRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationLocator(_ClosedValue): """Writer claims identify what to read. They are never verification proof.""" external_id: str binding_fingerprint: str resource: ReportingResourceRecord def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.external_id) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower())
class ReportingDestinationPage (reporting_revision_id: str,
rows: tuple[bytes, ...],
total_count: int,
has_more: bool,
cursor: str | None,
format: ReportingFormat | None,
verification_path: VerificationPath,
native_version_ref: str | None = None)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref)ReportingDestinationPage(reporting_revision_id: 'str', rows: 'tuple[bytes, …]', total_count: 'int', has_more: 'bool', cursor: 'str | None', format: 'ReportingFormat | None', verification_path: 'VerificationPath', native_version_ref: 'str | None' = None)
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var cursor : str | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref) var format : Literal['jsonl', 'csv', 'parquet', 'avro', 'orc'] | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref) var has_more : bool-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref) var native_version_ref : str | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref) var reporting_revision_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref) var rows : tuple[bytes, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref) var total_count : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref) var verification_path : Literal['producer', 'representative_consumer', 'destination']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationPage(_ClosedValue): reporting_revision_id: str rows: tuple[bytes, ...] = field(repr=False) total_count: int has_more: bool cursor: str | None format: ReportingFormat | None verification_path: VerificationPath native_version_ref: str | None = None def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id) if self.total_count < 0 or self.has_more != (self.cursor is not None): raise ValueError("destination page requires a paired cursor and valid total") if self.cursor is not None: reporting_identifier(self.cursor, maximum=2048) if self.native_version_ref is not None: native_version_reference(self.native_version_ref)
class ReportingDestinationRequest (principal: ReportingDeliveryPrincipal,
generation: ReportingConfigurationGenerationKey,
destination_ref: str,
trusted_binding_ref: str,
binding_fingerprint: str,
verification_key: ReportingVerificationKey,
reporting_obligation_id: str,
reporting_revision_id: str,
reporting_materialization_id: str,
attempt: int)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest()Trusted resolver input. Aliases must already resolve to the canonical consumer.
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Static methods
def from_binding(binding: ReportingDestinationBinding,
attempt: ReportingMaterializationAttempt,
key: ReportingVerificationKey) ‑> ReportingDestinationRequest
Instance variables
var attempt : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var binding_fingerprint : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var destination_ref : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() prop external_id : str-
Expand source code
@property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest()Stable across pending retries; isolated across tenants and revision attempts.
var generation : ReportingConfigurationGenerationKey-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var principal : ReportingDeliveryPrincipal-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var reporting_materialization_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var reporting_obligation_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var reporting_revision_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var trusted_binding_ref : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest() var verification_key : ReportingVerificationKey-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingDestinationRequest(_ClosedValue): """Trusted resolver input. Aliases must already resolve to the canonical consumer.""" principal: ReportingDeliveryPrincipal generation: ReportingConfigurationGenerationKey destination_ref: str trusted_binding_ref: str = field(repr=False) binding_fingerprint: str verification_key: ReportingVerificationKey reporting_obligation_id: str reporting_revision_id: str reporting_materialization_id: str attempt: int def __post_init__(self) -> None: _freeze_fields(self) if self.principal.account_id != self.generation.account_id or self.attempt < 1: raise ValueError("destination request requires an exact account and attempt") reporting_identifier(self.generation.delivery_config_id) if ( type(self.generation.delivery_config_version) is not int or self.generation.delivery_config_version < 1 ): raise ValueError("destination request requires an exact generation") destination_reference(self.destination_ref) destination_reference(self.trusted_binding_ref) sha256_value(self.binding_fingerprint) object.__setattr__(self, "binding_fingerprint", self.binding_fingerprint.lower()) for value in ( self.reporting_obligation_id, self.reporting_revision_id, self.reporting_materialization_id, ): reporting_identifier(value) @classmethod def from_binding( cls, binding: ReportingDestinationBinding, attempt: ReportingMaterializationAttempt, key: ReportingVerificationKey, ) -> ReportingDestinationRequest: cap = key.capability if ( binding.principal != attempt.scope.principal or binding.generation_key != attempt.scope.generation_key or (binding.method, binding.transport, binding.format, binding.verification_profile) != (cap.method, cap.transport, cap.format, cap.verification_profile) or (binding.success_status == "delivered" and cap.verification_path != "destination") ): raise failure("BINDING_MISMATCH") return cls( binding.principal, binding.generation_key, binding.destination_ref, binding.trusted_binding_ref, binding_fingerprint(binding), key, attempt.scope.reporting_obligation_id, attempt.reporting_revision_id, attempt.reporting_materialization_id, attempt.attempt, ) @property def external_id(self) -> str: """Stable across pending retries; isolated across tenants and revision attempts.""" return "rwm_" + hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest()
class ReportingDestinationResolver (*args, **kwargs)-
Expand source code
class ReportingDestinationResolver(Protocol): """Return an unopened session. _open reauthorizes each phase independently. Resolution must match every request coordinate, including consumer URL, immutable trusted binding, definition/canonicalization and capability. A resolver must never allocate resources before constructing the session. """ def resolve( self, request: ReportingDestinationRequest, *, phase: ReportingIOPhase, context: ReportingIOContext, ) -> ReportingDestinationSession: ...Return an unopened session. _open reauthorizes each phase independently.
Resolution must match every request coordinate, including consumer URL, immutable trusted binding, definition/canonicalization and capability. A resolver must never allocate resources before constructing the session.
Ancestors
- typing.Protocol
- typing.Generic
Subclasses
Methods
def resolve(self,
request: ReportingDestinationRequest,
*,
phase: ReportingIOPhase,
context: ReportingIOContext) ‑> ReportingDestinationSession-
Expand source code
def resolve( self, request: ReportingDestinationRequest, *, phase: ReportingIOPhase, context: ReportingIOContext, ) -> ReportingDestinationSession: ...
class ReportingDestinationSession (request: ReportingDestinationRequest,
phase: ReportingIOPhase,
context: ReportingIOContext)-
Expand source code
class ReportingDestinationSession(ABC): """Single-use SDK-owned lifecycle around an adopter's private provider session. Override only _open/_close and I/O methods. _close must tolerate partial _open, finish promptly, and remove any adopter-owned temporary spool. The SDK invokes it exactly once, shields cancellation, and joins all its tasks. No credentials may be placed on the public request, descriptor or results. """ def __init_subclass__(cls, **kwargs: Any) -> None: super().__init_subclass__(**kwargs) owned = { "__repr__", "__str__", "__reduce__", "__reduce_ex__", "__getstate__", "__aenter__", "__aexit__", "aclose", "request", "phase", "context", } if owned.intersection(cls.__dict__): raise TypeError("destination session lifecycle and redaction belong to the SDK") def __init__( self, request: ReportingDestinationRequest, phase: ReportingIOPhase, context: ReportingIOContext, ) -> None: if ( type(request) is not ReportingDestinationRequest or phase not in ("write", "readback") or type(context) is not ReportingIOContext ): raise failure("BINDING_MISMATCH") self._request, self._phase, self._context = request, phase, context self._entered = False self._closed = False def __repr__(self) -> str: return "<ReportingDestinationSession redacted>" __str__ = __repr__ @property def request(self) -> ReportingDestinationRequest: return self._request @property def phase(self) -> ReportingIOPhase: return self._phase @property def context(self) -> ReportingIOContext: return self._context def __reduce__(self) -> tuple[Any, ...]: raise TypeError("destination sessions cannot be persisted") @abstractmethod async def _open(self) -> None: ... @abstractmethod async def _close(self) -> None: ... async def __aenter__(self) -> ReportingDestinationSession: if self._entered or self._closed: raise failure("BINDING_MISMATCH") self._entered = True problem: ReportingWriterFailure | None = None canceled = False try: await self.context.run(self._open) except asyncio.CancelledError: canceled = True except ReportingWriterError as exc: problem = exc.failure if canceled or problem is not None: try: await self.aclose() except asyncio.CancelledError: canceled = True except ReportingWriterError: pass if canceled or self.context.cancel.is_set(): raise asyncio.CancelledError assert problem is not None raise ReportingWriterError(problem) return self async def __aexit__( self, kind: type[BaseException] | None, value: BaseException | None, traceback: TracebackType | None, ) -> None: try: await self.aclose() except ReportingWriterError: if kind is None: raise if self.context.cancel.is_set(): raise asyncio.CancelledError if kind is None and datetime.now(timezone.utc) >= self.context.deadline_at: raise ReportingWriterError( ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", "unknown") ) async def aclose(self) -> None: if self._closed: return self._closed = True await _close_owned(self._close, self.context.close_timeout_seconds) async def write(self, content: ReportingPreparedRevision) -> ReportingDestinationLocator: raise failure("UNSUPPORTED_VERIFICATION") async def read_rows( self, locator: ReportingDestinationLocator, *, cursor: str | None, limit: int ) -> ReportingDestinationPage: raise failure("UNSUPPORTED_VERIFICATION") async def read_manifest(self, locator: ReportingDestinationLocator) -> bytes: raise failure("UNSUPPORTED_VERIFICATION") async def list_objects(self, locator: ReportingDestinationLocator) -> tuple[str, ...]: raise failure("UNSUPPORTED_VERIFICATION") def read_object( self, locator: ReportingDestinationLocator, *, object_ref: str ) -> AsyncIterator[bytes]: raise failure("UNSUPPORTED_VERIFICATION") async def observe_native_version( self, locator: ReportingDestinationLocator ) -> ReportingNativeObservation: raise failure("UNSUPPORTED_VERIFICATION")Single-use SDK-owned lifecycle around an adopter's private provider session.
Override only _open/_close and I/O methods. _close must tolerate partial _open, finish promptly, and remove any adopter-owned temporary spool. The SDK invokes it exactly once, shields cancellation, and joins all its tasks. No credentials may be placed on the public request, descriptor or results.
Ancestors
- abc.ABC
Subclasses
- adcp.reporting.materializer.reference._Session
Instance variables
prop context : ReportingIOContext-
Expand source code
@property def context(self) -> ReportingIOContext: return self._context prop phase : ReportingIOPhase-
Expand source code
@property def phase(self) -> ReportingIOPhase: return self._phase prop request : ReportingDestinationRequest-
Expand source code
@property def request(self) -> ReportingDestinationRequest: return self._request
Methods
async def aclose(self) ‑> None-
Expand source code
async def aclose(self) -> None: if self._closed: return self._closed = True await _close_owned(self._close, self.context.close_timeout_seconds) async def list_objects(self,
locator: ReportingDestinationLocator) ‑> tuple[str, ...]-
Expand source code
async def list_objects(self, locator: ReportingDestinationLocator) -> tuple[str, ...]: raise failure("UNSUPPORTED_VERIFICATION") async def observe_native_version(self,
locator: ReportingDestinationLocator) ‑> ReportingNativeObservation-
Expand source code
async def observe_native_version( self, locator: ReportingDestinationLocator ) -> ReportingNativeObservation: raise failure("UNSUPPORTED_VERIFICATION") async def read_manifest(self,
locator: ReportingDestinationLocator) ‑> bytes-
Expand source code
async def read_manifest(self, locator: ReportingDestinationLocator) -> bytes: raise failure("UNSUPPORTED_VERIFICATION") def read_object(self,
locator: ReportingDestinationLocator,
*,
object_ref: str) ‑> AsyncIterator[bytes]-
Expand source code
def read_object( self, locator: ReportingDestinationLocator, *, object_ref: str ) -> AsyncIterator[bytes]: raise failure("UNSUPPORTED_VERIFICATION") async def read_rows(self,
locator: ReportingDestinationLocator,
*,
cursor: str | None,
limit: int) ‑> ReportingDestinationPage-
Expand source code
async def read_rows( self, locator: ReportingDestinationLocator, *, cursor: str | None, limit: int ) -> ReportingDestinationPage: raise failure("UNSUPPORTED_VERIFICATION") async def write(self,
content: ReportingPreparedRevision) ‑> ReportingDestinationLocator-
Expand source code
async def write(self, content: ReportingPreparedRevision) -> ReportingDestinationLocator: raise failure("UNSUPPORTED_VERIFICATION")
class ReportingDestinationWriter (*args, **kwargs)-
Expand source code
class ReportingDestinationWriter(Protocol): @property def capabilities(self) -> tuple[ReportingWriterCapability, ...]: ... @property def production_eligible(self) -> bool: ...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
Subclasses
Instance variables
prop capabilities : tuple[ReportingWriterCapability, ...]-
Expand source code
@property def capabilities(self) -> tuple[ReportingWriterCapability, ...]: ... prop production_eligible : bool-
Expand source code
@property def production_eligible(self) -> bool: ...
class ReportingHeartbeat (*args, **kwargs)-
Expand source code
class ReportingHeartbeat(Protocol): """B2 owns any parallel lease heartbeat. B1 only calls this checkpoint.""" async def checkpoint(self) -> None: ...B2 owns any parallel lease heartbeat. B1 only calls this checkpoint.
Ancestors
- typing.Protocol
- typing.Generic
Methods
async def checkpoint(self) ‑> None-
Expand source code
async def checkpoint(self) -> None: ...
class ReportingIOContext (deadline_at: datetime,
cancel: asyncio.Event,
heartbeat: ReportingHeartbeat | None = None,
close_timeout_seconds: float = 5.0)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingIOContext: deadline_at: datetime cancel: asyncio.Event = field(repr=False) heartbeat: ReportingHeartbeat | None = field(default=None, repr=False) close_timeout_seconds: float = 5.0 def __post_init__(self) -> None: object.__setattr__(self, "deadline_at", aware_utc(self.deadline_at)) if not math.isfinite(self.close_timeout_seconds) or self.close_timeout_seconds <= 0: raise ValueError("close timeout must be finite and positive") async def run( self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect = "not_started" ) -> T: if self.heartbeat is not None: await self._run(self.heartbeat.checkpoint, effect=effect) return await self._run(call, effect=effect) async def _run(self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect) -> T: if self.cancel.is_set(): raise asyncio.CancelledError remaining = (self.deadline_at - datetime.now(timezone.utc)).total_seconds() if remaining <= 0: raise ReportingWriterError( ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) ) async def invoke() -> T: return await call() task = asyncio.create_task(invoke()) canceled = asyncio.create_task(self.cancel.wait()) problem: ReportingWriterFailure | None = None was_canceled = False try: done, _ = await asyncio.wait( (task, canceled), timeout=remaining, return_when=asyncio.FIRST_COMPLETED ) if canceled in done: was_canceled = True elif task not in done: problem = ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) else: try: result = task.result() except ReportingWriterError as exc: problem = exc.failure if ( effect == "unknown" and problem.effect == "not_started" and problem.code in {"RESOURCE_UNAVAILABLE", "DEADLINE_EXCEEDED", "LEASE_LOST"} ): problem = replace(problem, effect="unknown", retry="same_identity") except Exception: problem = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", effect ) except asyncio.CancelledError: was_canceled = True finally: task.cancel() canceled.cancel() if await _join_tasks(task, canceled): was_canceled = True if was_canceled: raise asyncio.CancelledError # Outside the except suite: no provider __context__, even when inspected. if problem is not None: raise ReportingWriterError(problem) return resultReportingIOContext(deadline_at: 'datetime', cancel: 'asyncio.Event', heartbeat: 'ReportingHeartbeat | None' = None, close_timeout_seconds: 'float' = 5.0)
Instance variables
var cancel : asyncio.locks.Event-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingIOContext: deadline_at: datetime cancel: asyncio.Event = field(repr=False) heartbeat: ReportingHeartbeat | None = field(default=None, repr=False) close_timeout_seconds: float = 5.0 def __post_init__(self) -> None: object.__setattr__(self, "deadline_at", aware_utc(self.deadline_at)) if not math.isfinite(self.close_timeout_seconds) or self.close_timeout_seconds <= 0: raise ValueError("close timeout must be finite and positive") async def run( self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect = "not_started" ) -> T: if self.heartbeat is not None: await self._run(self.heartbeat.checkpoint, effect=effect) return await self._run(call, effect=effect) async def _run(self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect) -> T: if self.cancel.is_set(): raise asyncio.CancelledError remaining = (self.deadline_at - datetime.now(timezone.utc)).total_seconds() if remaining <= 0: raise ReportingWriterError( ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) ) async def invoke() -> T: return await call() task = asyncio.create_task(invoke()) canceled = asyncio.create_task(self.cancel.wait()) problem: ReportingWriterFailure | None = None was_canceled = False try: done, _ = await asyncio.wait( (task, canceled), timeout=remaining, return_when=asyncio.FIRST_COMPLETED ) if canceled in done: was_canceled = True elif task not in done: problem = ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) else: try: result = task.result() except ReportingWriterError as exc: problem = exc.failure if ( effect == "unknown" and problem.effect == "not_started" and problem.code in {"RESOURCE_UNAVAILABLE", "DEADLINE_EXCEEDED", "LEASE_LOST"} ): problem = replace(problem, effect="unknown", retry="same_identity") except Exception: problem = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", effect ) except asyncio.CancelledError: was_canceled = True finally: task.cancel() canceled.cancel() if await _join_tasks(task, canceled): was_canceled = True if was_canceled: raise asyncio.CancelledError # Outside the except suite: no provider __context__, even when inspected. if problem is not None: raise ReportingWriterError(problem) return result var close_timeout_seconds : float-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingIOContext: deadline_at: datetime cancel: asyncio.Event = field(repr=False) heartbeat: ReportingHeartbeat | None = field(default=None, repr=False) close_timeout_seconds: float = 5.0 def __post_init__(self) -> None: object.__setattr__(self, "deadline_at", aware_utc(self.deadline_at)) if not math.isfinite(self.close_timeout_seconds) or self.close_timeout_seconds <= 0: raise ValueError("close timeout must be finite and positive") async def run( self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect = "not_started" ) -> T: if self.heartbeat is not None: await self._run(self.heartbeat.checkpoint, effect=effect) return await self._run(call, effect=effect) async def _run(self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect) -> T: if self.cancel.is_set(): raise asyncio.CancelledError remaining = (self.deadline_at - datetime.now(timezone.utc)).total_seconds() if remaining <= 0: raise ReportingWriterError( ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) ) async def invoke() -> T: return await call() task = asyncio.create_task(invoke()) canceled = asyncio.create_task(self.cancel.wait()) problem: ReportingWriterFailure | None = None was_canceled = False try: done, _ = await asyncio.wait( (task, canceled), timeout=remaining, return_when=asyncio.FIRST_COMPLETED ) if canceled in done: was_canceled = True elif task not in done: problem = ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) else: try: result = task.result() except ReportingWriterError as exc: problem = exc.failure if ( effect == "unknown" and problem.effect == "not_started" and problem.code in {"RESOURCE_UNAVAILABLE", "DEADLINE_EXCEEDED", "LEASE_LOST"} ): problem = replace(problem, effect="unknown", retry="same_identity") except Exception: problem = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", effect ) except asyncio.CancelledError: was_canceled = True finally: task.cancel() canceled.cancel() if await _join_tasks(task, canceled): was_canceled = True if was_canceled: raise asyncio.CancelledError # Outside the except suite: no provider __context__, even when inspected. if problem is not None: raise ReportingWriterError(problem) return result var deadline_at : datetime.datetime-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingIOContext: deadline_at: datetime cancel: asyncio.Event = field(repr=False) heartbeat: ReportingHeartbeat | None = field(default=None, repr=False) close_timeout_seconds: float = 5.0 def __post_init__(self) -> None: object.__setattr__(self, "deadline_at", aware_utc(self.deadline_at)) if not math.isfinite(self.close_timeout_seconds) or self.close_timeout_seconds <= 0: raise ValueError("close timeout must be finite and positive") async def run( self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect = "not_started" ) -> T: if self.heartbeat is not None: await self._run(self.heartbeat.checkpoint, effect=effect) return await self._run(call, effect=effect) async def _run(self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect) -> T: if self.cancel.is_set(): raise asyncio.CancelledError remaining = (self.deadline_at - datetime.now(timezone.utc)).total_seconds() if remaining <= 0: raise ReportingWriterError( ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) ) async def invoke() -> T: return await call() task = asyncio.create_task(invoke()) canceled = asyncio.create_task(self.cancel.wait()) problem: ReportingWriterFailure | None = None was_canceled = False try: done, _ = await asyncio.wait( (task, canceled), timeout=remaining, return_when=asyncio.FIRST_COMPLETED ) if canceled in done: was_canceled = True elif task not in done: problem = ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) else: try: result = task.result() except ReportingWriterError as exc: problem = exc.failure if ( effect == "unknown" and problem.effect == "not_started" and problem.code in {"RESOURCE_UNAVAILABLE", "DEADLINE_EXCEEDED", "LEASE_LOST"} ): problem = replace(problem, effect="unknown", retry="same_identity") except Exception: problem = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", effect ) except asyncio.CancelledError: was_canceled = True finally: task.cancel() canceled.cancel() if await _join_tasks(task, canceled): was_canceled = True if was_canceled: raise asyncio.CancelledError # Outside the except suite: no provider __context__, even when inspected. if problem is not None: raise ReportingWriterError(problem) return result var heartbeat : ReportingHeartbeat | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingIOContext: deadline_at: datetime cancel: asyncio.Event = field(repr=False) heartbeat: ReportingHeartbeat | None = field(default=None, repr=False) close_timeout_seconds: float = 5.0 def __post_init__(self) -> None: object.__setattr__(self, "deadline_at", aware_utc(self.deadline_at)) if not math.isfinite(self.close_timeout_seconds) or self.close_timeout_seconds <= 0: raise ValueError("close timeout must be finite and positive") async def run( self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect = "not_started" ) -> T: if self.heartbeat is not None: await self._run(self.heartbeat.checkpoint, effect=effect) return await self._run(call, effect=effect) async def _run(self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect) -> T: if self.cancel.is_set(): raise asyncio.CancelledError remaining = (self.deadline_at - datetime.now(timezone.utc)).total_seconds() if remaining <= 0: raise ReportingWriterError( ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) ) async def invoke() -> T: return await call() task = asyncio.create_task(invoke()) canceled = asyncio.create_task(self.cancel.wait()) problem: ReportingWriterFailure | None = None was_canceled = False try: done, _ = await asyncio.wait( (task, canceled), timeout=remaining, return_when=asyncio.FIRST_COMPLETED ) if canceled in done: was_canceled = True elif task not in done: problem = ReportingWriterFailure("DEADLINE_EXCEEDED", "same_identity", effect) else: try: result = task.result() except ReportingWriterError as exc: problem = exc.failure if ( effect == "unknown" and problem.effect == "not_started" and problem.code in {"RESOURCE_UNAVAILABLE", "DEADLINE_EXCEEDED", "LEASE_LOST"} ): problem = replace(problem, effect="unknown", retry="same_identity") except Exception: problem = ReportingWriterFailure( "RESOURCE_UNAVAILABLE", "same_identity", effect ) except asyncio.CancelledError: was_canceled = True finally: task.cancel() canceled.cancel() if await _join_tasks(task, canceled): was_canceled = True if was_canceled: raise asyncio.CancelledError # Outside the except suite: no provider __context__, even when inspected. if problem is not None: raise ReportingWriterError(problem) return result
Methods
async def run(self,
call: Callable[[], Awaitable[T]],
*,
effect: ReportingExternalEffect = 'not_started') ‑> ~T-
Expand source code
async def run( self, call: Callable[[], Awaitable[T]], *, effect: ReportingExternalEffect = "not_started" ) -> T: if self.heartbeat is not None: await self._run(self.heartbeat.checkpoint, effect=effect) return await self._run(call, effect=effect)
class ReportingMaterializationAttempt (scope: ReportingDeliveryScope,
reporting_revision_id: str,
reporting_materialization_id: str,
attempt: int,
created_at: datetime,
*,
kind: "Literal['materialization_attempt']" = 'materialization_attempt')-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationAttempt(_ClosedValue): scope: ReportingDeliveryScope reporting_revision_id: str reporting_materialization_id: str attempt: int created_at: datetime kind: Literal["materialization_attempt"] = field( default="materialization_attempt", kw_only=True ) def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id, maximum=255) reporting_identifier(self.reporting_materialization_id, maximum=255) _positive(self.attempt) object.__setattr__(self, "created_at", aware_utc(self.created_at)) @property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id)ReportingMaterializationAttempt(scope: 'ReportingDeliveryScope', reporting_revision_id: 'str', reporting_materialization_id: 'str', attempt: 'int', created_at: 'datetime', *, kind: "Literal['materialization_attempt']" = 'materialization_attempt')
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var attempt : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationAttempt(_ClosedValue): scope: ReportingDeliveryScope reporting_revision_id: str reporting_materialization_id: str attempt: int created_at: datetime kind: Literal["materialization_attempt"] = field( default="materialization_attempt", kw_only=True ) def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id, maximum=255) reporting_identifier(self.reporting_materialization_id, maximum=255) _positive(self.attempt) object.__setattr__(self, "created_at", aware_utc(self.created_at)) @property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id) var created_at : datetime.datetime-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationAttempt(_ClosedValue): scope: ReportingDeliveryScope reporting_revision_id: str reporting_materialization_id: str attempt: int created_at: datetime kind: Literal["materialization_attempt"] = field( default="materialization_attempt", kw_only=True ) def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id, maximum=255) reporting_identifier(self.reporting_materialization_id, maximum=255) _positive(self.attempt) object.__setattr__(self, "created_at", aware_utc(self.created_at)) @property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id) prop key : ReportingMaterializationKey-
Expand source code
@property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id) var kind : Literal['materialization_attempt']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationAttempt(_ClosedValue): scope: ReportingDeliveryScope reporting_revision_id: str reporting_materialization_id: str attempt: int created_at: datetime kind: Literal["materialization_attempt"] = field( default="materialization_attempt", kw_only=True ) def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id, maximum=255) reporting_identifier(self.reporting_materialization_id, maximum=255) _positive(self.attempt) object.__setattr__(self, "created_at", aware_utc(self.created_at)) @property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id) var reporting_materialization_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationAttempt(_ClosedValue): scope: ReportingDeliveryScope reporting_revision_id: str reporting_materialization_id: str attempt: int created_at: datetime kind: Literal["materialization_attempt"] = field( default="materialization_attempt", kw_only=True ) def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id, maximum=255) reporting_identifier(self.reporting_materialization_id, maximum=255) _positive(self.attempt) object.__setattr__(self, "created_at", aware_utc(self.created_at)) @property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id) var reporting_revision_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationAttempt(_ClosedValue): scope: ReportingDeliveryScope reporting_revision_id: str reporting_materialization_id: str attempt: int created_at: datetime kind: Literal["materialization_attempt"] = field( default="materialization_attempt", kw_only=True ) def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id, maximum=255) reporting_identifier(self.reporting_materialization_id, maximum=255) _positive(self.attempt) object.__setattr__(self, "created_at", aware_utc(self.created_at)) @property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id) var scope : ReportingDeliveryScope-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationAttempt(_ClosedValue): scope: ReportingDeliveryScope reporting_revision_id: str reporting_materialization_id: str attempt: int created_at: datetime kind: Literal["materialization_attempt"] = field( default="materialization_attempt", kw_only=True ) def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.reporting_revision_id, maximum=255) reporting_identifier(self.reporting_materialization_id, maximum=255) _positive(self.attempt) object.__setattr__(self, "created_at", aware_utc(self.created_at)) @property def key(self) -> ReportingMaterializationKey: return ReportingMaterializationKey(self.scope.principal, self.reporting_materialization_id)
class ReportingMaterializerBoundary (caller: ReportingDeliveryPrincipal,
sequence: int,
account_sequence: int,
reporting_materialization_id: str,
as_of: datetime,
core: ReportingStatusSnapshot,
reconciliation: tuple[ReportingDeliveryRecord, ...])-
Expand source code
@dataclass(frozen=True) class ReportingMaterializerBoundary: """Private frozen projection inputs; never returned as a public wire blob.""" caller: ReportingDeliveryPrincipal sequence: int account_sequence: int reporting_materialization_id: str as_of: datetime core: ReportingStatusSnapshot = field(repr=False) reconciliation: tuple[ReportingDeliveryRecord, ...] = field(repr=False) def __post_init__(self) -> None: outcomes = tuple( record for record in self.reconciliation if isinstance(record, ReportingMaterializationRecord) and record.reporting_materialization_id == self.reporting_materialization_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(outcomes) != 1 or outcomes[0].completed_at != self.as_of ): raise ReportingNotificationError("materializer_boundary_invalid") reporting_identifier(self.reporting_materialization_id, maximum=255) object.__setattr__(self, "as_of", aware_utc(self.as_of)) def to_storage(self) -> dict[str, Any]: return { "version": 1, "account_id": self.caller.account_id, "consumer_id": self.caller.consumer_id, "sequence": self.sequence, "account_sequence": self.account_sequence, "reporting_materialization_id": self.reporting_materialization_id, "as_of": self.as_of.isoformat(), "core": _core_adapter().dump_python(self.core, mode="json"), "reconciliation": [payload(r) for r in self.reconciliation], }Private frozen projection inputs; never returned as a public wire blob.
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_materialization_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, "account_id": self.caller.account_id, "consumer_id": self.caller.consumer_id, "sequence": self.sequence, "account_sequence": self.account_sequence, "reporting_materialization_id": self.reporting_materialization_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 ReportingMaterializerLease (scope: ReportingDeliveryScope,
generation: int,
token: str,
expires_at: datetime,
attempt: ReportingMaterializationAttempt,
request: ReportingDestinationRequest,
context: MaterializerContext,
notifications_enabled: bool,
admission_epoch: int = 0)-
Expand source code
@dataclass(frozen=True) class ReportingMaterializerLease: """A short fenced reservation. The token grants no destination authorization.""" scope: ReportingDeliveryScope generation: int token: str = field(repr=False) expires_at: datetime attempt: ReportingMaterializationAttempt request: ReportingDestinationRequest context: MaterializerContext notifications_enabled: bool admission_epoch: int = 0 def __post_init__(self) -> None: invalid = False try: invalid = ( type(self.generation) is not int or self.generation < 1 or type(self.notifications_enabled) is not bool or type(self.admission_epoch) is not int or self.admission_epoch not in {0, 2} or UUID(self.token).version != 4 or self.attempt.scope != self.scope or self.context.scope != self.scope or self.request != ReportingDestinationRequest.from_binding( self.context.binding, self.attempt, self.request.verification_key ) ) object.__setattr__(self, "expires_at", aware_utc(self.expires_at)) except (ValueError, TypeError, AttributeError): invalid = True if invalid: raise failure("BINDING_MISMATCH")A short fenced reservation. The token grants no destination authorization.
Instance variables
var admission_epoch : intvar attempt : ReportingMaterializationAttemptvar context : MaterializerContextvar expires_at : datetime.datetimevar generation : intvar notifications_enabled : boolvar request : ReportingDestinationRequestvar scope : ReportingDeliveryScopevar token : str
class ReportingMaterializerService (store: ReportingMaterializerStore,
io: ReportingDestinationIO,
writer: ReportingDestinationWriter,
lease_seconds: int = 30,
io_timeout_seconds: int = 300)-
Expand source code
@dataclass(frozen=True) class ReportingMaterializerService: """One autonomous turn; schedule repeatedly, including after process restart. The destination must implement its advertised conditional/idempotent write semantics durably. A timeout or interrupted outcome commit retains the same immutable attempt. The SDK opens fresh authorization sessions for write and paginated readback; the service checks current ledger authority before each. """ store: ReportingMaterializerStore = field(repr=False) io: ReportingDestinationIO = field(repr=False) writer: ReportingDestinationWriter = field(repr=False) lease_seconds: int = 30 io_timeout_seconds: int = 300 def __post_init__(self) -> None: validate_lease_seconds(self.lease_seconds) if type(self.io_timeout_seconds) is not int or not 1 <= self.io_timeout_seconds <= 3600: raise ValueError("materializer I/O deadline requires 1..3600 seconds") if type(self.io) is not ReportingDestinationIO: raise failure("UNSUPPORTED_VERIFICATION") @materializer_errors async def run_once(self) -> ReportingMaterializerTurn: keys = tuple( v.key for v in self.io.registry.verifiers if v.key.capability in self.writer.capabilities ) reserved = await self.store.claim_materialization( keys=keys, lease_seconds=self.lease_seconds ) if isinstance(reserved, ReportingMaterializerTurn): return reserved cancel = asyncio.Event() heartbeat = _LeaseHeartbeat(self.store, reserved, self.lease_seconds, cancel) task = asyncio.create_task(heartbeat.run()) context = ReportingIOContext( datetime.now(timezone.utc) + timedelta(seconds=self.io_timeout_seconds), cancel, heartbeat, ) canceled = False turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) try: try: await self.store.authorize_materialization(reserved) target = reserved.context if target.delivery is None: raise failure("BINDING_MISMATCH") prepared = await self.io.registry.prepare( key=reserved.request.verification_key, binding=target.binding, delivery=target.delivery, obligation=target.obligation, revisions=target.revisions, attempt=reserved.attempt, reader=self.store, context=context, ) await self.store.authorize_materialization(reserved) locator = await self.io.write(prepared, context=context) await self.store.authorize_materialization(reserved) verified = await self.io.verify(prepared, locator, context=context) turn = await self.store.finish_materialization( reserved, prepared=prepared, verified=verified ) except ReportingWriterError as exc: record = exc.failure # The writer owns whether a failure is known and terminal. All # uncertain I/O and lost leases keep the external identity. if record.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"}: record = ReportingWriterFailure(record.code, "same_identity", "unknown") turn = await self.store.finish_materialization(reserved, error=record) except asyncio.CancelledError: canceled = not heartbeat.lost except Exception: # Includes commit/connection loss and failed atomic outbox enqueue. # The lease expires; a restart reuses the original identity. turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) finally: heartbeat.stopped.set() task.cancel() canceled = await _join_tasks(task) or canceled if canceled: raise asyncio.CancelledError return turnOne autonomous turn; schedule repeatedly, including after process restart.
The destination must implement its advertised conditional/idempotent write semantics durably. A timeout or interrupted outcome commit retains the same immutable attempt. The SDK opens fresh authorization sessions for write and paginated readback; the service checks current ledger authority before each.
Instance variables
var io : ReportingDestinationIOvar io_timeout_seconds : intvar lease_seconds : intvar store : ReportingMaterializerStorevar writer : ReportingDestinationWriter
Methods
async def run_once(self) ‑> ReportingMaterializerTurn-
Expand source code
@materializer_errors async def run_once(self) -> ReportingMaterializerTurn: keys = tuple( v.key for v in self.io.registry.verifiers if v.key.capability in self.writer.capabilities ) reserved = await self.store.claim_materialization( keys=keys, lease_seconds=self.lease_seconds ) if isinstance(reserved, ReportingMaterializerTurn): return reserved cancel = asyncio.Event() heartbeat = _LeaseHeartbeat(self.store, reserved, self.lease_seconds, cancel) task = asyncio.create_task(heartbeat.run()) context = ReportingIOContext( datetime.now(timezone.utc) + timedelta(seconds=self.io_timeout_seconds), cancel, heartbeat, ) canceled = False turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) try: try: await self.store.authorize_materialization(reserved) target = reserved.context if target.delivery is None: raise failure("BINDING_MISMATCH") prepared = await self.io.registry.prepare( key=reserved.request.verification_key, binding=target.binding, delivery=target.delivery, obligation=target.obligation, revisions=target.revisions, attempt=reserved.attempt, reader=self.store, context=context, ) await self.store.authorize_materialization(reserved) locator = await self.io.write(prepared, context=context) await self.store.authorize_materialization(reserved) verified = await self.io.verify(prepared, locator, context=context) turn = await self.store.finish_materialization( reserved, prepared=prepared, verified=verified ) except ReportingWriterError as exc: record = exc.failure # The writer owns whether a failure is known and terminal. All # uncertain I/O and lost leases keep the external identity. if record.code in {"DEADLINE_EXCEEDED", "LEASE_LOST"}: record = ReportingWriterFailure(record.code, "same_identity", "unknown") turn = await self.store.finish_materialization(reserved, error=record) except asyncio.CancelledError: canceled = not heartbeat.lost except Exception: # Includes commit/connection loss and failed atomic outbox enqueue. # The lease expires; a restart reuses the original identity. turn = ReportingMaterializerTurn( "pending", "effect_unknown", reserved.attempt.reporting_materialization_id ) finally: heartbeat.stopped.set() task.cancel() canceled = await _join_tasks(task) or canceled if canceled: raise asyncio.CancelledError return turn
class ReportingMaterializerStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingMaterializerStore(ReportingRevisionRowReader, Protocol): """Opt-in service primitives. Old structural store implementations stay valid.""" async def claim_materialization( self, *, keys: tuple[ReportingVerificationKey, ...], lease_seconds: int = 30 ) -> ReportingMaterializerLease | ReportingMaterializerTurn: ... async def renew_materialization( self, lease: ReportingMaterializerLease, *, lease_seconds: int ) -> bool: ... async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None: ... async def finish_materialization( self, lease: ReportingMaterializerLease, *, prepared: ReportingPreparedRevision | None = None, verified: ReportingVerifiedDestination | None = None, error: ReportingWriterFailure | None = None, ) -> ReportingMaterializerTurn: ...Opt-in service primitives. Old structural store implementations stay valid.
Ancestors
- ReportingRevisionRowReader
- typing.Protocol
- typing.Generic
Methods
-
Expand source code
async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None: ... async def claim_materialization(self,
*,
keys: tuple[ReportingVerificationKey, ...],
lease_seconds: int = 30) ‑> ReportingMaterializerLease | ReportingMaterializerTurn-
Expand source code
async def claim_materialization( self, *, keys: tuple[ReportingVerificationKey, ...], lease_seconds: int = 30 ) -> ReportingMaterializerLease | ReportingMaterializerTurn: ... async def finish_materialization(self,
lease: ReportingMaterializerLease,
*,
prepared: ReportingPreparedRevision | None = None,
verified: ReportingVerifiedDestination | None = None,
error: ReportingWriterFailure | None = None) ‑> ReportingMaterializerTurn-
Expand source code
async def finish_materialization( self, lease: ReportingMaterializerLease, *, prepared: ReportingPreparedRevision | None = None, verified: ReportingVerifiedDestination | None = None, error: ReportingWriterFailure | None = None, ) -> ReportingMaterializerTurn: ... async def renew_materialization(self,
lease: ReportingMaterializerLease,
*,
lease_seconds: int) ‑> bool-
Expand source code
async def renew_materialization( self, lease: ReportingMaterializerLease, *, lease_seconds: int ) -> bool: ...
class ReportingMaterializerTurn (state: "Literal['idle', 'discovered', 'parked', 'pending', 'verified', 'failed']",
reason: MaterializerReason | None = None,
reporting_materialization_id: str | None = None)-
Expand source code
@dataclass(frozen=True) class ReportingMaterializerTurn: state: Literal["idle", "discovered", "parked", "pending", "verified", "failed"] reason: MaterializerReason | None = None reporting_materialization_id: str | None = None def __post_init__(self) -> None: if self.state not in {"idle", "discovered", "parked", "pending", "verified", "failed"} or ( self.reason is not None and self.reason not in get_args(MaterializerReason) ): raise ValueError("invalid materializer turn") if self.reporting_materialization_id is not None: reporting_identifier(self.reporting_materialization_id, maximum=255)ReportingMaterializerTurn(state: "Literal['idle', 'discovered', 'parked', 'pending', 'verified', 'failed']", reason: 'MaterializerReason | None' = None, reporting_materialization_id: 'str | None' = None)
Instance variables
var reason : Literal['ready', 'verified', 'retry', 'inactive', 'revision_not_ready', 'revision_unreadable', 'target_changed', 'history_corrupt', 'legacy_pending', 'legacy_terminal', 'component_unavailable', 'binding_changed', 'effect_unknown', 'operator_required'] | Nonevar reporting_materialization_id : str | Nonevar state : Literal['idle', 'discovered', 'parked', 'pending', 'verified', 'failed']
class ReportingNativeObservation (location: str,
native_version_ref: str,
verification_path: "Literal['representative_consumer', 'destination']")-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingNativeObservation(_ClosedValue): location: str native_version_ref: str verification_path: Literal["representative_consumer", "destination"] def __post_init__(self) -> None: _freeze_fields(self) from adcp.reporting.evidence import resource_location resource_location(self.location) native_version_reference(self.native_version_ref)ReportingNativeObservation(location: 'str', native_version_ref: 'str', verification_path: "Literal['representative_consumer', 'destination']")
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var location : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingNativeObservation(_ClosedValue): location: str native_version_ref: str verification_path: Literal["representative_consumer", "destination"] def __post_init__(self) -> None: _freeze_fields(self) from adcp.reporting.evidence import resource_location resource_location(self.location) native_version_reference(self.native_version_ref) var native_version_ref : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingNativeObservation(_ClosedValue): location: str native_version_ref: str verification_path: Literal["representative_consumer", "destination"] def __post_init__(self) -> None: _freeze_fields(self) from adcp.reporting.evidence import resource_location resource_location(self.location) native_version_reference(self.native_version_ref) var verification_path : Literal['representative_consumer', 'destination']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingNativeObservation(_ClosedValue): location: str native_version_ref: str verification_path: Literal["representative_consumer", "destination"] def __post_init__(self) -> None: _freeze_fields(self) from adcp.reporting.evidence import resource_location resource_location(self.location) native_version_reference(self.native_version_ref)
class ReportingObligationDeliveryRecord (scope: ReportingDeliveryScope,
currency: str,
resource_retained_until: datetime,
created_at: datetime,
*,
kind: "Literal['obligation_delivery']" = 'obligation_delivery')-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingObligationDeliveryRecord(_ClosedValue): scope: ReportingDeliveryScope currency: str resource_retained_until: datetime created_at: datetime kind: Literal["obligation_delivery"] = field(default="obligation_delivery", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) validate_currency(self.currency) object.__setattr__(self, "resource_retained_until", aware_utc(self.resource_retained_until)) object.__setattr__(self, "created_at", aware_utc(self.created_at))ReportingObligationDeliveryRecord(scope: 'ReportingDeliveryScope', currency: 'str', resource_retained_until: 'datetime', created_at: 'datetime', *, kind: "Literal['obligation_delivery']" = 'obligation_delivery')
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var created_at : datetime.datetime-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingObligationDeliveryRecord(_ClosedValue): scope: ReportingDeliveryScope currency: str resource_retained_until: datetime created_at: datetime kind: Literal["obligation_delivery"] = field(default="obligation_delivery", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) validate_currency(self.currency) object.__setattr__(self, "resource_retained_until", aware_utc(self.resource_retained_until)) object.__setattr__(self, "created_at", aware_utc(self.created_at)) var currency : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingObligationDeliveryRecord(_ClosedValue): scope: ReportingDeliveryScope currency: str resource_retained_until: datetime created_at: datetime kind: Literal["obligation_delivery"] = field(default="obligation_delivery", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) validate_currency(self.currency) object.__setattr__(self, "resource_retained_until", aware_utc(self.resource_retained_until)) object.__setattr__(self, "created_at", aware_utc(self.created_at)) var kind : Literal['obligation_delivery']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingObligationDeliveryRecord(_ClosedValue): scope: ReportingDeliveryScope currency: str resource_retained_until: datetime created_at: datetime kind: Literal["obligation_delivery"] = field(default="obligation_delivery", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) validate_currency(self.currency) object.__setattr__(self, "resource_retained_until", aware_utc(self.resource_retained_until)) object.__setattr__(self, "created_at", aware_utc(self.created_at)) var resource_retained_until : datetime.datetime-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingObligationDeliveryRecord(_ClosedValue): scope: ReportingDeliveryScope currency: str resource_retained_until: datetime created_at: datetime kind: Literal["obligation_delivery"] = field(default="obligation_delivery", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) validate_currency(self.currency) object.__setattr__(self, "resource_retained_until", aware_utc(self.resource_retained_until)) object.__setattr__(self, "created_at", aware_utc(self.created_at)) var scope : ReportingDeliveryScope-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingObligationDeliveryRecord(_ClosedValue): scope: ReportingDeliveryScope currency: str resource_retained_until: datetime created_at: datetime kind: Literal["obligation_delivery"] = field(default="obligation_delivery", kw_only=True) def __post_init__(self) -> None: _freeze_fields(self) validate_currency(self.currency) object.__setattr__(self, "resource_retained_until", aware_utc(self.resource_retained_until)) object.__setattr__(self, "created_at", aware_utc(self.created_at))
class ReportingPreparedRevision (request: ReportingDestinationRequest,
obligation: ReportingObligationRecord,
revision: ReportingRevisionRecord,
delivery: ReportingObligationDeliveryRecord,
binding: ReportingDestinationBinding,
rows: tuple[bytes, ...])-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingPreparedRevision(_ClosedValue): """Bounded immutable canonical row bytes, produced before resolver I/O. This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly. """ request: ReportingDestinationRequest obligation: ReportingObligationRecord = field(repr=False) revision: ReportingRevisionRecord = field(repr=False) delivery: ReportingObligationDeliveryRecord = field(repr=False) binding: ReportingDestinationBinding = field(repr=False) rows: tuple[bytes, ...] = field(repr=False) def __post_init__(self) -> None: _freeze_fields(self)Bounded immutable canonical row bytes, produced before resolver I/O.
This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly.
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var binding : ReportingDestinationBinding-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingPreparedRevision(_ClosedValue): """Bounded immutable canonical row bytes, produced before resolver I/O. This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly. """ request: ReportingDestinationRequest obligation: ReportingObligationRecord = field(repr=False) revision: ReportingRevisionRecord = field(repr=False) delivery: ReportingObligationDeliveryRecord = field(repr=False) binding: ReportingDestinationBinding = field(repr=False) rows: tuple[bytes, ...] = field(repr=False) def __post_init__(self) -> None: _freeze_fields(self) var delivery : ReportingObligationDeliveryRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingPreparedRevision(_ClosedValue): """Bounded immutable canonical row bytes, produced before resolver I/O. This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly. """ request: ReportingDestinationRequest obligation: ReportingObligationRecord = field(repr=False) revision: ReportingRevisionRecord = field(repr=False) delivery: ReportingObligationDeliveryRecord = field(repr=False) binding: ReportingDestinationBinding = field(repr=False) rows: tuple[bytes, ...] = field(repr=False) def __post_init__(self) -> None: _freeze_fields(self) var obligation : ReportingObligationRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingPreparedRevision(_ClosedValue): """Bounded immutable canonical row bytes, produced before resolver I/O. This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly. """ request: ReportingDestinationRequest obligation: ReportingObligationRecord = field(repr=False) revision: ReportingRevisionRecord = field(repr=False) delivery: ReportingObligationDeliveryRecord = field(repr=False) binding: ReportingDestinationBinding = field(repr=False) rows: tuple[bytes, ...] = field(repr=False) def __post_init__(self) -> None: _freeze_fields(self) var request : ReportingDestinationRequest-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingPreparedRevision(_ClosedValue): """Bounded immutable canonical row bytes, produced before resolver I/O. This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly. """ request: ReportingDestinationRequest obligation: ReportingObligationRecord = field(repr=False) revision: ReportingRevisionRecord = field(repr=False) delivery: ReportingObligationDeliveryRecord = field(repr=False) binding: ReportingDestinationBinding = field(repr=False) rows: tuple[bytes, ...] = field(repr=False) def __post_init__(self) -> None: _freeze_fields(self) var revision : ReportingRevisionRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingPreparedRevision(_ClosedValue): """Bounded immutable canonical row bytes, produced before resolver I/O. This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly. """ request: ReportingDestinationRequest obligation: ReportingObligationRecord = field(repr=False) revision: ReportingRevisionRecord = field(repr=False) delivery: ReportingObligationDeliveryRecord = field(repr=False) binding: ReportingDestinationBinding = field(repr=False) rows: tuple[bytes, ...] = field(repr=False) def __post_init__(self) -> None: _freeze_fields(self) var rows : tuple[bytes, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingPreparedRevision(_ClosedValue): """Bounded immutable canonical row bytes, produced before resolver I/O. This is not a durable reservation. Keep the original attempt identity when an external effect is unknown. Never automatically import legacy pending attempts: drain old writers and recover/import their identities explicitly. """ request: ReportingDestinationRequest obligation: ReportingObligationRecord = field(repr=False) revision: ReportingRevisionRecord = field(repr=False) delivery: ReportingObligationDeliveryRecord = field(repr=False) binding: ReportingDestinationBinding = field(repr=False) rows: tuple[bytes, ...] = field(repr=False) def __post_init__(self) -> None: _freeze_fields(self)
class ReportingRevisionRowReader (*args, **kwargs)-
Expand source code
class ReportingRevisionRowReader(Protocol): async def read_revision_rows( self, *, account_id: str, reporting_revision_id: str, cursor: str | None = None, limit: int = 500, ) -> ReportingRowPage: ...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
Subclasses
Methods
async def read_revision_rows(self,
*,
account_id: str,
reporting_revision_id: str,
cursor: str | None = None,
limit: int = 500) ‑> ReportingRowPage-
Expand source code
async def read_revision_rows( self, *, account_id: str, reporting_revision_id: str, cursor: str | None = None, limit: int = 500, ) -> ReportingRowPage: ...
class ReportingRevisionVerifier (key: ReportingVerificationKey,
definition_bytes: bytes,
schema_bytes: bytes,
canonicalization_bytes: bytes,
limits: ReportingVerificationLimits = <factory>)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows)One exact installed contract using the SDK's safe-integer JCS subset.
Sum totals are declared on schema properties with
x-adcp-control-total(value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction.Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var canonicalization_bytes : bytes-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var definition_bytes : bytes-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var key : ReportingVerificationKey-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var limits : adcp.reporting.materializer._json.ReportingVerificationLimits-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows) var schema_bytes : bytes-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifier(_ClosedValue): """One exact installed contract using the SDK's safe-integer JCS subset. Sum totals are declared on schema properties with ``x-adcp-control-total`` (value_type, optional unit, decimal scale). This explicit pinned subset has no expression evaluator, custom callbacks, mutable registration or network. Unsupported formats and definition semantics fail during construction. """ key: ReportingVerificationKey definition_bytes: bytes = field(repr=False) schema_bytes: bytes = field(repr=False) canonicalization_bytes: bytes = field(repr=False) limits: ReportingVerificationLimits = field(default_factory=ReportingVerificationLimits) def __post_init__(self) -> None: _freeze_fields(self) error = False try: self._runtime() except Exception: error = True if error: raise failure("UNSUPPORTED_VERIFICATION") def _runtime(self) -> _Canonicalizer: key, definition = self.key, self.key.definition if key.capability.format not in {"jsonl", None}: raise failure("UNSUPPORTED_VERIFICATION") if ( key.capability.method == "file_transfer" and key.capability.immutability != "immutable_location" ): raise failure("UNSUPPORTED_VERIFICATION") if ( definition.schema_dialect != "https://json-schema.org/draft/2020-12/schema" or definition.schema_ref_policy != "local_fragment_only" ): raise failure("UNSUPPORTED_VERIFICATION") for raw, expected in ( (self.definition_bytes, definition.report_definition_sha256), (self.schema_bytes, definition.schema_sha256), (self.canonicalization_bytes, key.canonicalization.canonicalization_sha256), ): if hashlib.sha256(raw).hexdigest() != expected.lower(): raise failure("UNSUPPORTED_VERIFICATION") report = _mapping(parse_reporting_json(self.definition_bytes, self.limits)) schema = _mapping(parse_reporting_json(self.schema_bytes, self.limits)) contract = _mapping(parse_reporting_json(self.canonicalization_bytes, self.limits)) if ( report.get("report_definition_id") != key.report_definition_id or report.get("reporting_profile") != key.reporting_profile or schema.get("$schema") != definition.schema_dialect or set(contract) != { "contract_version", "media_type", "algorithm", "schema_sha256", "primary_keys", "golden_vectors", } or contract["contract_version"] != "1.0" or contract["media_type"] != "application/vnd.adcp.reporting-canonicalization+json" or contract["algorithm"] != "adcp_jcs_rows_v1" or contract["schema_sha256"].lower() != definition.schema_sha256.lower() ): raise failure("UNSUPPORTED_VERIFICATION") _local_schema(schema) Draft202012Validator.check_schema(schema) keys = contract["primary_keys"] if ( type(keys) is not list or not keys or any(type(k) is not str for k in keys) or len(set(keys)) != len(keys) ): raise failure("UNSUPPORTED_VERIFICATION") totals: list[_Total] = [] properties = _mapping(schema.get("properties")) metrics = report.get("metrics") if type(metrics) is not list or not metrics: raise failure("UNSUPPORTED_VERIFICATION") for metric in metrics: metric = _mapping(metric) name = metric["name"] prop = _mapping(properties[name]) total = _mapping(prop["x-adcp-control-total"]) if ( metric.get("source_expression") != name or metric.get("aggregation") != "sum" or set(total) - {"value_type", "unit", "scale"} or total.get("unit") != metric.get("unit") or (total["value_type"], prop.get("type")) not in {("integer", "integer"), ("decimal", "string")} ): raise failure("UNSUPPORTED_VERIFICATION") scale = total.get("scale", 0) if ( type(scale) is not int or not 0 <= scale <= 18 or (total["value_type"] == "integer" and scale != 0) ): raise failure("UNSUPPORTED_VERIFICATION") totals.append(_Total(name, total["value_type"], total.get("unit"), scale)) if len({t.name for t in totals}) != len(totals): raise failure("UNSUPPORTED_VERIFICATION") if {t.name for t in totals} != { name for name, prop in properties.items() if type(prop) is dict and "x-adcp-control-total" in prop }: raise failure("UNSUPPORTED_VERIFICATION") if not set(keys).union(t.name for t in totals) <= set(schema.get("required", [])): raise failure("UNSUPPORTED_VERIFICATION") units = {t.name: t.unit for t in totals} if any( units.get(name) != unit for name, unit in ( *definition.monetary_metric_units, *definition.monetary_control_total_units, ) ): raise failure("UNSUPPORTED_VERIFICATION") runtime = _Canonicalizer( tuple(keys), tuple(totals), Draft202012Validator(schema), self.limits ) vectors = _mapping(contract["golden_vectors"]) if not {"empty_report", "ordering_encoding"} <= set(vectors) or set(vectors) - { "empty_report", "ordering_encoding", "additional", }: raise failure("UNSUPPORTED_VERIFICATION") seen: set[str] = set() for vector in [ vectors["empty_report"], vectors["ordering_encoding"], *vectors.get("additional", []), ]: vector = _mapping(vector) if ( set(vector) != {"name", "purpose", "input_rows", "canonical_utf8_base64", "sha256"} or vector["name"] in seen ): raise failure("UNSUPPORTED_VERIFICATION") seen.add(vector["name"]) rows, _ = runtime.canonicalize(vector["input_rows"]) golden_bytes = base64.b64decode(vector["canonical_utf8_base64"], validate=True) if ( b"[" + b",".join(rows) + b"]" != golden_bytes or hashlib.sha256(golden_bytes).hexdigest() != vector["sha256"].lower() ): raise failure("UNSUPPORTED_VERIFICATION") if ( vectors["empty_report"]["input_rows"] != [] or vectors["empty_report"]["purpose"] != "empty_report" ): raise failure("UNSUPPORTED_VERIFICATION") ordering = vectors["ordering_encoding"] order_rows = ordering["input_rows"] if ( ordering["purpose"] != "ordering_encoding" or len(order_rows) < 2 or [runtime.order_key(r) for r in order_rows] == sorted(runtime.order_key(r) for r in order_rows) or not _unordered_members(order_rows) ): raise failure("UNSUPPORTED_VERIFICATION") return runtime def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows)
Methods
def canonicalize(self, rows: Sequence[object]) ‑> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]-
Expand source code
def canonicalize( self, rows: Sequence[object] ) -> tuple[tuple[bytes, ...], tuple[ReportingControlTotalRecord, ...]]: """Useful to trusted publishers preparing expected evidence before storage.""" return self._runtime().canonicalize(rows)Useful to trusted publishers preparing expected evidence before storage.
class ReportingRevisionVerifierRegistry (verifiers: tuple[ReportingRevisionVerifier, ...])-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifierRegistry(_ClosedValue): verifiers: tuple[ReportingRevisionVerifier, ...] def __post_init__(self) -> None: _freeze_fields(self) if len({v.key for v in self.verifiers}) != len(self.verifiers): raise failure("UNSUPPORTED_VERIFICATION") def require(self, key: ReportingVerificationKey) -> ReportingRevisionVerifier: for verifier in self.verifiers: if verifier.key == key: return verifier raise failure("UNSUPPORTED_VERIFICATION") async def prepare( self, *, key: ReportingVerificationKey, binding: ReportingDestinationBinding, delivery: ReportingObligationDeliveryRecord, obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], attempt: ReportingMaterializationAttempt, reader: ReportingRevisionRowReader, context: ReportingIOContext, ) -> ReportingPreparedRevision: verifier = self.require(key) # Unsupported tuples fail before any source/resolver I/O. selection = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if selection.kind != "selected": raise failure("REVISION_NOT_READY") revision = selection.revision _revision_metadata(revision) if not revision.readable: raise failure("REVISION_NOT_READY") request = ReportingDestinationRequest.from_binding(binding, attempt, key) if ( not _same_definition(key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (key.report_definition_id, key.reporting_profile) or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != revision.reporting_revision_id or delivery.scope != attempt.scope or delivery.currency != obligation.currency ): raise failure("BINDING_MISMATCH") _expected_digest(revision, key) rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() budget = _Budget(verifier.limits) while True: budget.page() page = await context.run( lambda: reader.read_revision_rows( account_id=obligation.account_id, reporting_revision_id=revision.reporting_revision_id, cursor=cursor, limit=500, ) ) if ( type(page) is not ReportingRowPage or type(page.rows) is not tuple or page.reporting_revision_id != revision.reporting_revision_id ): raise failure("SOURCE_INVALID") _page( page.total_count, page.has_more, page.cursor, len(page.rows), len(rows), revision.row_count, seen, ) if page.cursor is not None: cursor_invalid = False try: cursor_invalid = revision_row_offset( page.cursor, revision.reporting_revision_id, 500 ) != len(rows) + len(page.rows) except (ValueError, LedgerConflictError): cursor_invalid = True if cursor_invalid: raise failure("SOURCE_INVALID") for row in page.rows: encoded = strict_reporting_json(row, verifier.limits) budget.add(encoded) rows.append(_mapping(parse_reporting_json(encoded, verifier.limits))) if not page.has_more: break cursor = page.cursor if ( revision_content_sha256( reporting_revision_id=revision.reporting_revision_id, row_count=revision.row_count, control_totals=revision.control_totals, reporting_rows=rows, control_total_evidence=revision.managed_control_totals, ) != revision.revision_content_sha256.lower() ): raise failure("SOURCE_INVALID") canonical_rows, totals = verifier.canonicalize(rows) _verify_content(revision, key, canonical_rows, totals) await context.run(_checkpoint) return ReportingPreparedRevision( request, obligation, revision, delivery, binding, canonical_rows )ReportingRevisionVerifierRegistry(verifiers: 'tuple[ReportingRevisionVerifier, …]')
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var verifiers : tuple[ReportingRevisionVerifier, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingRevisionVerifierRegistry(_ClosedValue): verifiers: tuple[ReportingRevisionVerifier, ...] def __post_init__(self) -> None: _freeze_fields(self) if len({v.key for v in self.verifiers}) != len(self.verifiers): raise failure("UNSUPPORTED_VERIFICATION") def require(self, key: ReportingVerificationKey) -> ReportingRevisionVerifier: for verifier in self.verifiers: if verifier.key == key: return verifier raise failure("UNSUPPORTED_VERIFICATION") async def prepare( self, *, key: ReportingVerificationKey, binding: ReportingDestinationBinding, delivery: ReportingObligationDeliveryRecord, obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], attempt: ReportingMaterializationAttempt, reader: ReportingRevisionRowReader, context: ReportingIOContext, ) -> ReportingPreparedRevision: verifier = self.require(key) # Unsupported tuples fail before any source/resolver I/O. selection = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if selection.kind != "selected": raise failure("REVISION_NOT_READY") revision = selection.revision _revision_metadata(revision) if not revision.readable: raise failure("REVISION_NOT_READY") request = ReportingDestinationRequest.from_binding(binding, attempt, key) if ( not _same_definition(key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (key.report_definition_id, key.reporting_profile) or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != revision.reporting_revision_id or delivery.scope != attempt.scope or delivery.currency != obligation.currency ): raise failure("BINDING_MISMATCH") _expected_digest(revision, key) rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() budget = _Budget(verifier.limits) while True: budget.page() page = await context.run( lambda: reader.read_revision_rows( account_id=obligation.account_id, reporting_revision_id=revision.reporting_revision_id, cursor=cursor, limit=500, ) ) if ( type(page) is not ReportingRowPage or type(page.rows) is not tuple or page.reporting_revision_id != revision.reporting_revision_id ): raise failure("SOURCE_INVALID") _page( page.total_count, page.has_more, page.cursor, len(page.rows), len(rows), revision.row_count, seen, ) if page.cursor is not None: cursor_invalid = False try: cursor_invalid = revision_row_offset( page.cursor, revision.reporting_revision_id, 500 ) != len(rows) + len(page.rows) except (ValueError, LedgerConflictError): cursor_invalid = True if cursor_invalid: raise failure("SOURCE_INVALID") for row in page.rows: encoded = strict_reporting_json(row, verifier.limits) budget.add(encoded) rows.append(_mapping(parse_reporting_json(encoded, verifier.limits))) if not page.has_more: break cursor = page.cursor if ( revision_content_sha256( reporting_revision_id=revision.reporting_revision_id, row_count=revision.row_count, control_totals=revision.control_totals, reporting_rows=rows, control_total_evidence=revision.managed_control_totals, ) != revision.revision_content_sha256.lower() ): raise failure("SOURCE_INVALID") canonical_rows, totals = verifier.canonicalize(rows) _verify_content(revision, key, canonical_rows, totals) await context.run(_checkpoint) return ReportingPreparedRevision( request, obligation, revision, delivery, binding, canonical_rows )
Methods
async def prepare(self,
*,
key: ReportingVerificationKey,
binding: ReportingDestinationBinding,
delivery: ReportingObligationDeliveryRecord,
obligation: ReportingObligationRecord,
revisions: Sequence[ReportingRevisionRecord],
attempt: ReportingMaterializationAttempt,
reader: ReportingRevisionRowReader,
context: ReportingIOContext) ‑> ReportingPreparedRevision-
Expand source code
async def prepare( self, *, key: ReportingVerificationKey, binding: ReportingDestinationBinding, delivery: ReportingObligationDeliveryRecord, obligation: ReportingObligationRecord, revisions: Sequence[ReportingRevisionRecord], attempt: ReportingMaterializationAttempt, reader: ReportingRevisionRowReader, context: ReportingIOContext, ) -> ReportingPreparedRevision: verifier = self.require(key) # Unsupported tuples fail before any source/resolver I/O. selection = select_reporting_revision( revisions, account_id=obligation.account_id, reporting_obligation_id=obligation.reporting_obligation_id, required_finality=obligation.required_finality, ) if selection.kind == "corrupt": raise failure("HISTORY_CORRUPT") if selection.kind != "selected": raise failure("REVISION_NOT_READY") revision = selection.revision _revision_metadata(revision) if not revision.readable: raise failure("REVISION_NOT_READY") request = ReportingDestinationRequest.from_binding(binding, attempt, key) if ( not _same_definition(key, obligation.definition) or (obligation.report_definition_id, obligation.reporting_profile) != (key.report_definition_id, key.reporting_profile) or request.generation != obligation.generation_key or request.reporting_obligation_id != obligation.reporting_obligation_id or request.reporting_revision_id != revision.reporting_revision_id or delivery.scope != attempt.scope or delivery.currency != obligation.currency ): raise failure("BINDING_MISMATCH") _expected_digest(revision, key) rows: list[dict[str, Any]] = [] cursor: str | None = None seen: set[str] = set() budget = _Budget(verifier.limits) while True: budget.page() page = await context.run( lambda: reader.read_revision_rows( account_id=obligation.account_id, reporting_revision_id=revision.reporting_revision_id, cursor=cursor, limit=500, ) ) if ( type(page) is not ReportingRowPage or type(page.rows) is not tuple or page.reporting_revision_id != revision.reporting_revision_id ): raise failure("SOURCE_INVALID") _page( page.total_count, page.has_more, page.cursor, len(page.rows), len(rows), revision.row_count, seen, ) if page.cursor is not None: cursor_invalid = False try: cursor_invalid = revision_row_offset( page.cursor, revision.reporting_revision_id, 500 ) != len(rows) + len(page.rows) except (ValueError, LedgerConflictError): cursor_invalid = True if cursor_invalid: raise failure("SOURCE_INVALID") for row in page.rows: encoded = strict_reporting_json(row, verifier.limits) budget.add(encoded) rows.append(_mapping(parse_reporting_json(encoded, verifier.limits))) if not page.has_more: break cursor = page.cursor if ( revision_content_sha256( reporting_revision_id=revision.reporting_revision_id, row_count=revision.row_count, control_totals=revision.control_totals, reporting_rows=rows, control_total_evidence=revision.managed_control_totals, ) != revision.revision_content_sha256.lower() ): raise failure("SOURCE_INVALID") canonical_rows, totals = verifier.canonicalize(rows) _verify_content(revision, key, canonical_rows, totals) await context.run(_checkpoint) return ReportingPreparedRevision( request, obligation, revision, delivery, binding, canonical_rows ) def require(self,
key: ReportingVerificationKey) ‑> ReportingRevisionVerifier-
Expand source code
def require(self, key: ReportingVerificationKey) -> ReportingRevisionVerifier: for verifier in self.verifiers: if verifier.key == key: return verifier raise failure("UNSUPPORTED_VERIFICATION")
class ReportingVerificationKey (report_definition_id: str,
reporting_profile: str,
definition: ReportingDefinitionBinding,
canonicalization: ReportingCanonicalization,
capability: ReportingWriterCapability)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationKey(_ClosedValue): """Complete frozen identity, including method/format/profile/path and schema.""" report_definition_id: str reporting_profile: str definition: ReportingDefinitionBinding canonicalization: ReportingCanonicalization capability: ReportingWriterCapability def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.report_definition_id) reporting_identifier(self.reporting_profile, maximum=128) for value in (self.definition.report_definition_sha256, self.definition.schema_sha256): if type(value) is not str: raise ValueError("definition digests require exact strings") sha256_value(value) for uri in (self.definition.report_definition_uri, self.definition.schema_uri): # Reuse the strict public HTTPS contract screen; never dereference it. ReportingCanonicalDigest("0" * 64, "definition", uri, "0" * 64) for value in (self.definition.schema_version, self.definition.schema_ref_policy): reporting_identifier(value, maximum=128) ReportingCanonicalDigest("0" * 64, "dialect", self.definition.schema_dialect, "0" * 64) for units in ( self.definition.monetary_metric_units, self.definition.monetary_control_total_units, ): for name, unit in units: reporting_identifier(name, maximum=128) reporting_identifier(unit, maximum=32) object.__setattr__( self, "definition", replace( self.definition, report_definition_sha256=self.definition.report_definition_sha256.lower(), schema_sha256=self.definition.schema_sha256.lower(), ), )Complete frozen identity, including method/format/profile/path and schema.
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var canonicalization : ReportingCanonicalization-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationKey(_ClosedValue): """Complete frozen identity, including method/format/profile/path and schema.""" report_definition_id: str reporting_profile: str definition: ReportingDefinitionBinding canonicalization: ReportingCanonicalization capability: ReportingWriterCapability def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.report_definition_id) reporting_identifier(self.reporting_profile, maximum=128) for value in (self.definition.report_definition_sha256, self.definition.schema_sha256): if type(value) is not str: raise ValueError("definition digests require exact strings") sha256_value(value) for uri in (self.definition.report_definition_uri, self.definition.schema_uri): # Reuse the strict public HTTPS contract screen; never dereference it. ReportingCanonicalDigest("0" * 64, "definition", uri, "0" * 64) for value in (self.definition.schema_version, self.definition.schema_ref_policy): reporting_identifier(value, maximum=128) ReportingCanonicalDigest("0" * 64, "dialect", self.definition.schema_dialect, "0" * 64) for units in ( self.definition.monetary_metric_units, self.definition.monetary_control_total_units, ): for name, unit in units: reporting_identifier(name, maximum=128) reporting_identifier(unit, maximum=32) object.__setattr__( self, "definition", replace( self.definition, report_definition_sha256=self.definition.report_definition_sha256.lower(), schema_sha256=self.definition.schema_sha256.lower(), ), ) var capability : ReportingWriterCapability-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationKey(_ClosedValue): """Complete frozen identity, including method/format/profile/path and schema.""" report_definition_id: str reporting_profile: str definition: ReportingDefinitionBinding canonicalization: ReportingCanonicalization capability: ReportingWriterCapability def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.report_definition_id) reporting_identifier(self.reporting_profile, maximum=128) for value in (self.definition.report_definition_sha256, self.definition.schema_sha256): if type(value) is not str: raise ValueError("definition digests require exact strings") sha256_value(value) for uri in (self.definition.report_definition_uri, self.definition.schema_uri): # Reuse the strict public HTTPS contract screen; never dereference it. ReportingCanonicalDigest("0" * 64, "definition", uri, "0" * 64) for value in (self.definition.schema_version, self.definition.schema_ref_policy): reporting_identifier(value, maximum=128) ReportingCanonicalDigest("0" * 64, "dialect", self.definition.schema_dialect, "0" * 64) for units in ( self.definition.monetary_metric_units, self.definition.monetary_control_total_units, ): for name, unit in units: reporting_identifier(name, maximum=128) reporting_identifier(unit, maximum=32) object.__setattr__( self, "definition", replace( self.definition, report_definition_sha256=self.definition.report_definition_sha256.lower(), schema_sha256=self.definition.schema_sha256.lower(), ), ) var definition : ReportingDefinitionBinding-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationKey(_ClosedValue): """Complete frozen identity, including method/format/profile/path and schema.""" report_definition_id: str reporting_profile: str definition: ReportingDefinitionBinding canonicalization: ReportingCanonicalization capability: ReportingWriterCapability def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.report_definition_id) reporting_identifier(self.reporting_profile, maximum=128) for value in (self.definition.report_definition_sha256, self.definition.schema_sha256): if type(value) is not str: raise ValueError("definition digests require exact strings") sha256_value(value) for uri in (self.definition.report_definition_uri, self.definition.schema_uri): # Reuse the strict public HTTPS contract screen; never dereference it. ReportingCanonicalDigest("0" * 64, "definition", uri, "0" * 64) for value in (self.definition.schema_version, self.definition.schema_ref_policy): reporting_identifier(value, maximum=128) ReportingCanonicalDigest("0" * 64, "dialect", self.definition.schema_dialect, "0" * 64) for units in ( self.definition.monetary_metric_units, self.definition.monetary_control_total_units, ): for name, unit in units: reporting_identifier(name, maximum=128) reporting_identifier(unit, maximum=32) object.__setattr__( self, "definition", replace( self.definition, report_definition_sha256=self.definition.report_definition_sha256.lower(), schema_sha256=self.definition.schema_sha256.lower(), ), ) var report_definition_id : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationKey(_ClosedValue): """Complete frozen identity, including method/format/profile/path and schema.""" report_definition_id: str reporting_profile: str definition: ReportingDefinitionBinding canonicalization: ReportingCanonicalization capability: ReportingWriterCapability def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.report_definition_id) reporting_identifier(self.reporting_profile, maximum=128) for value in (self.definition.report_definition_sha256, self.definition.schema_sha256): if type(value) is not str: raise ValueError("definition digests require exact strings") sha256_value(value) for uri in (self.definition.report_definition_uri, self.definition.schema_uri): # Reuse the strict public HTTPS contract screen; never dereference it. ReportingCanonicalDigest("0" * 64, "definition", uri, "0" * 64) for value in (self.definition.schema_version, self.definition.schema_ref_policy): reporting_identifier(value, maximum=128) ReportingCanonicalDigest("0" * 64, "dialect", self.definition.schema_dialect, "0" * 64) for units in ( self.definition.monetary_metric_units, self.definition.monetary_control_total_units, ): for name, unit in units: reporting_identifier(name, maximum=128) reporting_identifier(unit, maximum=32) object.__setattr__( self, "definition", replace( self.definition, report_definition_sha256=self.definition.report_definition_sha256.lower(), schema_sha256=self.definition.schema_sha256.lower(), ), ) var reporting_profile : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationKey(_ClosedValue): """Complete frozen identity, including method/format/profile/path and schema.""" report_definition_id: str reporting_profile: str definition: ReportingDefinitionBinding canonicalization: ReportingCanonicalization capability: ReportingWriterCapability def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.report_definition_id) reporting_identifier(self.reporting_profile, maximum=128) for value in (self.definition.report_definition_sha256, self.definition.schema_sha256): if type(value) is not str: raise ValueError("definition digests require exact strings") sha256_value(value) for uri in (self.definition.report_definition_uri, self.definition.schema_uri): # Reuse the strict public HTTPS contract screen; never dereference it. ReportingCanonicalDigest("0" * 64, "definition", uri, "0" * 64) for value in (self.definition.schema_version, self.definition.schema_ref_policy): reporting_identifier(value, maximum=128) ReportingCanonicalDigest("0" * 64, "dialect", self.definition.schema_dialect, "0" * 64) for units in ( self.definition.monetary_metric_units, self.definition.monetary_control_total_units, ): for name, unit in units: reporting_identifier(name, maximum=128) reporting_identifier(unit, maximum=32) object.__setattr__( self, "definition", replace( self.definition, report_definition_sha256=self.definition.report_definition_sha256.lower(), schema_sha256=self.definition.schema_sha256.lower(), ), )
class ReportingVerificationLimits (max_depth: int = 32,
max_value_bytes: int = 1048576,
max_total_bytes: int = 67108864,
max_items: int = 1000000,
max_rows: int = 100000,
max_pages: int = 10000,
max_objects: int = 10000,
max_chunks: int = 1000000)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128")ReportingVerificationLimits(max_depth: 'int' = 32, max_value_bytes: 'int' = 1048576, max_total_bytes: 'int' = 67108864, max_items: 'int' = 1000000, max_rows: 'int' = 100000, max_pages: 'int' = 10000, max_objects: 'int' = 10000, max_chunks: 'int' = 1000000)
Instance variables
var max_chunks : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128") var max_depth : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128") var max_items : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128") var max_objects : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128") var max_pages : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128") var max_rows : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128") var max_total_bytes : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128") var max_value_bytes : int-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerificationLimits: max_depth: int = 32 max_value_bytes: int = 1024 * 1024 max_total_bytes: int = 64 * 1024 * 1024 max_items: int = 1_000_000 max_rows: int = 100_000 max_pages: int = 10_000 max_objects: int = 10_000 max_chunks: int = 1_000_000 def __post_init__(self) -> None: if any( type(getattr(self, f.name)) is not int or getattr(self, f.name) < 1 for f in fields(self) ): raise ValueError("verification limits require positive integer bounds") if self.max_depth > 128: raise ValueError("verification depth cannot exceed 128")
class ReportingVerifiedDestination (request: ReportingDestinationRequest,
resource: ReportingResourceRecord,
verification: ReportingVerificationRecord)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None)Immutable SDK observations. B2 alone owns atomic publication/readiness.
Ancestors
- adcp.reporting.materializer.verification._VerifiedProvenance
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var request : ReportingDestinationRequest-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None) var resource : ReportingResourceRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None) var verification : ReportingVerificationRecord-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingVerifiedDestination(_VerifiedProvenance): """Immutable SDK observations. B2 alone owns atomic publication/readiness.""" request: ReportingDestinationRequest resource: ReportingResourceRecord verification: ReportingVerificationRecord def __post_init__(self) -> None: _freeze_fields(self) object.__setattr__(self, "_sdk_seal", b"") object.__setattr__(self, "_materializer_fence", None)
class ReportingWriterCapability (method: DeliveryMethod,
transport: str,
format: ReportingFormat | None,
verification_profile: VerificationProfile,
verification_path: VerificationPath,
immutability: "Literal['immutable_location', 'native_version']",
checksum: "Literal['sha256']",
write_semantics: "Literal['conditional_create', 'idempotent']")-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path")ReportingWriterCapability(method: 'DeliveryMethod', transport: 'str', format: 'ReportingFormat | None', verification_profile: 'VerificationProfile', verification_path: 'VerificationPath', immutability: "Literal['immutable_location', 'native_version']", checksum: "Literal['sha256']", write_semantics: "Literal['conditional_create', 'idempotent']")
Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var checksum : Literal['sha256']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path") var format : Literal['jsonl', 'csv', 'parquet', 'avro', 'orc'] | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path") var immutability : Literal['immutable_location', 'native_version']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path") var method : Literal['file_transfer', 'dataset_share', 'warehouse_materialization']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path") var transport : str-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path") var verification_path : Literal['producer', 'representative_consumer', 'destination']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path") var verification_profile : Literal['canonical_digest', 'manifest_checksums', 'native_commit']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path") var write_semantics : Literal['conditional_create', 'idempotent']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterCapability(_ClosedValue): method: DeliveryMethod transport: str format: ReportingFormat | None verification_profile: VerificationProfile verification_path: VerificationPath immutability: Literal["immutable_location", "native_version"] checksum: Literal["sha256"] write_semantics: Literal["conditional_create", "idempotent"] def __post_init__(self) -> None: _freeze_fields(self) reporting_identifier(self.transport, maximum=64) if re.fullmatch(r"[a-z][a-z0-9_.-]{0,63}", self.transport) is None: raise ValueError("transport requires a public protocol label") if self.method == "file_transfer" and self.format is None: raise ValueError("file transfer requires a format") if self.verification_profile == "manifest_checksums" and self.method != "file_transfer": raise ValueError("manifest verification requires file transfer") if ( (self.method == "dataset_share" and self.verification_path != "representative_consumer") or ( self.method == "warehouse_materialization" and self.verification_path != "destination" ) or ( self.verification_profile == "native_commit" and self.immutability != "native_version" ) or (self.immutability == "native_version" and self.verification_path == "producer") ): raise ValueError("capability requires its exact immutable observation path")
class ReportingWriterError (failure: ReportingWriterFailure)-
Expand source code
class ReportingWriterError(Exception): """A closed failure; never pass provider prose or attach a provider cause.""" def __init__(self, failure: ReportingWriterFailure) -> None: if type(failure) is not ReportingWriterFailure: raise TypeError("writer errors require a closed failure") self.failure = failure super().__init__(failure.code) def __repr__(self) -> str: return f"ReportingWriterError({self.failure.code})"A closed failure; never pass provider prose or attach a provider cause.
Ancestors
- builtins.Exception
- builtins.BaseException
Subclasses
- adcp.reporting.production.contracts._SourceAuthorizationRevokedError
class ReportingWriterFailure (code: ReportingWriterFailureCode,
retry: ReportingWriterRetry = 'never',
effect: ReportingExternalEffect = 'not_started',
retry_after_seconds: int | None = None)-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterFailure(_ClosedValue): """Safe diagnostics only. Unknown effects must retain the original identity. ``new_attempt`` is usable by B2 only for its own known terminal failure. Public ledger persistence still allows N+1 after any immutable outcome. Consumer receipt rejection is outside this contract. """ code: ReportingWriterFailureCode retry: ReportingWriterRetry = "never" effect: ReportingExternalEffect = "not_started" retry_after_seconds: int | None = None def __post_init__(self) -> None: _freeze_fields(self) if self.effect == "unknown" and self.retry == "new_attempt": raise ValueError("unknown external effects require the original identity") if self.retry_after_seconds is not None and self.retry_after_seconds < 0: raise ValueError("retry delay must be nonnegative")Safe diagnostics only. Unknown effects must retain the original identity.
new_attemptis usable by B2 only for its own known terminal failure. Public ledger persistence still allows N+1 after any immutable outcome. Consumer receipt rejection is outside this contract.Ancestors
- adcp.reporting.ledger.delivery_models._ClosedValue
Instance variables
var code : Literal['UNSUPPORTED_VERIFICATION', 'HISTORY_CORRUPT', 'REVISION_NOT_READY', 'CURRENT_REVISION_CHANGED', 'AUTHORIZATION_DENIED', 'BINDING_MISMATCH', 'SOURCE_INVALID', 'DESTINATION_CORRUPT', 'RESOURCE_UNAVAILABLE', 'WRITE_FAILED', 'DEADLINE_EXCEEDED', 'LEASE_LOST', 'LIMIT_EXCEEDED']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterFailure(_ClosedValue): """Safe diagnostics only. Unknown effects must retain the original identity. ``new_attempt`` is usable by B2 only for its own known terminal failure. Public ledger persistence still allows N+1 after any immutable outcome. Consumer receipt rejection is outside this contract. """ code: ReportingWriterFailureCode retry: ReportingWriterRetry = "never" effect: ReportingExternalEffect = "not_started" retry_after_seconds: int | None = None def __post_init__(self) -> None: _freeze_fields(self) if self.effect == "unknown" and self.retry == "new_attempt": raise ValueError("unknown external effects require the original identity") if self.retry_after_seconds is not None and self.retry_after_seconds < 0: raise ValueError("retry delay must be nonnegative") var effect : Literal['not_started', 'applied', 'unknown']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterFailure(_ClosedValue): """Safe diagnostics only. Unknown effects must retain the original identity. ``new_attempt`` is usable by B2 only for its own known terminal failure. Public ledger persistence still allows N+1 after any immutable outcome. Consumer receipt rejection is outside this contract. """ code: ReportingWriterFailureCode retry: ReportingWriterRetry = "never" effect: ReportingExternalEffect = "not_started" retry_after_seconds: int | None = None def __post_init__(self) -> None: _freeze_fields(self) if self.effect == "unknown" and self.retry == "new_attempt": raise ValueError("unknown external effects require the original identity") if self.retry_after_seconds is not None and self.retry_after_seconds < 0: raise ValueError("retry delay must be nonnegative") var retry : Literal['never', 'same_identity', 'new_attempt']-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterFailure(_ClosedValue): """Safe diagnostics only. Unknown effects must retain the original identity. ``new_attempt`` is usable by B2 only for its own known terminal failure. Public ledger persistence still allows N+1 after any immutable outcome. Consumer receipt rejection is outside this contract. """ code: ReportingWriterFailureCode retry: ReportingWriterRetry = "never" effect: ReportingExternalEffect = "not_started" retry_after_seconds: int | None = None def __post_init__(self) -> None: _freeze_fields(self) if self.effect == "unknown" and self.retry == "new_attempt": raise ValueError("unknown external effects require the original identity") if self.retry_after_seconds is not None and self.retry_after_seconds < 0: raise ValueError("retry delay must be nonnegative") var retry_after_seconds : int | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingWriterFailure(_ClosedValue): """Safe diagnostics only. Unknown effects must retain the original identity. ``new_attempt`` is usable by B2 only for its own known terminal failure. Public ledger persistence still allows N+1 after any immutable outcome. Consumer receipt rejection is outside this contract. """ code: ReportingWriterFailureCode retry: ReportingWriterRetry = "never" effect: ReportingExternalEffect = "not_started" retry_after_seconds: int | None = None def __post_init__(self) -> None: _freeze_fields(self) if self.effect == "unknown" and self.retry == "new_attempt": raise ValueError("unknown external effects require the original identity") if self.retry_after_seconds is not None and self.retry_after_seconds < 0: raise ValueError("retry delay must be nonnegative")