Module adcp.reporting.ledger.delivery
Optional durable seller contracts for future destination writers and receipt handlers.
These stores persist evidence only. They do not resolve credentials, perform external writes, mount tasks, or advertise Managed Delivery/Reconciled Billing. The existing Core store protocol and construction contract remain unchanged.
Functions
def adjustment_to_wire(adjustment: ReportingAdjustmentRecord) ‑> dict[str, typing.Any]-
Expand source code
def adjustment_to_wire(adjustment: ReportingAdjustmentRecord) -> dict[str, Any]: """Complete immutable adjustment evidence, including reason_detail in its digest.""" return { **adjustment_payload(adjustment), "canonical_adjustment_sha256": adjustment_sha256(adjustment), }Complete immutable adjustment evidence, including reason_detail in its digest.
def materialization_to_wire(view: ReportingMaterializationView,
*,
obligation: ReportingObligationRecord) ‑> dict[str, typing.Any]-
Expand source code
def materialization_to_wire( view: ReportingMaterializationView, *, obligation: ReportingObligationRecord ) -> dict[str, Any]: if view.attempt.scope.generation_key != obligation.generation_key or ( view.attempt.scope.reporting_obligation_id != obligation.reporting_obligation_id ): unavailable() result = view.to_wire() result["feed_purpose"] = obligation.feed_purpose return result def receipt_to_wire(record: ReportingReceiptRecord) ‑> dict[str, typing.Any]-
Expand source code
def receipt_to_wire(record: ReportingReceiptRecord) -> dict[str, Any]: result = payload(record) result.pop("scope") result.pop("kind") if isinstance(record, ReportingRevisionReceiptRecord): result["reporting_obligation_id"] = record.scope.reporting_obligation_id result["observed_control_totals"] = totals_to_wire(record.observed_control_totals) if record.observed_canonical_content_digest is not None: result["observed_canonical_content_digest"] = ( record.observed_canonical_content_digest.to_wire() ) return { key: value for key, value in result.items() if value is not None and not (key == "rejection_codes" and not value) } def revision_to_wire(revision: ReportingRevisionRecord, *, obligation: ReportingObligationRecord) ‑> dict[str, typing.Any]-
Expand source code
def revision_to_wire( revision: ReportingRevisionRecord, *, obligation: ReportingObligationRecord ) -> dict[str, Any]: """Opt-in evidence projection; Core's mounted handler remains unchanged.""" from adcp.reporting.ledger.status import _revision_to_wire if ( revision.account_id != obligation.account_id or revision.reporting_obligation_id != obligation.reporting_obligation_id ): unavailable() result = _revision_to_wire(revision, obligation) result["control_totals"] = totals_to_wire(revision_control_totals(revision, obligation)) if revision.canonical_content_digest is not None: result["canonical_content_digest"] = revision.canonical_content_digest.to_wire() return resultOpt-in evidence projection; Core's mounted handler remains unchanged.
Classes
class InMemoryReportingReconciliationStore (*, clock: Callable[[], datetime] | None = None, notifications: bool = False)-
Expand source code
class InMemoryReportingReconciliationStore(InMemoryReportingLedgerStore, _ReconciliationOperations): """Optional reference extension. No destination/receipt services at construction.""" _delivery_records: list[tuple[int, ReportingDeliveryPrincipal, ReportingDeliveryRecord]] async def _commit(self, record: RecordT) -> tuple[RecordT, bool]: candidate = cast(RecordT, decode_record(payload(record))) async with self._mutation(): return self._commit_record_unlocked(candidate) def _commit_record_unlocked( self, record: RecordT, *, notify: bool = True, dirty: bool = True ) -> tuple[RecordT, bool]: """Caller owns the memory mutation; no lock or callback is acquired here. ``notify`` and ``dirty`` are independent: suppressing a readiness event must never also drop the projection work that an ordinary public write has always produced. """ candidate = decode_record(payload(record)) who = principal(candidate) records = tuple(item.record for item in self._caller_changes(who)) existing = replay(candidate, records) if existing is not None: return cast(RecordT, existing), False context = self._delivery_context(candidate) stored = validate_transition(candidate, records, context, self._clock()) self._append_reconciliation_change(stored) if (notify or dirty) and self._notification_state is not None: from adcp.reporting.ledger.notification_events import ( delivery_dirty, materialization_event, ) if notify: event = materialization_event( stored, records, context.obligation, context.revision, context.configuration, self._clock(), ) if event is not None: self._record_notification(event) if dirty: scope, reason, evidence = delivery_dirty(stored, context.obligation) self._dirty_status(scope, reason, after=evidence) return cast(RecordT, stored), True def _append_reconciliation_change(self, record: ReportingDeliveryRecord) -> None: who = principal(record) retained = self._retained_delivery_records() sequence = sum(owner == who for _, owner, _ in retained) + 1 # One assignment publishes the record, its feed row, and its local head. # Core's sequence and change list never participate in this transaction. self._delivery_records = [*retained, (sequence, who, record)] def _retained_delivery_records( self, ) -> list[tuple[int, ReportingDeliveryPrincipal, ReportingDeliveryRecord]]: # Lazily allocated so construction keeps Core's exact component surface. if not hasattr(self, "_delivery_records"): self._delivery_records = [] return self._delivery_records def _caller_changes( self, caller: ReportingDeliveryPrincipal ) -> tuple[ReportingReconciliationChange, ...]: retained = [ item for item in self._retained_delivery_records() if item[1] == caller or principal(item[2]) == caller ] if any( type(sequence) is not int or sequence != ordinal or owner != caller or principal(record) != caller for ordinal, (sequence, owner, record) in enumerate(retained, start=1) ): fail("REPORTING_HISTORY_CORRUPT") return tuple( ReportingReconciliationChange(sequence, record) for sequence, _, record in retained ) def _delivery_context(self, record: ReportingDeliveryRecord) -> DeliveryContext: if isinstance(record, ReportingDestinationBinding): return DeliveryContext(configuration=self._configurations.get(record.generation_key)) obligation = self._obligations.get(record.scope.reporting_obligation_id) revision_id = getattr(record, "reporting_revision_id", None) if isinstance(record, ReportingAdjustmentReceiptRecord): revision_id = record.adjusts_reporting_revision_id revision = self._revisions.get(revision_id) if revision_id is not None else None return DeliveryContext( configuration=self._configurations.get(record.scope.generation_key), obligation=obligation, revision=revision, adjustment=( self._adjustments.get(record.reporting_adjustment_id) if isinstance(record, ReportingAdjustmentReceiptRecord) else None ), ) async def read_reconciliation_snapshot( self, *, caller: ReportingDeliveryPrincipal, boundary: ReportingReconciliationSnapshotToken | None = None, ) -> ReportingReconciliationSnapshot: requested = boundary is not None async with self._lock: changes = self._caller_changes(caller) if boundary is None: boundary = change_boundary(caller, len(changes), self._clock()) if ( type(boundary) is not ReportingReconciliationSnapshotToken or boundary.caller != caller or boundary.min_sequence != 0 or boundary.filters != ReportingReconciliationFilter() ): unavailable() validate_boundary(caller, boundary, len(changes)) if boundary.total_count != boundary.max_sequence: # A caller-presented boundary that disagrees with the retained feed # is a stale or hand-built token, never evidence that storage is # corrupt. Only a boundary this store opened can accuse itself. _boundary_unavailable(requested) return ReportingReconciliationSnapshot( caller, boundary, tuple(item.record for item in changes if item.sequence <= boundary.max_sequence), ) async def read_reconciliation_changes( self, *, caller: ReportingDeliveryPrincipal, changes_after: ReportingReconciliationCheckpoint | None = None, cursor: ReportingReconciliationCursor | None = None, limit: int = 100, filters: ReportingReconciliationFilter = ReportingReconciliationFilter(), ) -> ReportingReconciliationPage: after, boundary, last_key = read_position(caller, changes_after, cursor, limit, filters) continued = boundary is not None async with self._lock: records = self._caller_changes(caller) if boundary is None: boundary = change_boundary( caller, len(records), self._clock(), after=after, total_count=sum( item.sequence > after and filters.matches(item.record) for item in records ), filters=filters, ) validate_boundary(caller, boundary, len(records)) if last_key is not None and not any( item.sequence == after and change_id(item.record) == last_key and filters.matches(item.record) for item in records ): raise LedgerConflictError("INVALID_CHECKPOINT", "reconciliation key is unavailable") changes = tuple( item for item in records if boundary.min_sequence < item.sequence <= boundary.max_sequence and filters.matches(item.record) ) if len(changes) != boundary.total_count: _boundary_unavailable(continued) return change_page( caller, boundary, after, tuple(item for item in changes if item.sequence > after), limit, )Optional reference extension. No destination/receipt services at construction.
Ancestors
- InMemoryReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
Subclasses
Methods
async def read_reconciliation_changes(self,
*,
caller: ReportingDeliveryPrincipal,
changes_after: ReportingReconciliationCheckpoint | None = None,
cursor: ReportingReconciliationCursor | None = None,
limit: int = 100,
filters: ReportingReconciliationFilter = ReportingReconciliationFilter(record_kinds=(), reporting_obligation_id=None)) ‑> ReportingReconciliationPage-
Expand source code
async def read_reconciliation_changes( self, *, caller: ReportingDeliveryPrincipal, changes_after: ReportingReconciliationCheckpoint | None = None, cursor: ReportingReconciliationCursor | None = None, limit: int = 100, filters: ReportingReconciliationFilter = ReportingReconciliationFilter(), ) -> ReportingReconciliationPage: after, boundary, last_key = read_position(caller, changes_after, cursor, limit, filters) continued = boundary is not None async with self._lock: records = self._caller_changes(caller) if boundary is None: boundary = change_boundary( caller, len(records), self._clock(), after=after, total_count=sum( item.sequence > after and filters.matches(item.record) for item in records ), filters=filters, ) validate_boundary(caller, boundary, len(records)) if last_key is not None and not any( item.sequence == after and change_id(item.record) == last_key and filters.matches(item.record) for item in records ): raise LedgerConflictError("INVALID_CHECKPOINT", "reconciliation key is unavailable") changes = tuple( item for item in records if boundary.min_sequence < item.sequence <= boundary.max_sequence and filters.matches(item.record) ) if len(changes) != boundary.total_count: _boundary_unavailable(continued) return change_page( caller, boundary, after, tuple(item for item in changes if item.sequence > after), limit, ) async def read_reconciliation_snapshot(self,
*,
caller: ReportingDeliveryPrincipal,
boundary: ReportingReconciliationSnapshotToken | None = None) ‑> ReportingReconciliationSnapshot-
Expand source code
async def read_reconciliation_snapshot( self, *, caller: ReportingDeliveryPrincipal, boundary: ReportingReconciliationSnapshotToken | None = None, ) -> ReportingReconciliationSnapshot: requested = boundary is not None async with self._lock: changes = self._caller_changes(caller) if boundary is None: boundary = change_boundary(caller, len(changes), self._clock()) if ( type(boundary) is not ReportingReconciliationSnapshotToken or boundary.caller != caller or boundary.min_sequence != 0 or boundary.filters != ReportingReconciliationFilter() ): unavailable() validate_boundary(caller, boundary, len(changes)) if boundary.total_count != boundary.max_sequence: # A caller-presented boundary that disagrees with the retained feed # is a stale or hand-built token, never evidence that storage is # corrupt. Only a boundary this store opened can accuse itself. _boundary_unavailable(requested) return ReportingReconciliationSnapshot( caller, boundary, tuple(item.record for item in changes if item.sequence <= boundary.max_sequence), )
Inherited members
class ReportingDestinationStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingDestinationStore(Protocol): """Trusted configuration ingestion. Never accept a destination from a receipt body.""" async def put_destination_binding( self, record: ReportingDestinationBinding ) -> tuple[ReportingDestinationBinding, bool]: ... async def bind_obligation_delivery( self, record: ReportingObligationDeliveryRecord ) -> tuple[ReportingObligationDeliveryRecord, bool]: ... async def get_destination_binding( self, *, caller: ReportingDeliveryPrincipal, generation_key: ReportingConfigurationGenerationKey, ) -> ReportingDestinationBinding | None: ... async def get_obligation_delivery( self, scope: ReportingDeliveryScope ) -> ReportingObligationDeliveryRecord | None: ...Trusted configuration ingestion. Never accept a destination from a receipt body.
Ancestors
- typing.Protocol
- typing.Generic
Subclasses
Methods
async def bind_obligation_delivery(self, record: ReportingObligationDeliveryRecord) ‑> tuple[ReportingObligationDeliveryRecord, bool]-
Expand source code
async def bind_obligation_delivery( self, record: ReportingObligationDeliveryRecord ) -> tuple[ReportingObligationDeliveryRecord, bool]: ... async def get_destination_binding(self,
*,
caller: ReportingDeliveryPrincipal,
generation_key: ReportingConfigurationGenerationKey) ‑> ReportingDestinationBinding | None-
Expand source code
async def get_destination_binding( self, *, caller: ReportingDeliveryPrincipal, generation_key: ReportingConfigurationGenerationKey, ) -> ReportingDestinationBinding | None: ... async def get_obligation_delivery(self, scope: ReportingDeliveryScope) ‑> ReportingObligationDeliveryRecord | None-
Expand source code
async def get_obligation_delivery( self, scope: ReportingDeliveryScope ) -> ReportingObligationDeliveryRecord | None: ... async def put_destination_binding(self, record: ReportingDestinationBinding) ‑> tuple[ReportingDestinationBinding, bool]-
Expand source code
async def put_destination_binding( self, record: ReportingDestinationBinding ) -> tuple[ReportingDestinationBinding, bool]: ...
class ReportingMaterializationStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingMaterializationStore(Protocol): async def commit_materialization_attempt( self, record: ReportingMaterializationAttempt ) -> tuple[ReportingMaterializationAttempt, bool]: ... async def commit_materialization( self, record: ReportingMaterializationRecord ) -> tuple[ReportingMaterializationRecord, bool]: ... async def record_materialization_check( self, record: ReportingMaterializationCheck ) -> tuple[ReportingMaterializationCheck, bool]: ... async def get_materialization( self, key: ReportingMaterializationKey ) -> ReportingMaterializationView | None: ...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 commit_materialization(self, record: ReportingMaterializationRecord) ‑> tuple[ReportingMaterializationRecord, bool]-
Expand source code
async def commit_materialization( self, record: ReportingMaterializationRecord ) -> tuple[ReportingMaterializationRecord, bool]: ... async def commit_materialization_attempt(self, record: ReportingMaterializationAttempt) ‑> tuple[ReportingMaterializationAttempt, bool]-
Expand source code
async def commit_materialization_attempt( self, record: ReportingMaterializationAttempt ) -> tuple[ReportingMaterializationAttempt, bool]: ... async def get_materialization(self, key: ReportingMaterializationKey) ‑> ReportingMaterializationView | None-
Expand source code
async def get_materialization( self, key: ReportingMaterializationKey ) -> ReportingMaterializationView | None: ... async def record_materialization_check(self, record: ReportingMaterializationCheck) ‑> tuple[ReportingMaterializationCheck, bool]-
Expand source code
async def record_materialization_check( self, record: ReportingMaterializationCheck ) -> tuple[ReportingMaterializationCheck, bool]: ...
class ReportingMaterializationView (attempt: ReportingMaterializationAttempt,
binding: ReportingDestinationBinding,
outcome: ReportingMaterializationRecord | None,
checks: tuple[ReportingMaterializationCheck, ...])-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationView: attempt: ReportingMaterializationAttempt binding: ReportingDestinationBinding outcome: ReportingMaterializationRecord | None checks: tuple[ReportingMaterializationCheck, ...] @property def check(self) -> ReportingMaterializationCheck | None: return max(self.checks, key=lambda item: item.checked_at, default=None) def readable_at(self, at: datetime) -> bool: outcome = self.outcome check = max( (item for item in self.checks if item.checked_at <= at), key=lambda item: item.checked_at, default=None, ) return bool( outcome is not None and outcome.status in {"available", "delivered"} and outcome.resource is not None and outcome.completed_at <= at < outcome.resource.expires_at and (check is None or check.state == "readable") ) def to_wire(self) -> dict[str, Any]: """An inert projection for future handlers. Never exposes the trusted reference.""" attempt, binding, outcome = self.attempt, self.binding, self.outcome generation = attempt.scope.generation_key result: dict[str, Any] = { "reporting_materialization_id": attempt.reporting_materialization_id, "reporting_revision_id": attempt.reporting_revision_id, "reporting_obligation_id": attempt.scope.reporting_obligation_id, "delivery_config_id": generation.delivery_config_id, "delivery_config_version": generation.delivery_config_version, "destination_ref": binding.destination_ref, "feed_purpose": binding.feed_purpose, "method": binding.method, "transport": binding.transport, "attempt": attempt.attempt, "status": outcome.status if outcome is not None else "pending", "created_at": iso(attempt.created_at), } if outcome is None: return result if outcome.status == "failed": result.update(failed_at=iso(outcome.completed_at), failure_code=outcome.failure_code) return result resource, verification = outcome.resource, outcome.verification assert resource is not None and verification is not None descriptor = { key: value for key, value in asdict(resource).items() if value is not None and key != "object_refs" } descriptor["expires_at"] = iso(resource.expires_at) descriptor["reader_compatibility"] = list(resource.reader_compatibility) if resource.kind == "manifest": descriptor["manifest_version"] = "1.0" evidence: dict[str, Any] = { "verified_at": iso(verification.verified_at), "verification_path": verification.verification_path, "verification_profile": verification.verification_profile, "row_count": verification.row_count, "control_totals": totals_to_wire(verification.control_totals), } if verification.canonical_content_digest is not None: evidence["canonical_content_digest"] = verification.canonical_content_digest.to_wire() if verification.physical_checksums: evidence["physical_checksums"] = [ asdict(item) for item in verification.physical_checksums ] if verification.native_version_ref is not None: evidence["native_commit_evidence"] = { "native_version_ref": verification.native_version_ref, "observed_through": verification.native_observed_through, } result.update( ready_at=iso(outcome.completed_at), resource=descriptor, verification=evidence ) return resultReportingMaterializationView(attempt: 'ReportingMaterializationAttempt', binding: 'ReportingDestinationBinding', outcome: 'ReportingMaterializationRecord | None', checks: 'tuple[ReportingMaterializationCheck, …]')
Instance variables
var attempt : ReportingMaterializationAttempt-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationView: attempt: ReportingMaterializationAttempt binding: ReportingDestinationBinding outcome: ReportingMaterializationRecord | None checks: tuple[ReportingMaterializationCheck, ...] @property def check(self) -> ReportingMaterializationCheck | None: return max(self.checks, key=lambda item: item.checked_at, default=None) def readable_at(self, at: datetime) -> bool: outcome = self.outcome check = max( (item for item in self.checks if item.checked_at <= at), key=lambda item: item.checked_at, default=None, ) return bool( outcome is not None and outcome.status in {"available", "delivered"} and outcome.resource is not None and outcome.completed_at <= at < outcome.resource.expires_at and (check is None or check.state == "readable") ) def to_wire(self) -> dict[str, Any]: """An inert projection for future handlers. Never exposes the trusted reference.""" attempt, binding, outcome = self.attempt, self.binding, self.outcome generation = attempt.scope.generation_key result: dict[str, Any] = { "reporting_materialization_id": attempt.reporting_materialization_id, "reporting_revision_id": attempt.reporting_revision_id, "reporting_obligation_id": attempt.scope.reporting_obligation_id, "delivery_config_id": generation.delivery_config_id, "delivery_config_version": generation.delivery_config_version, "destination_ref": binding.destination_ref, "feed_purpose": binding.feed_purpose, "method": binding.method, "transport": binding.transport, "attempt": attempt.attempt, "status": outcome.status if outcome is not None else "pending", "created_at": iso(attempt.created_at), } if outcome is None: return result if outcome.status == "failed": result.update(failed_at=iso(outcome.completed_at), failure_code=outcome.failure_code) return result resource, verification = outcome.resource, outcome.verification assert resource is not None and verification is not None descriptor = { key: value for key, value in asdict(resource).items() if value is not None and key != "object_refs" } descriptor["expires_at"] = iso(resource.expires_at) descriptor["reader_compatibility"] = list(resource.reader_compatibility) if resource.kind == "manifest": descriptor["manifest_version"] = "1.0" evidence: dict[str, Any] = { "verified_at": iso(verification.verified_at), "verification_path": verification.verification_path, "verification_profile": verification.verification_profile, "row_count": verification.row_count, "control_totals": totals_to_wire(verification.control_totals), } if verification.canonical_content_digest is not None: evidence["canonical_content_digest"] = verification.canonical_content_digest.to_wire() if verification.physical_checksums: evidence["physical_checksums"] = [ asdict(item) for item in verification.physical_checksums ] if verification.native_version_ref is not None: evidence["native_commit_evidence"] = { "native_version_ref": verification.native_version_ref, "observed_through": verification.native_observed_through, } result.update( ready_at=iso(outcome.completed_at), resource=descriptor, verification=evidence ) return result var binding : ReportingDestinationBinding-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationView: attempt: ReportingMaterializationAttempt binding: ReportingDestinationBinding outcome: ReportingMaterializationRecord | None checks: tuple[ReportingMaterializationCheck, ...] @property def check(self) -> ReportingMaterializationCheck | None: return max(self.checks, key=lambda item: item.checked_at, default=None) def readable_at(self, at: datetime) -> bool: outcome = self.outcome check = max( (item for item in self.checks if item.checked_at <= at), key=lambda item: item.checked_at, default=None, ) return bool( outcome is not None and outcome.status in {"available", "delivered"} and outcome.resource is not None and outcome.completed_at <= at < outcome.resource.expires_at and (check is None or check.state == "readable") ) def to_wire(self) -> dict[str, Any]: """An inert projection for future handlers. Never exposes the trusted reference.""" attempt, binding, outcome = self.attempt, self.binding, self.outcome generation = attempt.scope.generation_key result: dict[str, Any] = { "reporting_materialization_id": attempt.reporting_materialization_id, "reporting_revision_id": attempt.reporting_revision_id, "reporting_obligation_id": attempt.scope.reporting_obligation_id, "delivery_config_id": generation.delivery_config_id, "delivery_config_version": generation.delivery_config_version, "destination_ref": binding.destination_ref, "feed_purpose": binding.feed_purpose, "method": binding.method, "transport": binding.transport, "attempt": attempt.attempt, "status": outcome.status if outcome is not None else "pending", "created_at": iso(attempt.created_at), } if outcome is None: return result if outcome.status == "failed": result.update(failed_at=iso(outcome.completed_at), failure_code=outcome.failure_code) return result resource, verification = outcome.resource, outcome.verification assert resource is not None and verification is not None descriptor = { key: value for key, value in asdict(resource).items() if value is not None and key != "object_refs" } descriptor["expires_at"] = iso(resource.expires_at) descriptor["reader_compatibility"] = list(resource.reader_compatibility) if resource.kind == "manifest": descriptor["manifest_version"] = "1.0" evidence: dict[str, Any] = { "verified_at": iso(verification.verified_at), "verification_path": verification.verification_path, "verification_profile": verification.verification_profile, "row_count": verification.row_count, "control_totals": totals_to_wire(verification.control_totals), } if verification.canonical_content_digest is not None: evidence["canonical_content_digest"] = verification.canonical_content_digest.to_wire() if verification.physical_checksums: evidence["physical_checksums"] = [ asdict(item) for item in verification.physical_checksums ] if verification.native_version_ref is not None: evidence["native_commit_evidence"] = { "native_version_ref": verification.native_version_ref, "observed_through": verification.native_observed_through, } result.update( ready_at=iso(outcome.completed_at), resource=descriptor, verification=evidence ) return result prop check : ReportingMaterializationCheck | None-
Expand source code
@property def check(self) -> ReportingMaterializationCheck | None: return max(self.checks, key=lambda item: item.checked_at, default=None) var checks : tuple[ReportingMaterializationCheck, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationView: attempt: ReportingMaterializationAttempt binding: ReportingDestinationBinding outcome: ReportingMaterializationRecord | None checks: tuple[ReportingMaterializationCheck, ...] @property def check(self) -> ReportingMaterializationCheck | None: return max(self.checks, key=lambda item: item.checked_at, default=None) def readable_at(self, at: datetime) -> bool: outcome = self.outcome check = max( (item for item in self.checks if item.checked_at <= at), key=lambda item: item.checked_at, default=None, ) return bool( outcome is not None and outcome.status in {"available", "delivered"} and outcome.resource is not None and outcome.completed_at <= at < outcome.resource.expires_at and (check is None or check.state == "readable") ) def to_wire(self) -> dict[str, Any]: """An inert projection for future handlers. Never exposes the trusted reference.""" attempt, binding, outcome = self.attempt, self.binding, self.outcome generation = attempt.scope.generation_key result: dict[str, Any] = { "reporting_materialization_id": attempt.reporting_materialization_id, "reporting_revision_id": attempt.reporting_revision_id, "reporting_obligation_id": attempt.scope.reporting_obligation_id, "delivery_config_id": generation.delivery_config_id, "delivery_config_version": generation.delivery_config_version, "destination_ref": binding.destination_ref, "feed_purpose": binding.feed_purpose, "method": binding.method, "transport": binding.transport, "attempt": attempt.attempt, "status": outcome.status if outcome is not None else "pending", "created_at": iso(attempt.created_at), } if outcome is None: return result if outcome.status == "failed": result.update(failed_at=iso(outcome.completed_at), failure_code=outcome.failure_code) return result resource, verification = outcome.resource, outcome.verification assert resource is not None and verification is not None descriptor = { key: value for key, value in asdict(resource).items() if value is not None and key != "object_refs" } descriptor["expires_at"] = iso(resource.expires_at) descriptor["reader_compatibility"] = list(resource.reader_compatibility) if resource.kind == "manifest": descriptor["manifest_version"] = "1.0" evidence: dict[str, Any] = { "verified_at": iso(verification.verified_at), "verification_path": verification.verification_path, "verification_profile": verification.verification_profile, "row_count": verification.row_count, "control_totals": totals_to_wire(verification.control_totals), } if verification.canonical_content_digest is not None: evidence["canonical_content_digest"] = verification.canonical_content_digest.to_wire() if verification.physical_checksums: evidence["physical_checksums"] = [ asdict(item) for item in verification.physical_checksums ] if verification.native_version_ref is not None: evidence["native_commit_evidence"] = { "native_version_ref": verification.native_version_ref, "observed_through": verification.native_observed_through, } result.update( ready_at=iso(outcome.completed_at), resource=descriptor, verification=evidence ) return result var outcome : ReportingMaterializationRecord | None-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingMaterializationView: attempt: ReportingMaterializationAttempt binding: ReportingDestinationBinding outcome: ReportingMaterializationRecord | None checks: tuple[ReportingMaterializationCheck, ...] @property def check(self) -> ReportingMaterializationCheck | None: return max(self.checks, key=lambda item: item.checked_at, default=None) def readable_at(self, at: datetime) -> bool: outcome = self.outcome check = max( (item for item in self.checks if item.checked_at <= at), key=lambda item: item.checked_at, default=None, ) return bool( outcome is not None and outcome.status in {"available", "delivered"} and outcome.resource is not None and outcome.completed_at <= at < outcome.resource.expires_at and (check is None or check.state == "readable") ) def to_wire(self) -> dict[str, Any]: """An inert projection for future handlers. Never exposes the trusted reference.""" attempt, binding, outcome = self.attempt, self.binding, self.outcome generation = attempt.scope.generation_key result: dict[str, Any] = { "reporting_materialization_id": attempt.reporting_materialization_id, "reporting_revision_id": attempt.reporting_revision_id, "reporting_obligation_id": attempt.scope.reporting_obligation_id, "delivery_config_id": generation.delivery_config_id, "delivery_config_version": generation.delivery_config_version, "destination_ref": binding.destination_ref, "feed_purpose": binding.feed_purpose, "method": binding.method, "transport": binding.transport, "attempt": attempt.attempt, "status": outcome.status if outcome is not None else "pending", "created_at": iso(attempt.created_at), } if outcome is None: return result if outcome.status == "failed": result.update(failed_at=iso(outcome.completed_at), failure_code=outcome.failure_code) return result resource, verification = outcome.resource, outcome.verification assert resource is not None and verification is not None descriptor = { key: value for key, value in asdict(resource).items() if value is not None and key != "object_refs" } descriptor["expires_at"] = iso(resource.expires_at) descriptor["reader_compatibility"] = list(resource.reader_compatibility) if resource.kind == "manifest": descriptor["manifest_version"] = "1.0" evidence: dict[str, Any] = { "verified_at": iso(verification.verified_at), "verification_path": verification.verification_path, "verification_profile": verification.verification_profile, "row_count": verification.row_count, "control_totals": totals_to_wire(verification.control_totals), } if verification.canonical_content_digest is not None: evidence["canonical_content_digest"] = verification.canonical_content_digest.to_wire() if verification.physical_checksums: evidence["physical_checksums"] = [ asdict(item) for item in verification.physical_checksums ] if verification.native_version_ref is not None: evidence["native_commit_evidence"] = { "native_version_ref": verification.native_version_ref, "observed_through": verification.native_observed_through, } result.update( ready_at=iso(outcome.completed_at), resource=descriptor, verification=evidence ) return result
Methods
def readable_at(self, at: datetime) ‑> bool-
Expand source code
def readable_at(self, at: datetime) -> bool: outcome = self.outcome check = max( (item for item in self.checks if item.checked_at <= at), key=lambda item: item.checked_at, default=None, ) return bool( outcome is not None and outcome.status in {"available", "delivered"} and outcome.resource is not None and outcome.completed_at <= at < outcome.resource.expires_at and (check is None or check.state == "readable") ) def to_wire(self) ‑> dict[str, typing.Any]-
Expand source code
def to_wire(self) -> dict[str, Any]: """An inert projection for future handlers. Never exposes the trusted reference.""" attempt, binding, outcome = self.attempt, self.binding, self.outcome generation = attempt.scope.generation_key result: dict[str, Any] = { "reporting_materialization_id": attempt.reporting_materialization_id, "reporting_revision_id": attempt.reporting_revision_id, "reporting_obligation_id": attempt.scope.reporting_obligation_id, "delivery_config_id": generation.delivery_config_id, "delivery_config_version": generation.delivery_config_version, "destination_ref": binding.destination_ref, "feed_purpose": binding.feed_purpose, "method": binding.method, "transport": binding.transport, "attempt": attempt.attempt, "status": outcome.status if outcome is not None else "pending", "created_at": iso(attempt.created_at), } if outcome is None: return result if outcome.status == "failed": result.update(failed_at=iso(outcome.completed_at), failure_code=outcome.failure_code) return result resource, verification = outcome.resource, outcome.verification assert resource is not None and verification is not None descriptor = { key: value for key, value in asdict(resource).items() if value is not None and key != "object_refs" } descriptor["expires_at"] = iso(resource.expires_at) descriptor["reader_compatibility"] = list(resource.reader_compatibility) if resource.kind == "manifest": descriptor["manifest_version"] = "1.0" evidence: dict[str, Any] = { "verified_at": iso(verification.verified_at), "verification_path": verification.verification_path, "verification_profile": verification.verification_profile, "row_count": verification.row_count, "control_totals": totals_to_wire(verification.control_totals), } if verification.canonical_content_digest is not None: evidence["canonical_content_digest"] = verification.canonical_content_digest.to_wire() if verification.physical_checksums: evidence["physical_checksums"] = [ asdict(item) for item in verification.physical_checksums ] if verification.native_version_ref is not None: evidence["native_commit_evidence"] = { "native_version_ref": verification.native_version_ref, "observed_through": verification.native_observed_through, } result.update( ready_at=iso(outcome.completed_at), resource=descriptor, verification=evidence ) return resultAn inert projection for future handlers. Never exposes the trusted reference.
class ReportingReceiptStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingReceiptStore(Protocol): """Transport-derived scope is mandatory. Exact retries precede temporal checks. ``received_at`` is assigned by the store; callers retry the immutable input or the returned record. Both receipt kinds share the same scoped ID namespace. """ async def record_revision_receipt( self, record: ReportingRevisionReceiptRecord ) -> tuple[ReportingRevisionReceiptRecord, bool]: ... async def record_adjustment_receipt( self, record: ReportingAdjustmentReceiptRecord ) -> tuple[ReportingAdjustmentReceiptRecord, bool]: ... async def get_receipt(self, key: ReportingReceiptKey) -> ReportingReceiptRecord | None: ...Transport-derived scope is mandatory. Exact retries precede temporal checks.
received_atis assigned by the store; callers retry the immutable input or the returned record. Both receipt kinds share the same scoped ID namespace.Ancestors
- typing.Protocol
- typing.Generic
Subclasses
Methods
async def get_receipt(self, key: ReportingReceiptKey) ‑> ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord | None-
Expand source code
async def get_receipt(self, key: ReportingReceiptKey) -> ReportingReceiptRecord | None: ... async def record_adjustment_receipt(self, record: ReportingAdjustmentReceiptRecord) ‑> tuple[ReportingAdjustmentReceiptRecord, bool]-
Expand source code
async def record_adjustment_receipt( self, record: ReportingAdjustmentReceiptRecord ) -> tuple[ReportingAdjustmentReceiptRecord, bool]: ... async def record_revision_receipt(self, record: ReportingRevisionReceiptRecord) ‑> tuple[ReportingRevisionReceiptRecord, bool]-
Expand source code
async def record_revision_receipt( self, record: ReportingRevisionReceiptRecord ) -> tuple[ReportingRevisionReceiptRecord, bool]: ...
class ReportingReconciliationSnapshot (caller: ReportingDeliveryPrincipal,
boundary: ReportingReconciliationSnapshotToken,
records: tuple[ReportingDeliveryRecord, ...])-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshot: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken records: tuple[ReportingDeliveryRecord, ...] def materialization( self, key: ReportingMaterializationKey ) -> ReportingMaterializationView | None: if key.principal != self.caller: return None attempt = next( ( item for item in self.records if isinstance(item, ReportingMaterializationAttempt) and item.key == key ), None, ) if attempt is None: return None binding = next( ( item for item in self.records if isinstance(item, ReportingDestinationBinding) and item.generation_key == attempt.scope.generation_key ), None, ) if binding is None: unavailable() outcome = next( ( item for item in self.records if isinstance(item, ReportingMaterializationRecord) and item.key == key ), None, ) return ReportingMaterializationView( attempt, binding, outcome, tuple( item for item in self.records if isinstance(item, ReportingMaterializationCheck) and item.reporting_materialization_id == key.reporting_materialization_id ), ) @property def current_receipts(self) -> tuple[ReportingReceiptRecord, ...]: receipts = tuple( item for item in self.records if isinstance(item, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)) ) return tuple(item for item in receipts if current_receipt(self.records, item) == item) @property def terminal_acceptances(self) -> tuple[ReportingReceiptKey, ...]: return tuple(item.key for item in self.current_receipts if item.status == "accepted")ReportingReconciliationSnapshot(caller: 'ReportingDeliveryPrincipal', boundary: 'ReportingReconciliationSnapshotToken', records: 'tuple[ReportingDeliveryRecord, …]')
Instance variables
var boundary : ReportingReconciliationSnapshotToken-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshot: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken records: tuple[ReportingDeliveryRecord, ...] def materialization( self, key: ReportingMaterializationKey ) -> ReportingMaterializationView | None: if key.principal != self.caller: return None attempt = next( ( item for item in self.records if isinstance(item, ReportingMaterializationAttempt) and item.key == key ), None, ) if attempt is None: return None binding = next( ( item for item in self.records if isinstance(item, ReportingDestinationBinding) and item.generation_key == attempt.scope.generation_key ), None, ) if binding is None: unavailable() outcome = next( ( item for item in self.records if isinstance(item, ReportingMaterializationRecord) and item.key == key ), None, ) return ReportingMaterializationView( attempt, binding, outcome, tuple( item for item in self.records if isinstance(item, ReportingMaterializationCheck) and item.reporting_materialization_id == key.reporting_materialization_id ), ) @property def current_receipts(self) -> tuple[ReportingReceiptRecord, ...]: receipts = tuple( item for item in self.records if isinstance(item, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)) ) return tuple(item for item in receipts if current_receipt(self.records, item) == item) @property def terminal_acceptances(self) -> tuple[ReportingReceiptKey, ...]: return tuple(item.key for item in self.current_receipts if item.status == "accepted") var caller : ReportingDeliveryPrincipal-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshot: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken records: tuple[ReportingDeliveryRecord, ...] def materialization( self, key: ReportingMaterializationKey ) -> ReportingMaterializationView | None: if key.principal != self.caller: return None attempt = next( ( item for item in self.records if isinstance(item, ReportingMaterializationAttempt) and item.key == key ), None, ) if attempt is None: return None binding = next( ( item for item in self.records if isinstance(item, ReportingDestinationBinding) and item.generation_key == attempt.scope.generation_key ), None, ) if binding is None: unavailable() outcome = next( ( item for item in self.records if isinstance(item, ReportingMaterializationRecord) and item.key == key ), None, ) return ReportingMaterializationView( attempt, binding, outcome, tuple( item for item in self.records if isinstance(item, ReportingMaterializationCheck) and item.reporting_materialization_id == key.reporting_materialization_id ), ) @property def current_receipts(self) -> tuple[ReportingReceiptRecord, ...]: receipts = tuple( item for item in self.records if isinstance(item, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)) ) return tuple(item for item in receipts if current_receipt(self.records, item) == item) @property def terminal_acceptances(self) -> tuple[ReportingReceiptKey, ...]: return tuple(item.key for item in self.current_receipts if item.status == "accepted") prop current_receipts : tuple[ReportingReceiptRecord, ...]-
Expand source code
@property def current_receipts(self) -> tuple[ReportingReceiptRecord, ...]: receipts = tuple( item for item in self.records if isinstance(item, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)) ) return tuple(item for item in receipts if current_receipt(self.records, item) == item) var records : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...]-
Expand source code
@dataclass(frozen=True, slots=True) class ReportingReconciliationSnapshot: caller: ReportingDeliveryPrincipal boundary: ReportingReconciliationSnapshotToken records: tuple[ReportingDeliveryRecord, ...] def materialization( self, key: ReportingMaterializationKey ) -> ReportingMaterializationView | None: if key.principal != self.caller: return None attempt = next( ( item for item in self.records if isinstance(item, ReportingMaterializationAttempt) and item.key == key ), None, ) if attempt is None: return None binding = next( ( item for item in self.records if isinstance(item, ReportingDestinationBinding) and item.generation_key == attempt.scope.generation_key ), None, ) if binding is None: unavailable() outcome = next( ( item for item in self.records if isinstance(item, ReportingMaterializationRecord) and item.key == key ), None, ) return ReportingMaterializationView( attempt, binding, outcome, tuple( item for item in self.records if isinstance(item, ReportingMaterializationCheck) and item.reporting_materialization_id == key.reporting_materialization_id ), ) @property def current_receipts(self) -> tuple[ReportingReceiptRecord, ...]: receipts = tuple( item for item in self.records if isinstance(item, (ReportingRevisionReceiptRecord, ReportingAdjustmentReceiptRecord)) ) return tuple(item for item in receipts if current_receipt(self.records, item) == item) @property def terminal_acceptances(self) -> tuple[ReportingReceiptKey, ...]: return tuple(item.key for item in self.current_receipts if item.status == "accepted") prop terminal_acceptances : tuple[ReportingReceiptKey, ...]-
Expand source code
@property def terminal_acceptances(self) -> tuple[ReportingReceiptKey, ...]: return tuple(item.key for item in self.current_receipts if item.status == "accepted")
Methods
def materialization(self, key: ReportingMaterializationKey) ‑> ReportingMaterializationView | None-
Expand source code
def materialization( self, key: ReportingMaterializationKey ) -> ReportingMaterializationView | None: if key.principal != self.caller: return None attempt = next( ( item for item in self.records if isinstance(item, ReportingMaterializationAttempt) and item.key == key ), None, ) if attempt is None: return None binding = next( ( item for item in self.records if isinstance(item, ReportingDestinationBinding) and item.generation_key == attempt.scope.generation_key ), None, ) if binding is None: unavailable() outcome = next( ( item for item in self.records if isinstance(item, ReportingMaterializationRecord) and item.key == key ), None, ) return ReportingMaterializationView( attempt, binding, outcome, tuple( item for item in self.records if isinstance(item, ReportingMaterializationCheck) and item.reporting_materialization_id == key.reporting_materialization_id ), )
class ReportingReconciliationStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingReconciliationStore( ReportingDestinationStore, ReportingMaterializationStore, ReportingReceiptStore, Protocol ): async def read_reconciliation_snapshot( self, *, caller: ReportingDeliveryPrincipal, boundary: ReportingReconciliationSnapshotToken | None = None, ) -> ReportingReconciliationSnapshot: ...Trusted configuration ingestion. Never accept a destination from a receipt body.
Ancestors
- ReportingDestinationStore
- ReportingMaterializationStore
- ReportingReceiptStore
- typing.Protocol
- typing.Generic
Methods
async def read_reconciliation_snapshot(self,
*,
caller: ReportingDeliveryPrincipal,
boundary: ReportingReconciliationSnapshotToken | None = None) ‑> ReportingReconciliationSnapshot-
Expand source code
async def read_reconciliation_snapshot( self, *, caller: ReportingDeliveryPrincipal, boundary: ReportingReconciliationSnapshotToken | None = None, ) -> ReportingReconciliationSnapshot: ...