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 result

Opt-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

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 check

See PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::

class GenProto(Protocol[T]):
    def meth(self) -> T:
        ...

Ancestors

  • typing.Protocol
  • typing.Generic

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 result

ReportingMaterializationView(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 result

An 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_at is 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

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: ...