Module adcp.reporting.ledger.store

The durable seam under the reporting producer and the status handler.

:class:ReportingLedgerStore is deliberately small. Everything interesting – health, coverage roll-up, issue identity, period derivation – is computed from what the store returns, not stored by it. A store's whole job is: keep immutable records, keep them ordered, hand back a consistent snapshot, and refuse to break the invariants that make the records evidence.

The invariants a conforming store MUST enforce, because they cannot be checked after the fact:

  • Obligation identity is unique per (config generation, period). Two obligations for one logical period would let a seller quietly publish twice and pick a winner.
  • Revisions are immutable. Re-committing an existing reporting_revision_id with different content is an idempotency conflict, not an update.
  • At most one official revision per obligation. It is terminal; a later restatement is an adjustment.
  • Supersession is atomic and exact. A snapshot restatement names the current leaf; naming a stale leaf must fail rather than fork the chain.
  • The change feed is gapless per account. Every immutable write appends in the same transaction as the record it describes.

:class:InMemoryReportingLedgerStore is the reference implementation and the conformance target. It is genuinely usable in tests and single-process pilots, and it is not durable – Reliable Reporting's whole premise is retained evidence, so production wants :class:~adcp.reporting.ledger.pg.PgReportingLedgerStore or your own.

Functions

def check_issue_state_transition(current: str, requested: str) ‑> None
Expand source code
def check_issue_state_transition(current: str, requested: str) -> None:
    """Refuse any transition the spec's forward-only lifecycle forbids.

    Raises :class:`LedgerConflictError` rather than returning a bool so neither
    store can forget to act on the answer.
    """
    if requested == current:
        # Idempotent: re-acknowledging an acknowledged issue is a no-op, which
        # is what makes an operator retry safe.
        return
    if requested not in _ISSUE_STATE_SUCCESSORS.get(current, frozenset()):
        raise LedgerConflictError(
            "ISSUE_STATE_TRANSITION_INVALID",
            "issue_state moves only forward through open, acknowledged, then resolved or "
            f"waived; {current!r} -> {requested!r} is not permitted. A recurrence gets a new "
            "occurrence rather than reviving a retired one",
        )

Refuse any transition the spec's forward-only lifecycle forbids.

Raises :class:LedgerConflictError rather than returning a bool so neither store can forget to act on the answer.

def decode_cursor(cursor: str) ‑> dict[str, typing.Any]
Expand source code
def decode_cursor(cursor: str) -> dict[str, Any]:
    padded = cursor + "=" * (-len(cursor) % 4)

    def unique_pairs(pairs: list[tuple[str, Any]]) -> dict[str, Any]:
        result: dict[str, Any] = {}
        for key, value in pairs:
            if key in result:
                raise ValueError("duplicate cursor member")
            result[key] = value
        return result

    try:
        payload = json.loads(
            base64.urlsafe_b64decode(padded.encode("ascii")), object_pairs_hook=unique_pairs
        )
    except Exception:
        payload = None
    if not isinstance(payload, dict):
        raise LedgerConflictError("INVALID_CURSOR", "the pagination cursor is not readable")
    return payload
def encode_cursor(payload: dict[str, Any]) ‑> str
Expand source code
def encode_cursor(payload: dict[str, Any]) -> str:
    """Opaque, self-describing cursor.

    Base64url over canonical JSON: opaque to the consumer, cheap for the seller
    to validate, and stable enough that the same logical position encodes
    identically on every page.  It is *not* signed -- the store re-validates
    caller, account, and snapshot on use, which is the check that actually
    matters.
    """
    return base64.urlsafe_b64encode(canonical_json_utf8_v1(payload)).decode("ascii").rstrip("=")

Opaque, self-describing cursor.

Base64url over canonical JSON: opaque to the consumer, cheap for the seller to validate, and stable enough that the same logical position encodes identically on every page. It is not signed – the store re-validates caller, account, and snapshot on use, which is the check that actually matters.

def issue_is_retirable(current: str) ‑> bool
Expand source code
def issue_is_retirable(current: str) -> bool:
    """Whether the projection may move this state to ``resolved``.

    ``waived`` is not retirable: it is already out of the projection by
    agreement, and overwriting that readable act with ``resolved`` is an edge
    the forward-only lifecycle forbids. Both stores consult this rather than
    each remembering the rule.
    """
    return "resolved" in _ISSUE_STATE_SUCCESSORS.get(current, frozenset())

Whether the projection may move this state to resolved.

waived is not retirable: it is already out of the projection by agreement, and overwriting that readable act with resolved is an edge the forward-only lifecycle forbids. Both stores consult this rather than each remembering the rule.

def reject_reserved_authoritative_party(configuration: ReportingConfiguration) ‑> None
Expand source code
def reject_reserved_authoritative_party(configuration: ReportingConfiguration) -> None:
    """Refuse ``authoritative_party: consumer`` before the generation is stored.

    AdCP 3.2.0-rc.3 reserves the value for a buyer-deposited billing revision
    task scoped to a later minor. No released minor defines that task, so a
    3.2 seller must reject it with ``UNSUPPORTED_FEATURE`` *before* the
    generation becomes ready and before any obligation exists -- and must not
    silently coerce it to ``seller``.

    Coercing would be the dangerous option: the buyer asked to be the
    authoritative counter for a billing feed and would get a seller-authoritative
    one, with every obligation, revision, and receipt in this ledger quietly
    attributed the wrong way round.
    """
    if configuration.authoritative_party == "consumer":
        raise LedgerConflictError(
            "UNSUPPORTED_FEATURE",
            "authoritative_party 'consumer' is reserved for the buyer-deposited billing "
            "revision task, which no released AdCP minor defines. This seller produces "
            "every revision; omit the field or set it to 'seller'. See "
            "https://github.com/adcontextprotocol/adcp/issues/7440",
        )

Refuse authoritative_party: consumer before the generation is stored.

AdCP 3.2.0-rc.3 reserves the value for a buyer-deposited billing revision task scoped to a later minor. No released minor defines that task, so a 3.2 seller must reject it with UNSUPPORTED_FEATURE before the generation becomes ready and before any obligation exists – and must not silently coerce it to seller.

Coercing would be the dangerous option: the buyer asked to be the authoritative counter for a billing feed and would get a seller-authoritative one, with every obligation, revision, and receipt in this ledger quietly attributed the wrong way round.

Classes

class InMemoryReportingLedgerStore (*, clock: Callable[[], datetime] | None = None, notifications: bool = False)
Expand source code
class InMemoryReportingLedgerStore:
    """Process-local reference store. Correct, ordered, and not durable.

    ``clock`` supplies change timestamps and the snapshot observation boundary,
    standing in for the database clock a durable store reads. Override it to
    place a test's ledger boundary at a deliberate instant.
    """

    def __init__(
        self, *, clock: Callable[[], datetime] | None = None, notifications: bool = False
    ) -> None:
        self._clock = clock or (lambda: datetime.now(timezone.utc))
        self._lock = asyncio.Lock()
        self._configurations: dict[ReportingConfigurationGenerationKey, ReportingConfiguration] = {}
        self._obligations: dict[str, ReportingObligationRecord] = {}
        self._obligation_by_period: dict[
            tuple[ReportingConfigurationGenerationKey, str, str], str
        ] = {}
        self._revisions: dict[str, ReportingRevisionRecord] = {}
        self._revision_identity: dict[str, str] = {}
        self._rows: dict[str, tuple[dict[str, Any], ...]] = {}
        self._restatement_checkpoints: dict[str, RestatementCheckpoint] = {}
        self._retry_schedules: dict[str, RetryScheduleEntry] = {}
        self._provisional_acquisitions: dict[tuple[str, str, int], ProvisionalAcquisition] = {}
        self._provisional_observations: dict[tuple[str, str, int], ProvisionalObservation] = {}
        self._adjustments: dict[str, ReportingAdjustmentRecord] = {}
        self._statuses: dict[tuple[str, str, str], ConsumerStatusRecord] = {}
        self._status_identity: dict[tuple[str, str, str], str] = {}
        self._changes: list[tuple[int, str, LedgerRecordKind, str, datetime]] = []
        self._change_owners: dict[int, str] = {}
        self._sequence = 0
        self._leases: dict[ReportingConfigurationGenerationKey, tuple[str, datetime]] = {}
        # When each generation was last handed to a worker, so releasing a
        # lease sends that generation to the back of the queue instead of
        # letting it win every turn.
        self._lease_turns: dict[ReportingConfigurationGenerationKey, int] = {}
        self._lease_turn = 0
        # Live occurrence per (account, issue_key), plus the retired generation
        # high-water mark so a recurrence never reuses an id.
        self._issues: dict[tuple[str, str], ReportingIssueLifecycle] = {}
        self._issue_generations: dict[tuple[str, str], int] = {}
        self._issue_status_scopes: dict[tuple[str, str], ReportingStatusScope] = {}
        self._notification_state: NotificationState | None = None
        self._status_notification_state: Any = None
        if notifications:
            from adcp.reporting.outbox.memory import NotificationState

            self._notification_state = NotificationState()
            self._notification_state.issue_scopes = self._issue_status_scopes

    @asynccontextmanager
    async def _mutation(self) -> AsyncIterator[None]:
        """Publish domain changes and notifications under one rollback boundary.

        Rollback also applies with notifications disabled. Newly initialized
        collections and sequence heads belong to the transaction too.
        """
        owner = (id(self), asyncio.current_task())
        if _MEMORY_TRANSACTION.get() == owner:
            yield
            return
        async with self._lock:
            before = deepcopy(
                {key: value for key, value in vars(self).items() if key not in {"_lock", "_clock"}}
            )
            token = _MEMORY_TRANSACTION.set(owner)
            try:
                dirty_start = len(self._notification_state.dirty) if self._notification_state else 0
                yield
                if self._notification_state is not None:
                    from adcp.reporting.ledger.status_snapshot import settle_memory_snapshot
                    from adcp.reporting.outbox.status import StatusBoundary

                    dirty = tuple(self._notification_state.dirty[dirty_start:])
                    transaction_id = str(uuid4())
                    for account_id in sorted({d.scope.account_id for d in dirty}):
                        snapshot = settle_memory_snapshot(self, account_id)
                        account_dirty = tuple(d for d in dirty if d.scope.account_id == account_id)
                        self._notification_state.boundaries.append(
                            StatusBoundary(
                                transaction_id,
                                max(d.sequence for d in account_dirty),
                                account_dirty,
                                snapshot,
                            )
                        )
            except BaseException:
                for key in set(vars(self)) - {"_lock", "_clock"} - set(before):
                    del vars(self)[key]
                vars(self).update(before)
                raise
            finally:
                _MEMORY_TRANSACTION.reset(token)

    @asynccontextmanager
    async def transaction(self) -> AsyncIterator[InMemoryReportingLedgerStore]:
        """Group several source mutations into one atomic, replayable boundary."""
        async with self._mutation():
            yield self

    @asynccontextmanager
    async def _source_publication(
        self, account_id: str, *, seals: ReportingSealStore | None = None
    ) -> AsyncIterator[None]:
        # The memory mutation lock serializes all accounts, including seals.
        async with self._mutation():
            yield

    def _record_notification(self, event: ReportingDomainEvent) -> None:
        if self._notification_state is not None:
            self._notification_state.enqueue(event)

    def _dirty_status(
        self,
        scope: ReportingStatusScope,
        reason: DirtyReason,
        before: ReportingStatusEvidence | None = None,
        after: ReportingStatusEvidence | None = None,
    ) -> None:
        if self._notification_state is not None:
            self._notification_state.mark_dirty(scope, reason, self._clock(), before, after)

    def _resolve_issue_scope(
        self,
        *,
        account_id: str,
        consumer_id: str | None,
        issue_id: str,
        status_scope: ReportingStatusScope | None,
    ) -> ReportingStatusScope:
        """Validate a requested scope refinement without touching retained state.

        Default-off stores deliberately keep no rollback copy, so every caller
        that moves an issue record must clear this check *before* mutating:
        PostgreSQL rolls the statement back, the reference store cannot.
        """
        existing = self._issue_status_scopes.get((account_id, issue_id))
        scope = (
            status_scope or existing or ReportingStatusScope(account_id, consumer_id=consumer_id)
        )
        if scope.account_id != account_id or scope.consumer_id != consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        validate_scope_refinement(existing, scope)
        if scope.generation_key is not None:
            configuration = self._configurations.get(scope.generation_key)
            if configuration is None or (
                scope.feed_purpose is not None and configuration.feed_purpose != scope.feed_purpose
            ):
                raise ReportingNotificationError("invalid_status_scope")
        if scope.reporting_obligation_id is not None:
            obligation = self._obligations.get(scope.reporting_obligation_id)
            if (
                obligation is None
                or obligation.account_id != scope.account_id
                or (
                    scope.generation_key is not None
                    and obligation.generation_key != scope.generation_key
                )
                or (
                    scope.feed_purpose is not None and obligation.feed_purpose != scope.feed_purpose
                )
            ):
                raise ReportingNotificationError("invalid_status_scope")
        return scope

    def _dirty_issue(
        self,
        issue: ReportingIssueLifecycle,
        status_scope: ReportingStatusScope | None,
        before: ReportingIssueLifecycle | None = None,
        *,
        enqueue: bool = True,
    ) -> None:
        key = (issue.account_id, issue.issue_id)
        existing = self._issue_status_scopes.get(key)
        scope = self._resolve_issue_scope(
            account_id=issue.account_id,
            consumer_id=issue.consumer_id,
            issue_id=issue.issue_id,
            status_scope=status_scope,
        )
        self._issue_status_scopes[key] = scope
        if enqueue and (before != issue or existing != scope):
            self._dirty_status(
                scope, "issue", issue_evidence(before) if before else None, issue_evidence(issue)
            )

    async def create_schema(self) -> None:
        return None

    def _append(
        self, account_id: str, kind: LedgerRecordKind, record_id: str, *, consumer_id: str
    ) -> None:
        self._sequence += 1
        self._changes.append((self._sequence, account_id, kind, record_id, self._clock()))
        self._change_owners[self._sequence] = consumer_id

    # -- configurations --------------------------------------------------

    async def put_configuration(self, configuration: ReportingConfiguration) -> None:
        reject_reserved_authoritative_party(configuration)
        async with self._mutation():
            key = configuration.generation_key
            existing = self._configurations.get(key)
            if existing is not None and existing.quarantined:
                raise LedgerConflictError(
                    "REPORTING_GENERATION_QUARANTINED", "operator reconciliation is required"
                )
            if existing is not None and _fingerprint(_config_payload(existing)) != _fingerprint(
                _config_payload(configuration)
            ):
                raise LedgerConflictError(
                    "CONFIGURATION_GENERATION_IMMUTABLE",
                    f"configuration {key.delivery_config_id}@{key.delivery_config_version} "
                    "already exists with different content for this account; "
                    "publish a new version instead of editing a retained generation",
                )
            self._configurations[key] = configuration
            changed = existing is None or configuration_lifecycle(
                existing
            ) != configuration_lifecycle(configuration)
            if changed and self._notification_state is not None:
                self._dirty_status(
                    ReportingStatusScope(configuration.account_id, configuration.generation_key),
                    "configuration",
                    before=configuration_evidence(existing) if existing is not None else None,
                    after=configuration_evidence(configuration),
                )

    async def list_configurations(
        self, *, caller: ReportingCaller, delivery_config_ids: Sequence[str] | None = None
    ) -> tuple[ReportingConfiguration, ...]:
        wanted = set(delivery_config_ids) if delivery_config_ids else None
        return tuple(
            configuration
            for configuration in self._configurations.values()
            if configuration.account_id == caller.account_id
            and configuration.consumer_id == caller.consumer_id
            and (wanted is None or configuration.delivery_config_id in wanted)
        )

    async def list_all_configurations(self) -> tuple[ReportingConfiguration, ...]:
        return tuple(
            sorted(
                (c for c in self._configurations.values() if not c.quarantined),
                key=lambda item: (
                    item.account_id,
                    item.consumer_id,
                    item.delivery_config_id,
                    item.delivery_config_version,
                ),
            )
        )

    # -- obligations -----------------------------------------------------

    async def commit_obligation(
        self, obligation: ReportingObligationRecord
    ) -> ReportingObligationRecord:
        async with self._mutation():
            key = (
                obligation.generation_key,
                _utc(obligation.period.start).isoformat(),
                _utc(obligation.period.end).isoformat(),
            )
            existing_id = self._obligation_by_period.get(key)
            if existing_id is not None:
                return self._obligations[existing_id]
            if obligation.reporting_obligation_id in self._obligations:
                raise LedgerConflictError(
                    "OBLIGATION_IDENTITY_CONFLICT",
                    "the obligation identifier already belongs to a different logical period",
                )
            if obligation.generation_key not in self._configurations:
                raise LedgerConflictError(
                    "UNKNOWN_CONFIGURATION_GENERATION", "generation is unavailable"
                )
            if self._configurations[obligation.generation_key].quarantined:
                raise LedgerConflictError(
                    "REPORTING_GENERATION_QUARANTINED",
                    "establish a new owned generation after operator reconciliation",
                )
            require_frozen_currency(obligation.currency)
            self._obligations[obligation.reporting_obligation_id] = obligation
            self._obligation_by_period[key] = obligation.reporting_obligation_id
            self._append(
                obligation.account_id,
                "obligation",
                obligation.reporting_obligation_id,
                consumer_id=obligation.consumer_id,
            )
            if self._notification_state is not None:
                self._dirty_status(
                    ReportingStatusScope.for_obligation(obligation),
                    "obligation",
                    after=ReportingStatusEvidence("obligation", obligation.reporting_obligation_id),
                )
            return obligation

    async def get_obligation(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> ReportingObligationRecord | None:
        found = self._obligations.get(reporting_obligation_id)
        return found if found and found.account_id == account_id else None

    async def find_obligation(
        self,
        *,
        account_id: str,
        consumer_id: str,
        delivery_config_id: str,
        delivery_config_version: int,
        period_start: datetime,
        period_end: datetime,
    ) -> ReportingObligationRecord | None:
        key = (
            ReportingConfigurationGenerationKey(
                account_id=account_id,
                consumer_id=consumer_id,
                delivery_config_id=delivery_config_id,
                delivery_config_version=delivery_config_version,
            ),
            _utc(period_start).isoformat(),
            _utc(period_end).isoformat(),
        )
        found = self._obligation_by_period.get(key)
        return (
            await self.get_obligation(account_id=account_id, reporting_obligation_id=found)
            if found
            else None
        )

    # -- revisions -------------------------------------------------------

    async def commit_revision(
        self, revision: ReportingRevisionRecord, rows: Sequence[dict[str, Any]]
    ) -> ReportingRevisionRecord:
        rows = tuple(deepcopy(row) for row in rows)
        # Checked first, exactly as PgReportingLedgerStore does: a caller whose
        # rows do not match its own declared count must get the same code from
        # both stores, not whichever invariant that store happens to reach.
        if revision.row_count != len(rows):
            raise LedgerConflictError(
                "ROW_COUNT_MISMATCH",
                f"revision declares {revision.row_count} rows but {len(rows)} were supplied",
            )
        validate_managed_revision_rows(revision, rows)
        async with self._mutation():
            identity = _revision_identity(revision)
            existing = self._revisions.get(revision.reporting_revision_id)
            if existing is not None:
                if existing.account_id != revision.account_id:
                    raise LedgerConflictError(
                        "REVISION_NOT_FOUND", "no such revision for this account"
                    )
                if self._revision_identity[revision.reporting_revision_id] != identity:
                    raise LedgerConflictError(
                        "REVISION_IMMUTABLE",
                        f"revision {revision.reporting_revision_id} already exists with "
                        "different content; a restatement is a new revision",
                    )
                return existing
            obligation = self._obligations.get(revision.reporting_obligation_id)
            if obligation is None or obligation.account_id != revision.account_id:
                raise LedgerConflictError(
                    "OBLIGATION_NOT_FOUND",
                    "a revision must attach to an obligation committed at the period close",
                )
            if self._configurations[obligation.generation_key].quarantined:
                raise LedgerConflictError(
                    "REPORTING_GENERATION_QUARANTINED", "legacy evidence is read-only"
                )
            validate_revision_currency(obligation, revision, rows)
            siblings = [
                item
                for item in self._revisions.values()
                if item.reporting_obligation_id == revision.reporting_obligation_id
            ]
            if revision.finality == "official" and any(
                item.finality == "official" for item in siblings
            ):
                raise LedgerConflictError(
                    "OFFICIAL_REVISION_TERMINAL",
                    "an official revision already exists for this obligation; publish a later "
                    "source correction as an adjustment",
                )
            if revision.supersedes_reporting_revision_id:
                self._require_current_leaf(revision, siblings)
            self._revisions[revision.reporting_revision_id] = revision
            self._revision_identity[revision.reporting_revision_id] = identity
            self._rows[revision.reporting_revision_id] = tuple(dict(row) for row in rows)
            self._append(
                revision.account_id,
                "revision",
                revision.reporting_revision_id,
                consumer_id=obligation.consumer_id,
            )
            if self._notification_state is not None:
                from adcp.reporting.ledger.notification_events import revision_event

                self._record_notification(
                    revision_event(revision, self._clock(), consumer_id=obligation.consumer_id)
                )
                self._dirty_status(
                    ReportingStatusScope.for_obligation(obligation),
                    "revision",
                    after=ReportingStatusEvidence(
                        "revision",
                        revision.reporting_revision_id,
                        readable=revision.readable,
                        supersedes_id=revision.supersedes_reporting_revision_id,
                    ),
                )
            return revision

    def _require_current_leaf(
        self, revision: ReportingRevisionRecord, siblings: Sequence[ReportingRevisionRecord]
    ) -> None:
        target = revision.supersedes_reporting_revision_id
        known = {item.reporting_revision_id for item in siblings}
        if target not in known:
            raise LedgerConflictError(
                "SUPERSEDES_UNKNOWN",
                f"revision {target} is not part of this obligation's chain",
            )
        already = {
            item.supersedes_reporting_revision_id
            for item in siblings
            if item.supersedes_reporting_revision_id
        }
        if target in already:
            raise LedgerConflictError(
                "SUPERSEDES_STALE",
                f"revision {target} has already been superseded; a stale pointer would fork "
                "the chain and let a successful retry erase a recorded restatement",
            )

    async def list_revisions(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> tuple[ReportingRevisionRecord, ...]:
        return tuple(
            item
            for item in self._revisions.values()
            if item.account_id == account_id
            and item.reporting_obligation_id == reporting_obligation_id
        )

    async def get_provisional_acquisition(
        self, *, account_id: str, reporting_obligation_id: str, ordinal: int
    ) -> ProvisionalAcquisition | None:
        return self._provisional_acquisitions.get((account_id, reporting_obligation_id, ordinal))

    async def reserve_provisional_acquisition(
        self, acquisition: ProvisionalAcquisition
    ) -> ProvisionalAcquisition:
        # Canonical round-trip also detaches all request collections.
        acquisition = ProvisionalAcquisition.from_wire(acquisition.to_wire())
        key = (acquisition.account_id, acquisition.obligation_id, acquisition.ordinal)
        async with self._mutation():
            existing = self._provisional_acquisitions.get(key)
            if existing is not None:
                return existing
            obligation = self._obligations.get(acquisition.obligation_id)
            if obligation is None or obligation.account_id != acquisition.account_id:
                raise LedgerConflictError("OBLIGATION_NOT_FOUND", "unknown observation obligation")
            if not acquisition.binds(obligation):
                raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition generation differs")
            if any(
                item.account_id == acquisition.account_id
                and item.execution_key == acquisition.execution_key
                for item in self._provisional_acquisitions.values()
            ):
                raise LedgerConflictError(
                    "OBSERVATION_CONFLICT", "execution key is already reserved"
                )
            checkpoint = self._restatement_checkpoints.get(acquisition.obligation_id)
            expected = (
                checkpoint.next_observation
                if checkpoint
                else len(
                    await self.list_revisions(
                        account_id=acquisition.account_id,
                        reporting_obligation_id=acquisition.obligation_id,
                    )
                )
            )
            if acquisition.ordinal != expected:
                raise LedgerConflictError("OBSERVATION_CONFLICT", "observation ordinal changed")
            self._provisional_acquisitions[key] = acquisition
            return acquisition

    async def get_provisional_observation(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> ProvisionalObservation | None:
        observations = [
            value
            for (account, obligation, _), value in self._provisional_observations.items()
            if account == account_id and obligation == reporting_obligation_id
        ]
        return max(observations, key=lambda item: item.acquisition.ordinal, default=None)

    async def commit_provisional_observation(
        self,
        observation: ProvisionalObservation,
        revision: ReportingRevisionRecord,
        rows: Sequence[dict[str, Any]],
    ) -> ReportingRevisionRecord:
        acquisition = observation.acquisition
        key = (acquisition.account_id, acquisition.obligation_id, acquisition.ordinal)
        if (
            revision.account_id != acquisition.account_id
            or revision.reporting_obligation_id != acquisition.obligation_id
            or revision.reporting_revision_id != observation.revision_id
        ):
            raise LedgerConflictError("OBSERVATION_CONFLICT", "observation identity differs")
        async with self._mutation():
            existing = self._provisional_observations.get(key)
            if existing is not None:
                if (
                    existing.acquisition != acquisition
                    or existing.revision_id != observation.revision_id
                ):
                    raise LedgerConflictError("OBSERVATION_CONFLICT", "observation replay differs")
                if existing.revision_id not in self._revisions:
                    raise LedgerConflictError(
                        "HISTORY_UNAVAILABLE", "observation revision is missing"
                    )
                # A matching observation identity does not make conflicting
                # publication content an idempotent replay.
                return await self.commit_revision(revision, rows)
            retained = self._provisional_acquisitions.get(key)
            if retained != acquisition:
                raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition was not reserved")
            checkpoint = self._restatement_checkpoints.get(acquisition.obligation_id)
            if checkpoint is not None and checkpoint.next_observation != acquisition.ordinal:
                raise LedgerConflictError("OBSERVATION_CONFLICT", "observation ordinal changed")
            committed = await self.commit_revision(revision, rows)
            await self.record_restatement_checkpoint(
                RestatementCheckpoint(
                    acquisition.account_id,
                    acquisition.obligation_id,
                    observation.checked_at,
                    acquisition.ordinal + 1,
                    observation.provisional_until,
                )
            )
            self._provisional_observations[key] = observation
            return committed

    async def get_restatement_checkpoint(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> RestatementCheckpoint | None:
        checkpoint = self._restatement_checkpoints.get(reporting_obligation_id)
        return checkpoint if checkpoint and checkpoint.account_id == account_id else None

    async def get_retry_schedule(self, *, scope_key: str) -> RetryScheduleEntry | None:
        return self._retry_schedules.get(scope_key)

    async def record_retry_schedule(self, entry: RetryScheduleEntry) -> RetryScheduleEntry:
        async with self._mutation():
            existing = self._retry_schedules.get(entry.scope_key)
            if existing is not None:
                assert entry.recorded_at is not None and existing.recorded_at is not None
                if _utc(entry.recorded_at) < _utc(existing.recorded_at):
                    return existing
                if entry.attempt > 0 and not entry.blocked:
                    if existing.blocked and not entry.replayed:
                        return existing
                    if entry.attempt < existing.attempt and not entry.replayed:
                        return existing
                    if not entry.replayed and not existing.blocked:
                        entry = replace(
                            entry,
                            retry_not_before=max(
                                entry.retry_not_before, existing.retry_not_before, key=_utc
                            ),
                        )
            stored = replace(entry, replayed=False)
            self._retry_schedules[entry.scope_key] = stored
            return stored

    async def record_restatement_checkpoint(
        self, checkpoint: RestatementCheckpoint
    ) -> RestatementCheckpoint:
        async with self._mutation():
            obligation = self._obligations.get(checkpoint.reporting_obligation_id)
            if obligation is None or obligation.account_id != checkpoint.account_id:
                raise LedgerConflictError(
                    "OBLIGATION_NOT_FOUND",
                    "a restatement checkpoint must attach to an obligation for this account",
                )
            existing = self._restatement_checkpoints.get(checkpoint.reporting_obligation_id)
            if existing is not None:
                if checkpoint.next_observation < existing.next_observation:
                    return existing
                if checkpoint.next_observation == existing.next_observation and _utc(
                    checkpoint.checked_at
                ) <= _utc(existing.checked_at):
                    return existing
            self._restatement_checkpoints[checkpoint.reporting_obligation_id] = checkpoint
            return checkpoint

    async def get_revision(
        self, *, account_id: str, reporting_revision_id: str
    ) -> ReportingRevisionRecord | None:
        found = self._revisions.get(reporting_revision_id)
        return found if found and found.account_id == account_id else None

    async def read_revision_rows(
        self,
        *,
        account_id: str,
        reporting_revision_id: str,
        cursor: str | None = None,
        limit: int = 500,
    ) -> ReportingRowPage:
        revision = await self.get_revision(
            account_id=account_id, reporting_revision_id=reporting_revision_id
        )
        if revision is None:
            raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account")
        rows = self._rows.get(reporting_revision_id, ())
        offset = revision_row_offset(cursor, reporting_revision_id, limit)
        window = rows[offset : offset + limit]
        has_more = offset + limit < len(rows)
        return ReportingRowPage(
            reporting_revision_id=reporting_revision_id,
            rows=tuple(deepcopy(row) for row in window),
            total_count=len(rows),
            has_more=has_more,
            cursor=(
                encode_cursor(
                    {"ownership": 2, "revision": reporting_revision_id, "offset": offset + limit}
                )
                if has_more
                else None
            ),
        )

    async def set_revision_readable(
        self, *, account_id: str, reporting_revision_id: str, readable: bool
    ) -> None:
        async with self._mutation():
            existing = self._revisions.get(reporting_revision_id)
            if existing is None or existing.account_id != account_id:
                raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account")
            if existing.readable == readable:
                return
            self._revisions[reporting_revision_id] = replace(existing, readable=readable)
            if self._notification_state is not None:
                obligation = self._obligations[existing.reporting_obligation_id]
                self._dirty_status(
                    ReportingStatusScope.for_obligation(obligation),
                    "readability",
                    ReportingStatusEvidence(
                        "revision", reporting_revision_id, readable=existing.readable
                    ),
                    ReportingStatusEvidence("revision", reporting_revision_id, readable=readable),
                )

    # -- adjustments -----------------------------------------------------

    async def commit_adjustment(
        self, adjustment: ReportingAdjustmentRecord
    ) -> ReportingAdjustmentRecord:
        async with self._mutation():
            existing = self._adjustments.get(adjustment.reporting_adjustment_id)
            if existing is not None:
                if existing.account_id != adjustment.account_id:
                    raise LedgerConflictError("ADJUSTMENT_UNAVAILABLE", "adjustment is unavailable")
                if existing != adjustment:
                    raise LedgerConflictError(
                        "ADJUSTMENT_IMMUTABLE", "adjustment content is immutable"
                    )
                return existing
            revision = self._revisions.get(adjustment.adjusts_reporting_revision_id)
            if revision is None or revision.account_id != adjustment.account_id:
                raise LedgerConflictError(
                    "REVISION_NOT_FOUND", "an adjustment must name a committed revision"
                )
            if revision.finality != "official":
                raise LedgerConflictError(
                    "ADJUSTMENT_REQUIRES_OFFICIAL",
                    "adjustments correct an official revision; restate a snapshot with a "
                    "superseding snapshot revision instead",
                )
            obligation = self._obligations[revision.reporting_obligation_id]
            if self._configurations[obligation.generation_key].quarantined:
                raise LedgerConflictError(
                    "REPORTING_GENERATION_QUARANTINED", "legacy evidence is read-only"
                )
            validate_adjustment_currency(
                self._obligations[revision.reporting_obligation_id], adjustment
            )
            self._adjustments[adjustment.reporting_adjustment_id] = adjustment
            self._append(
                adjustment.account_id,
                "adjustment",
                adjustment.reporting_adjustment_id,
                consumer_id=obligation.consumer_id,
            )
            if self._notification_state is not None:
                from adcp.reporting.ledger.notification_events import adjustment_event

                self._record_notification(
                    adjustment_event(
                        adjustment,
                        self._clock(),
                        consumer_id=self._obligations[revision.reporting_obligation_id].consumer_id,
                    )
                )
                self._dirty_status(
                    ReportingStatusScope.for_obligation(
                        self._obligations[revision.reporting_obligation_id]
                    ),
                    "adjustment",
                    after=ReportingStatusEvidence("adjustment", adjustment.reporting_adjustment_id),
                )
            return adjustment

    async def list_adjustments(
        self, *, account_id: str, reporting_revision_ids: Sequence[str]
    ) -> tuple[ReportingAdjustmentRecord, ...]:
        wanted = set(reporting_revision_ids)
        return tuple(
            item
            for item in self._adjustments.values()
            if item.account_id == account_id and item.adjusts_reporting_revision_id in wanted
        )

    # -- consumer status -------------------------------------------------

    async def resolve_consumer_status_replay(
        self, status: ConsumerStatusRecord
    ) -> ConsumerStatusRecord | None:
        async with self._lock:
            return self._replay(status)

    def _replay(self, status: ConsumerStatusRecord) -> ConsumerStatusRecord | None:
        key = (status.account_id, status.consumer_id, status.reporting_status_id)
        existing = self._statuses.get(key)
        if existing is None:
            return None
        if self._status_identity[key] != _consumer_status_identity(status):
            raise LedgerConflictError(
                "STATUS_IDENTITY_CONFLICT",
                f"reporting_status_id {status.reporting_status_id} was already "
                "recorded with different content",
            )
        return existing

    async def record_consumer_status(
        self, status: ConsumerStatusRecord
    ) -> tuple[ConsumerStatusRecord, bool]:
        from adcp.reporting.evidence import consumer_reference

        consumer_reference(status.consumer_id)
        async with self._mutation():
            identity = _consumer_status_identity(status)
            existing = self._replay(status)
            if existing is not None:
                return existing, False
            from adcp.reporting.ledger.status_snapshot import (
                memory_snapshot,
                validate_status_evidence,
            )

            validate_status_evidence(status, memory_snapshot(self, status.account_id))
            leaf = self._current_status_leaf(status.chain_key)
            if status.supersedes_reporting_status_id:
                if (
                    leaf is None
                    or leaf.reporting_status_id != status.supersedes_reporting_status_id
                ):
                    raise LedgerConflictError(
                        "STATUS_SUPERSEDES_STALE",
                        "supersedes_reporting_status_id must name this chain's current leaf; "
                        "a stale pointer would let a successful retry erase a recorded outage",
                    )
                self._statuses[(leaf.account_id, leaf.consumer_id, leaf.reporting_status_id)] = (
                    replace(leaf, superseded=True)
                )
            elif leaf is not None:
                raise LedgerConflictError(
                    "STATUS_SUPERSEDES_REQUIRED",
                    "this chain already has a current statement; a new statement must "
                    "explicitly supersede it",
                )
            key = (status.account_id, status.consumer_id, status.reporting_status_id)
            self._statuses[key] = status
            self._status_identity[key] = identity
            self._append(
                status.account_id,
                "consumer_status",
                status.reporting_status_id,
                consumer_id=status.consumer_id,
            )
            from adcp.reporting.ledger.status_snapshot import settle_memory_snapshot

            settle_memory_snapshot(self, status.account_id)
            if self._notification_state is not None:
                self._dirty_status(
                    ReportingStatusScope(
                        status.account_id,
                        status.generation_key,
                        status.reporting_obligation_id,
                        status.consumer_id,
                    ),
                    "consumer_status",
                    (
                        ReportingStatusEvidence("consumer_status", leaf.reporting_status_id)
                        if leaf is not None
                        else None
                    ),
                    ReportingStatusEvidence(
                        "consumer_status",
                        status.reporting_status_id,
                        supersedes_id=status.supersedes_reporting_status_id,
                    ),
                )
            return status, True

    async def record_consumer_status_with_lifecycle(
        self, status: ConsumerStatusRecord
    ) -> tuple[ConsumerStatusRecord, bool]:
        return await self.record_consumer_status(status)

    async def read_status_snapshot(self, *, caller: ReportingCaller) -> ReportingStatusSnapshot:
        from adcp.reporting.ledger.status_snapshot import settle_memory_snapshot

        async with self._mutation():
            from adcp.reporting.ledger.delivery_models import ReportingDeliveryPrincipal
            from adcp.reporting.materializer.capture import private_snapshot

            return private_snapshot(
                settle_memory_snapshot(self, caller.account_id),
                ReportingDeliveryPrincipal(caller.account_id, caller.consumer_id),
            )

    def _current_status_leaf(self, chain_key: tuple[Any, ...]) -> ConsumerStatusRecord | None:
        for item in self._statuses.values():
            if item.chain_key == chain_key and not item.superseded:
                return item
        return None

    async def list_consumer_statuses(
        self,
        *,
        account_id: str,
        consumer_id: str,
        reporting_obligation_ids: Sequence[str] | None = None,
    ) -> tuple[ConsumerStatusRecord, ...]:
        wanted = set(reporting_obligation_ids) if reporting_obligation_ids is not None else None
        result = []
        for item in self._statuses.values():
            if item.account_id != account_id or item.consumer_id != consumer_id:
                continue
            if wanted is not None and not self._matches_obligations(item, wanted):
                continue
            result.append(item)
        return tuple(result)

    def _matches_obligations(self, status: ConsumerStatusRecord, wanted: set[str]) -> bool:
        """Attach a statement to an obligation by id *or* by logical period key.

        The logical-key path is what makes the loop repairable: a statement
        filed as ``obligation_missing`` predates the obligation it is about, and
        when the seller later creates that obligation the existing chain must
        attach to it rather than being lost, forked, or reset.
        """
        for obligation_id in wanted:
            obligation = self._obligations.get(obligation_id)
            if obligation is None or obligation.account_id != status.account_id:
                continue
            if status.reporting_obligation_id == obligation_id:
                return True
            if (
                obligation.generation_key == status.generation_key
                and obligation.report_definition_id == status.report_definition_id
                and _utc(obligation.period.start) == _utc(status.period_start)
                and _utc(obligation.period.end) == _utc(status.period_end)
            ):
                return True
        return False

    # -- issue lifecycle -------------------------------------------------

    async def ensure_issue_opened(
        self,
        *,
        issue_key: str,
        account_id: str,
        consumer_id: str | None,
        observed_at: datetime,
        status_scope: ReportingStatusScope | None = None,
    ) -> ReportingIssueLifecycle:
        async with self._mutation():
            key = (account_id, issue_key)
            live = self._issues.get(key)
            if live is not None:
                if live.consumer_id != consumer_id:
                    raise ReportingNotificationError("invalid_status_scope")
                previous_scope = self._issue_status_scopes.get((account_id, live.issue_id))
                if status_scope is not None:
                    validate_scope_refinement(previous_scope, status_scope)
                else:
                    status_scope = previous_scope
            if live is not None and live.live:
                self._dirty_issue(live, status_scope, live)
                return live
            generation = self._issue_generations.get(key, 0) + 1
            record = ReportingIssueLifecycle(
                issue_key=issue_key,
                issue_id=issue_id_for_occurrence(issue_key, generation),
                account_id=account_id,
                consumer_id=consumer_id,
                opened_at=_utc(observed_at),
                issue_state="open",
                generation=generation,
            )
            self._resolve_issue_scope(
                account_id=account_id,
                consumer_id=consumer_id,
                issue_id=record.issue_id,
                status_scope=status_scope,
            )
            self._issue_generations[key] = generation
            self._issues[key] = record
            self._dirty_issue(record, status_scope)
            return record

    async def set_issue_state(
        self,
        *,
        issue_key: str,
        account_id: str,
        state: Literal["acknowledged", "waived"],
        at: datetime,
        external_ref: str | None = None,
        status_scope: ReportingStatusScope | None = None,
    ) -> ReportingIssueLifecycle:
        async with self._mutation():
            if state not in {"acknowledged", "waived"}:
                raise LedgerConflictError(
                    "ISSUE_STATE_NOT_OPERATOR_SETTABLE",
                    f"issue_state {state!r} is not settable by an operator. 'resolved' is "
                    "reachable only when the condition actually clears -- the projection "
                    "retires it -- because a seller must not retire a mismatch out of a "
                    "degraded projection while the statement that caused it is still the "
                    "consumer's current leaf",
                )
            key = (account_id, issue_key)
            live = self._issues.get(key)
            if live is None or not live.live:
                raise LedgerConflictError(
                    "ISSUE_NOT_OPEN",
                    f"no open issue {issue_key!r} for this account; a retired issue cannot be "
                    "reopened, and a recurrence gets a new occurrence",
                )
            check_issue_state_transition(live.issue_state, state)
            if state == "waived" and live.issue_state != "waived":
                from adcp.reporting.ledger.status_projection import bind_mismatch_waiver
                from adcp.reporting.ledger.status_snapshot import memory_snapshot

                live = bind_mismatch_waiver(memory_snapshot(self, account_id), live)
            updated = replace(
                live,
                issue_state=state,
                external_ref=external_ref or live.external_ref,
                # Set on the way into a retired state and never cleared.
                retired_at=(
                    (
                        live.retired_at
                        if live.waived_reporting_status_id is not None
                        and live.retired_at is not None
                        else _utc(at)
                    )
                    if state == "waived"
                    else live.retired_at
                ),
            )
            self._resolve_issue_scope(
                account_id=account_id,
                consumer_id=live.consumer_id,
                issue_id=live.issue_id,
                status_scope=status_scope,
            )
            self._issues[key] = updated
            # Derive the no-op from the resulting record rather than predicting
            # it: an idempotent re-acknowledge changes nothing and enqueues
            # nothing, while anything that does move retained evidence stays
            # reconstructable for the projector.
            self._dirty_issue(updated, status_scope, live)
            return updated

    async def retire_issue(
        self,
        *,
        issue_key: str,
        account_id: str,
        at: datetime,
        status_scope: ReportingStatusScope | None = None,
    ) -> ReportingIssueLifecycle | None:
        async with self._mutation():
            key = (account_id, issue_key)
            live = self._issues.get(key)
            if live is None or not live.live or not issue_is_retirable(live.issue_state):
                # Convergent: nothing live, or already waived. A waived issue
                # is retired from the projection by agreement, and overwriting
                # that readable act with `resolved` is an edge the forward-only
                # lifecycle forbids -- enforced here rather than left to each
                # caller to remember.
                return None
            check_issue_state_transition(live.issue_state, "resolved")
            retired = replace(live, issue_state="resolved", retired_at=_utc(at))
            self._resolve_issue_scope(
                account_id=account_id,
                consumer_id=live.consumer_id,
                issue_id=live.issue_id,
                status_scope=status_scope,
            )
            self._issues[key] = retired
            self._dirty_issue(retired, status_scope, live)
            return retired

    async def get_issue(self, *, issue_key: str, account_id: str) -> ReportingIssueLifecycle | None:
        async with self._lock:
            # Live occurrences only, matching the Protocol docstring and the
            # Postgres store. Returning a retired row here would let the
            # projection re-emit a resolved issue's opened_at.
            live = self._issues.get((account_id, issue_key))
            return live if live is not None and live.live else None

    # -- snapshots -------------------------------------------------------

    async def open_snapshot(
        self, *, caller: ReportingCaller, filters_fingerprint: str
    ) -> LedgerSnapshot:
        account_id = caller.account_id
        async with self._lock:
            as_of = _utc(self._clock())
            max_sequence = max(
                (
                    sequence
                    for sequence, account, kind, record_id, _ in self._changes
                    if account == account_id and self._change_owners[sequence] == caller.consumer_id
                ),
                default=0,
            )
            return LedgerSnapshot(
                snapshot_id="rpls_"
                + _fingerprint([account_id, caller.consumer_id, filters_fingerprint, max_sequence])[
                    :32
                ],
                account_id=account_id,
                consumer_id=caller.consumer_id,
                ledger_as_of=as_of,
                max_sequence=max_sequence,
            )

    async def read_page(
        self,
        *,
        snapshot: LedgerSnapshot,
        consumer_id: str | None,
        delivery_config_ids: Sequence[str] | None,
        media_buy_ids: Sequence[str] | None,
        offset: int,
        limit: int,
        changes_after_sequence: int | None,
        feed_purposes: Sequence[str] | None = None,
        period_start: datetime | None = None,
        period_end: datetime | None = None,
    ) -> LedgerPage:
        if consumer_id not in {None, snapshot.consumer_id}:
            raise LedgerConflictError(
                "CURSOR_SNAPSHOT_MISMATCH", "snapshot belongs to another caller"
            )
        consumer_id = snapshot.consumer_id
        lower = changes_after_sequence or 0
        selected = [
            item
            for item in self._changes
            if item[1] == snapshot.account_id
            and self._change_owners[item[0]] == consumer_id
            and lower < item[0] <= snapshot.max_sequence
        ]
        config_filter = set(delivery_config_ids) if delivery_config_ids else None
        media_buy_filter = set(media_buy_ids) if media_buy_ids else None

        records: list[tuple[int, LedgerRecordKind, Any]] = []
        for sequence, _account, kind, record_id, _committed in sorted(selected):
            record = self._resolve(kind, record_id, snapshot.account_id, consumer_id)
            if record is None or record.account_id != snapshot.account_id:
                continue
            if not self._in_scope(
                kind,
                record,
                config_filter,
                media_buy_filter,
                consumer_id,
                feed_purposes,
                period_start,
                period_end,
            ):
                continue
            records.append((sequence, kind, record))

        window = records[offset : offset + limit]
        has_more = offset + limit < len(records)
        return LedgerPage(
            obligations=tuple(item[2] for item in window if item[1] == "obligation"),
            revisions=tuple(item[2] for item in window if item[1] == "revision"),
            adjustments=tuple(item[2] for item in window if item[1] == "adjustment"),
            consumer_statuses=tuple(item[2] for item in window if item[1] == "consumer_status"),
            total_count=len(records),
            has_more=has_more,
            cursor=(
                encode_cursor({"snapshot": snapshot.snapshot_id, "offset": offset + limit})
                if has_more
                else None
            ),
        )

    def _resolve(
        self, kind: str, record_id: str, account_id: str = "", consumer_id: str | None = None
    ) -> Any:
        if kind == "obligation":
            return self._obligations.get(record_id)
        if kind == "revision":
            return self._revisions.get(record_id)
        if kind == "adjustment":
            return self._adjustments.get(record_id)
        if kind == "consumer_status":
            return self._statuses.get((account_id, consumer_id or "", record_id))
        # Optional delivery records have their own retained snapshot reader.
        # Core status must not count or expose them, even in a shared ledger.
        return None

    def _in_scope(
        self,
        kind: LedgerRecordKind,
        record: Any,
        config_filter: set[str] | None,
        media_buy_filter: set[str] | None,
        consumer_id: str | None,
        feed_purposes: Sequence[str] | None = None,
        period_start: datetime | None = None,
        period_end: datetime | None = None,
    ) -> bool:
        from adcp.reporting.ledger.status_projection import configuration_selected, period_selected

        if kind == "consumer_status":
            # A caller sees only its own statements; another consumer's
            # operational status is not disclosed.
            if consumer_id is None or record.consumer_id != consumer_id:
                return False
            configuration = self._configurations.get(record.generation_key)
            return (
                configuration is not None
                and configuration_selected(
                    configuration,
                    delivery_config_ids=tuple(config_filter or ()),
                    media_buy_ids=tuple(media_buy_filter or ()),
                    feed_purposes=feed_purposes or (),
                )
                and period_selected(
                    record.period_start, record.period_end, period_start, period_end
                )
            )
        obligation = self._obligation_for(kind, record)
        if obligation is None or obligation.consumer_id != consumer_id:
            return False
        configuration = self._configurations.get(obligation.generation_key)
        if (
            configuration is None
            or not configuration_selected(
                configuration,
                delivery_config_ids=tuple(config_filter or ()),
                media_buy_ids=tuple(media_buy_filter or ()),
                feed_purposes=feed_purposes or (),
            )
            or not period_selected(
                obligation.period.start, obligation.period.end, period_start, period_end
            )
        ):
            return False
        if config_filter is not None and obligation.delivery_config_id not in config_filter:
            return False
        if media_buy_filter is not None and not media_buy_filter.intersection(
            obligation.media_buy_ids
        ):
            return False
        return True

    def _obligation_for(
        self, kind: LedgerRecordKind, record: Any
    ) -> ReportingObligationRecord | None:
        if kind == "obligation":
            found: ReportingObligationRecord = record
            return found
        if kind == "revision":
            return self._obligations.get(record.reporting_obligation_id)
        revision = self._revisions.get(record.adjusts_reporting_revision_id)
        return self._obligations.get(revision.reporting_obligation_id) if revision else None

    # -- leasing ---------------------------------------------------------

    def _configuration_lease_eligible(self, configuration: ReportingConfiguration) -> bool:
        return not configuration.quarantined

    async def lease_period_close(
        self, *, worker_id: str, now: datetime, lease_seconds: float
    ) -> LeasedConfiguration | None:
        from datetime import timedelta

        async with self._lock:
            moment = _utc(now)
            # Rank leasable generations exactly the way the SQL store's
            # `ORDER BY lease_turn, lease_expires_at NULLS FIRST, ...` does:
            # whichever generation went longest without a turn goes first.
            #
            # The turn has to be the *primary* term. The caller's `now` has
            # already excluded every live lease, so among the survivors the
            # expiry carries no fairness information -- and preferring unheld
            # over expired ahead of the turn starves a crashed generation
            # forever: a peer that is leased and released every turn is always
            # unheld, so it wins every comparison while the generation whose
            # worker died stays expired and never closes another period.
            # Expiry and the generation key only break exact turn ties, so the
            # order stays total and never depends on physical layout.
            ranked: list[
                tuple[
                    tuple[int, int, float, str, str, str, int],
                    ReportingConfigurationGenerationKey,
                ]
            ] = []
            for key in self._configurations:
                if not self._configuration_lease_eligible(self._configurations[key]):
                    continue
                turn = self._lease_turns.get(key, 0)
                held = self._leases.get(key)
                tail = (
                    key.account_id,
                    key.consumer_id,
                    key.delivery_config_id,
                    key.delivery_config_version,
                )
                if held is None:
                    ranked.append(((turn, 0, 0.0, *tail), key))
                elif _utc(held[1]) <= moment:
                    ranked.append(((turn, 1, _utc(held[1]).timestamp(), *tail), key))
            if not ranked:
                return None
            # The generation key is part of the rank, so the order is total and
            # both stores make the same choice. Ranking by `min` alone would
            # fall back to whichever generation this process happened to accept
            # first, which the SQL store cannot reproduce and no adopter can
            # observe consistently across a restart or a second worker.
            key = min(ranked, key=lambda item: item[0])[1]
            configuration = self._configurations[key]
            expires = moment + timedelta(seconds=lease_seconds)
            self._lease_turn += 1
            self._lease_turns[key] = self._lease_turn
            self._leases[key] = (worker_id, expires)
            return LeasedConfiguration(
                account_id=configuration.account_id,
                consumer_id=configuration.consumer_id,
                delivery_config_id=configuration.delivery_config_id,
                delivery_config_version=configuration.delivery_config_version,
                lease_expires_at=expires,
            )

    async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) -> None:
        async with self._lock:
            key = lease.generation_key
            held = self._leases.get(key)
            if held == (worker_id, _utc(lease.lease_expires_at)):
                del self._leases[key]

Process-local reference store. Correct, ordered, and not durable.

clock supplies change timestamps and the snapshot observation boundary, standing in for the database clock a durable store reads. Override it to place a test's ledger boundary at a deliberate instant.

Subclasses

Methods

async def commit_adjustment(self, adjustment: ReportingAdjustmentRecord) ‑> ReportingAdjustmentRecord
Expand source code
async def commit_adjustment(
    self, adjustment: ReportingAdjustmentRecord
) -> ReportingAdjustmentRecord:
    async with self._mutation():
        existing = self._adjustments.get(adjustment.reporting_adjustment_id)
        if existing is not None:
            if existing.account_id != adjustment.account_id:
                raise LedgerConflictError("ADJUSTMENT_UNAVAILABLE", "adjustment is unavailable")
            if existing != adjustment:
                raise LedgerConflictError(
                    "ADJUSTMENT_IMMUTABLE", "adjustment content is immutable"
                )
            return existing
        revision = self._revisions.get(adjustment.adjusts_reporting_revision_id)
        if revision is None or revision.account_id != adjustment.account_id:
            raise LedgerConflictError(
                "REVISION_NOT_FOUND", "an adjustment must name a committed revision"
            )
        if revision.finality != "official":
            raise LedgerConflictError(
                "ADJUSTMENT_REQUIRES_OFFICIAL",
                "adjustments correct an official revision; restate a snapshot with a "
                "superseding snapshot revision instead",
            )
        obligation = self._obligations[revision.reporting_obligation_id]
        if self._configurations[obligation.generation_key].quarantined:
            raise LedgerConflictError(
                "REPORTING_GENERATION_QUARANTINED", "legacy evidence is read-only"
            )
        validate_adjustment_currency(
            self._obligations[revision.reporting_obligation_id], adjustment
        )
        self._adjustments[adjustment.reporting_adjustment_id] = adjustment
        self._append(
            adjustment.account_id,
            "adjustment",
            adjustment.reporting_adjustment_id,
            consumer_id=obligation.consumer_id,
        )
        if self._notification_state is not None:
            from adcp.reporting.ledger.notification_events import adjustment_event

            self._record_notification(
                adjustment_event(
                    adjustment,
                    self._clock(),
                    consumer_id=self._obligations[revision.reporting_obligation_id].consumer_id,
                )
            )
            self._dirty_status(
                ReportingStatusScope.for_obligation(
                    self._obligations[revision.reporting_obligation_id]
                ),
                "adjustment",
                after=ReportingStatusEvidence("adjustment", adjustment.reporting_adjustment_id),
            )
        return adjustment
async def commit_obligation(self, obligation: ReportingObligationRecord) ‑> ReportingObligationRecord
Expand source code
async def commit_obligation(
    self, obligation: ReportingObligationRecord
) -> ReportingObligationRecord:
    async with self._mutation():
        key = (
            obligation.generation_key,
            _utc(obligation.period.start).isoformat(),
            _utc(obligation.period.end).isoformat(),
        )
        existing_id = self._obligation_by_period.get(key)
        if existing_id is not None:
            return self._obligations[existing_id]
        if obligation.reporting_obligation_id in self._obligations:
            raise LedgerConflictError(
                "OBLIGATION_IDENTITY_CONFLICT",
                "the obligation identifier already belongs to a different logical period",
            )
        if obligation.generation_key not in self._configurations:
            raise LedgerConflictError(
                "UNKNOWN_CONFIGURATION_GENERATION", "generation is unavailable"
            )
        if self._configurations[obligation.generation_key].quarantined:
            raise LedgerConflictError(
                "REPORTING_GENERATION_QUARANTINED",
                "establish a new owned generation after operator reconciliation",
            )
        require_frozen_currency(obligation.currency)
        self._obligations[obligation.reporting_obligation_id] = obligation
        self._obligation_by_period[key] = obligation.reporting_obligation_id
        self._append(
            obligation.account_id,
            "obligation",
            obligation.reporting_obligation_id,
            consumer_id=obligation.consumer_id,
        )
        if self._notification_state is not None:
            self._dirty_status(
                ReportingStatusScope.for_obligation(obligation),
                "obligation",
                after=ReportingStatusEvidence("obligation", obligation.reporting_obligation_id),
            )
        return obligation
async def commit_provisional_observation(self,
observation: ProvisionalObservation,
revision: ReportingRevisionRecord,
rows: Sequence[dict[str, Any]]) ‑> ReportingRevisionRecord
Expand source code
async def commit_provisional_observation(
    self,
    observation: ProvisionalObservation,
    revision: ReportingRevisionRecord,
    rows: Sequence[dict[str, Any]],
) -> ReportingRevisionRecord:
    acquisition = observation.acquisition
    key = (acquisition.account_id, acquisition.obligation_id, acquisition.ordinal)
    if (
        revision.account_id != acquisition.account_id
        or revision.reporting_obligation_id != acquisition.obligation_id
        or revision.reporting_revision_id != observation.revision_id
    ):
        raise LedgerConflictError("OBSERVATION_CONFLICT", "observation identity differs")
    async with self._mutation():
        existing = self._provisional_observations.get(key)
        if existing is not None:
            if (
                existing.acquisition != acquisition
                or existing.revision_id != observation.revision_id
            ):
                raise LedgerConflictError("OBSERVATION_CONFLICT", "observation replay differs")
            if existing.revision_id not in self._revisions:
                raise LedgerConflictError(
                    "HISTORY_UNAVAILABLE", "observation revision is missing"
                )
            # A matching observation identity does not make conflicting
            # publication content an idempotent replay.
            return await self.commit_revision(revision, rows)
        retained = self._provisional_acquisitions.get(key)
        if retained != acquisition:
            raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition was not reserved")
        checkpoint = self._restatement_checkpoints.get(acquisition.obligation_id)
        if checkpoint is not None and checkpoint.next_observation != acquisition.ordinal:
            raise LedgerConflictError("OBSERVATION_CONFLICT", "observation ordinal changed")
        committed = await self.commit_revision(revision, rows)
        await self.record_restatement_checkpoint(
            RestatementCheckpoint(
                acquisition.account_id,
                acquisition.obligation_id,
                observation.checked_at,
                acquisition.ordinal + 1,
                observation.provisional_until,
            )
        )
        self._provisional_observations[key] = observation
        return committed
async def commit_revision(self, revision: ReportingRevisionRecord, rows: Sequence[dict[str, Any]]) ‑> ReportingRevisionRecord
Expand source code
async def commit_revision(
    self, revision: ReportingRevisionRecord, rows: Sequence[dict[str, Any]]
) -> ReportingRevisionRecord:
    rows = tuple(deepcopy(row) for row in rows)
    # Checked first, exactly as PgReportingLedgerStore does: a caller whose
    # rows do not match its own declared count must get the same code from
    # both stores, not whichever invariant that store happens to reach.
    if revision.row_count != len(rows):
        raise LedgerConflictError(
            "ROW_COUNT_MISMATCH",
            f"revision declares {revision.row_count} rows but {len(rows)} were supplied",
        )
    validate_managed_revision_rows(revision, rows)
    async with self._mutation():
        identity = _revision_identity(revision)
        existing = self._revisions.get(revision.reporting_revision_id)
        if existing is not None:
            if existing.account_id != revision.account_id:
                raise LedgerConflictError(
                    "REVISION_NOT_FOUND", "no such revision for this account"
                )
            if self._revision_identity[revision.reporting_revision_id] != identity:
                raise LedgerConflictError(
                    "REVISION_IMMUTABLE",
                    f"revision {revision.reporting_revision_id} already exists with "
                    "different content; a restatement is a new revision",
                )
            return existing
        obligation = self._obligations.get(revision.reporting_obligation_id)
        if obligation is None or obligation.account_id != revision.account_id:
            raise LedgerConflictError(
                "OBLIGATION_NOT_FOUND",
                "a revision must attach to an obligation committed at the period close",
            )
        if self._configurations[obligation.generation_key].quarantined:
            raise LedgerConflictError(
                "REPORTING_GENERATION_QUARANTINED", "legacy evidence is read-only"
            )
        validate_revision_currency(obligation, revision, rows)
        siblings = [
            item
            for item in self._revisions.values()
            if item.reporting_obligation_id == revision.reporting_obligation_id
        ]
        if revision.finality == "official" and any(
            item.finality == "official" for item in siblings
        ):
            raise LedgerConflictError(
                "OFFICIAL_REVISION_TERMINAL",
                "an official revision already exists for this obligation; publish a later "
                "source correction as an adjustment",
            )
        if revision.supersedes_reporting_revision_id:
            self._require_current_leaf(revision, siblings)
        self._revisions[revision.reporting_revision_id] = revision
        self._revision_identity[revision.reporting_revision_id] = identity
        self._rows[revision.reporting_revision_id] = tuple(dict(row) for row in rows)
        self._append(
            revision.account_id,
            "revision",
            revision.reporting_revision_id,
            consumer_id=obligation.consumer_id,
        )
        if self._notification_state is not None:
            from adcp.reporting.ledger.notification_events import revision_event

            self._record_notification(
                revision_event(revision, self._clock(), consumer_id=obligation.consumer_id)
            )
            self._dirty_status(
                ReportingStatusScope.for_obligation(obligation),
                "revision",
                after=ReportingStatusEvidence(
                    "revision",
                    revision.reporting_revision_id,
                    readable=revision.readable,
                    supersedes_id=revision.supersedes_reporting_revision_id,
                ),
            )
        return revision
async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    return None
async def ensure_issue_opened(self,
*,
issue_key: str,
account_id: str,
consumer_id: str | None,
observed_at: datetime,
status_scope: ReportingStatusScope | None = None) ‑> ReportingIssueLifecycle
Expand source code
async def ensure_issue_opened(
    self,
    *,
    issue_key: str,
    account_id: str,
    consumer_id: str | None,
    observed_at: datetime,
    status_scope: ReportingStatusScope | None = None,
) -> ReportingIssueLifecycle:
    async with self._mutation():
        key = (account_id, issue_key)
        live = self._issues.get(key)
        if live is not None:
            if live.consumer_id != consumer_id:
                raise ReportingNotificationError("invalid_status_scope")
            previous_scope = self._issue_status_scopes.get((account_id, live.issue_id))
            if status_scope is not None:
                validate_scope_refinement(previous_scope, status_scope)
            else:
                status_scope = previous_scope
        if live is not None and live.live:
            self._dirty_issue(live, status_scope, live)
            return live
        generation = self._issue_generations.get(key, 0) + 1
        record = ReportingIssueLifecycle(
            issue_key=issue_key,
            issue_id=issue_id_for_occurrence(issue_key, generation),
            account_id=account_id,
            consumer_id=consumer_id,
            opened_at=_utc(observed_at),
            issue_state="open",
            generation=generation,
        )
        self._resolve_issue_scope(
            account_id=account_id,
            consumer_id=consumer_id,
            issue_id=record.issue_id,
            status_scope=status_scope,
        )
        self._issue_generations[key] = generation
        self._issues[key] = record
        self._dirty_issue(record, status_scope)
        return record
async def find_obligation(self,
*,
account_id: str,
consumer_id: str,
delivery_config_id: str,
delivery_config_version: int,
period_start: datetime,
period_end: datetime) ‑> ReportingObligationRecord | None
Expand source code
async def find_obligation(
    self,
    *,
    account_id: str,
    consumer_id: str,
    delivery_config_id: str,
    delivery_config_version: int,
    period_start: datetime,
    period_end: datetime,
) -> ReportingObligationRecord | None:
    key = (
        ReportingConfigurationGenerationKey(
            account_id=account_id,
            consumer_id=consumer_id,
            delivery_config_id=delivery_config_id,
            delivery_config_version=delivery_config_version,
        ),
        _utc(period_start).isoformat(),
        _utc(period_end).isoformat(),
    )
    found = self._obligation_by_period.get(key)
    return (
        await self.get_obligation(account_id=account_id, reporting_obligation_id=found)
        if found
        else None
    )
async def get_issue(self, *, issue_key: str, account_id: str) ‑> ReportingIssueLifecycle | None
Expand source code
async def get_issue(self, *, issue_key: str, account_id: str) -> ReportingIssueLifecycle | None:
    async with self._lock:
        # Live occurrences only, matching the Protocol docstring and the
        # Postgres store. Returning a retired row here would let the
        # projection re-emit a resolved issue's opened_at.
        live = self._issues.get((account_id, issue_key))
        return live if live is not None and live.live else None
async def get_obligation(self, *, account_id: str, reporting_obligation_id: str) ‑> ReportingObligationRecord | None
Expand source code
async def get_obligation(
    self, *, account_id: str, reporting_obligation_id: str
) -> ReportingObligationRecord | None:
    found = self._obligations.get(reporting_obligation_id)
    return found if found and found.account_id == account_id else None
async def get_provisional_acquisition(self, *, account_id: str, reporting_obligation_id: str, ordinal: int) ‑> ProvisionalAcquisition | None
Expand source code
async def get_provisional_acquisition(
    self, *, account_id: str, reporting_obligation_id: str, ordinal: int
) -> ProvisionalAcquisition | None:
    return self._provisional_acquisitions.get((account_id, reporting_obligation_id, ordinal))
async def get_provisional_observation(self, *, account_id: str, reporting_obligation_id: str) ‑> ProvisionalObservation | None
Expand source code
async def get_provisional_observation(
    self, *, account_id: str, reporting_obligation_id: str
) -> ProvisionalObservation | None:
    observations = [
        value
        for (account, obligation, _), value in self._provisional_observations.items()
        if account == account_id and obligation == reporting_obligation_id
    ]
    return max(observations, key=lambda item: item.acquisition.ordinal, default=None)
async def get_restatement_checkpoint(self, *, account_id: str, reporting_obligation_id: str) ‑> RestatementCheckpoint | None
Expand source code
async def get_restatement_checkpoint(
    self, *, account_id: str, reporting_obligation_id: str
) -> RestatementCheckpoint | None:
    checkpoint = self._restatement_checkpoints.get(reporting_obligation_id)
    return checkpoint if checkpoint and checkpoint.account_id == account_id else None
async def get_retry_schedule(self, *, scope_key: str) ‑> RetryScheduleEntry | None
Expand source code
async def get_retry_schedule(self, *, scope_key: str) -> RetryScheduleEntry | None:
    return self._retry_schedules.get(scope_key)
async def get_revision(self, *, account_id: str, reporting_revision_id: str) ‑> ReportingRevisionRecord | None
Expand source code
async def get_revision(
    self, *, account_id: str, reporting_revision_id: str
) -> ReportingRevisionRecord | None:
    found = self._revisions.get(reporting_revision_id)
    return found if found and found.account_id == account_id else None
async def lease_period_close(self, *, worker_id: str, now: datetime, lease_seconds: float) ‑> LeasedConfiguration | None
Expand source code
async def lease_period_close(
    self, *, worker_id: str, now: datetime, lease_seconds: float
) -> LeasedConfiguration | None:
    from datetime import timedelta

    async with self._lock:
        moment = _utc(now)
        # Rank leasable generations exactly the way the SQL store's
        # `ORDER BY lease_turn, lease_expires_at NULLS FIRST, ...` does:
        # whichever generation went longest without a turn goes first.
        #
        # The turn has to be the *primary* term. The caller's `now` has
        # already excluded every live lease, so among the survivors the
        # expiry carries no fairness information -- and preferring unheld
        # over expired ahead of the turn starves a crashed generation
        # forever: a peer that is leased and released every turn is always
        # unheld, so it wins every comparison while the generation whose
        # worker died stays expired and never closes another period.
        # Expiry and the generation key only break exact turn ties, so the
        # order stays total and never depends on physical layout.
        ranked: list[
            tuple[
                tuple[int, int, float, str, str, str, int],
                ReportingConfigurationGenerationKey,
            ]
        ] = []
        for key in self._configurations:
            if not self._configuration_lease_eligible(self._configurations[key]):
                continue
            turn = self._lease_turns.get(key, 0)
            held = self._leases.get(key)
            tail = (
                key.account_id,
                key.consumer_id,
                key.delivery_config_id,
                key.delivery_config_version,
            )
            if held is None:
                ranked.append(((turn, 0, 0.0, *tail), key))
            elif _utc(held[1]) <= moment:
                ranked.append(((turn, 1, _utc(held[1]).timestamp(), *tail), key))
        if not ranked:
            return None
        # The generation key is part of the rank, so the order is total and
        # both stores make the same choice. Ranking by `min` alone would
        # fall back to whichever generation this process happened to accept
        # first, which the SQL store cannot reproduce and no adopter can
        # observe consistently across a restart or a second worker.
        key = min(ranked, key=lambda item: item[0])[1]
        configuration = self._configurations[key]
        expires = moment + timedelta(seconds=lease_seconds)
        self._lease_turn += 1
        self._lease_turns[key] = self._lease_turn
        self._leases[key] = (worker_id, expires)
        return LeasedConfiguration(
            account_id=configuration.account_id,
            consumer_id=configuration.consumer_id,
            delivery_config_id=configuration.delivery_config_id,
            delivery_config_version=configuration.delivery_config_version,
            lease_expires_at=expires,
        )
async def list_adjustments(self, *, account_id: str, reporting_revision_ids: Sequence[str]) ‑> tuple[ReportingAdjustmentRecord, ...]
Expand source code
async def list_adjustments(
    self, *, account_id: str, reporting_revision_ids: Sequence[str]
) -> tuple[ReportingAdjustmentRecord, ...]:
    wanted = set(reporting_revision_ids)
    return tuple(
        item
        for item in self._adjustments.values()
        if item.account_id == account_id and item.adjusts_reporting_revision_id in wanted
    )
async def list_all_configurations(self) ‑> tuple[ReportingConfiguration, ...]
Expand source code
async def list_all_configurations(self) -> tuple[ReportingConfiguration, ...]:
    return tuple(
        sorted(
            (c for c in self._configurations.values() if not c.quarantined),
            key=lambda item: (
                item.account_id,
                item.consumer_id,
                item.delivery_config_id,
                item.delivery_config_version,
            ),
        )
    )
async def list_configurations(self,
*,
caller: ReportingCaller,
delivery_config_ids: Sequence[str] | None = None) ‑> tuple[ReportingConfiguration, ...]
Expand source code
async def list_configurations(
    self, *, caller: ReportingCaller, delivery_config_ids: Sequence[str] | None = None
) -> tuple[ReportingConfiguration, ...]:
    wanted = set(delivery_config_ids) if delivery_config_ids else None
    return tuple(
        configuration
        for configuration in self._configurations.values()
        if configuration.account_id == caller.account_id
        and configuration.consumer_id == caller.consumer_id
        and (wanted is None or configuration.delivery_config_id in wanted)
    )
async def list_consumer_statuses(self,
*,
account_id: str,
consumer_id: str,
reporting_obligation_ids: Sequence[str] | None = None) ‑> tuple[ConsumerStatusRecord, ...]
Expand source code
async def list_consumer_statuses(
    self,
    *,
    account_id: str,
    consumer_id: str,
    reporting_obligation_ids: Sequence[str] | None = None,
) -> tuple[ConsumerStatusRecord, ...]:
    wanted = set(reporting_obligation_ids) if reporting_obligation_ids is not None else None
    result = []
    for item in self._statuses.values():
        if item.account_id != account_id or item.consumer_id != consumer_id:
            continue
        if wanted is not None and not self._matches_obligations(item, wanted):
            continue
        result.append(item)
    return tuple(result)
async def list_revisions(self, *, account_id: str, reporting_obligation_id: str) ‑> tuple[ReportingRevisionRecord, ...]
Expand source code
async def list_revisions(
    self, *, account_id: str, reporting_obligation_id: str
) -> tuple[ReportingRevisionRecord, ...]:
    return tuple(
        item
        for item in self._revisions.values()
        if item.account_id == account_id
        and item.reporting_obligation_id == reporting_obligation_id
    )
async def open_snapshot(self, *, caller: ReportingCaller, filters_fingerprint: str) ‑> LedgerSnapshot
Expand source code
async def open_snapshot(
    self, *, caller: ReportingCaller, filters_fingerprint: str
) -> LedgerSnapshot:
    account_id = caller.account_id
    async with self._lock:
        as_of = _utc(self._clock())
        max_sequence = max(
            (
                sequence
                for sequence, account, kind, record_id, _ in self._changes
                if account == account_id and self._change_owners[sequence] == caller.consumer_id
            ),
            default=0,
        )
        return LedgerSnapshot(
            snapshot_id="rpls_"
            + _fingerprint([account_id, caller.consumer_id, filters_fingerprint, max_sequence])[
                :32
            ],
            account_id=account_id,
            consumer_id=caller.consumer_id,
            ledger_as_of=as_of,
            max_sequence=max_sequence,
        )
async def put_configuration(self, configuration: ReportingConfiguration) ‑> None
Expand source code
async def put_configuration(self, configuration: ReportingConfiguration) -> None:
    reject_reserved_authoritative_party(configuration)
    async with self._mutation():
        key = configuration.generation_key
        existing = self._configurations.get(key)
        if existing is not None and existing.quarantined:
            raise LedgerConflictError(
                "REPORTING_GENERATION_QUARANTINED", "operator reconciliation is required"
            )
        if existing is not None and _fingerprint(_config_payload(existing)) != _fingerprint(
            _config_payload(configuration)
        ):
            raise LedgerConflictError(
                "CONFIGURATION_GENERATION_IMMUTABLE",
                f"configuration {key.delivery_config_id}@{key.delivery_config_version} "
                "already exists with different content for this account; "
                "publish a new version instead of editing a retained generation",
            )
        self._configurations[key] = configuration
        changed = existing is None or configuration_lifecycle(
            existing
        ) != configuration_lifecycle(configuration)
        if changed and self._notification_state is not None:
            self._dirty_status(
                ReportingStatusScope(configuration.account_id, configuration.generation_key),
                "configuration",
                before=configuration_evidence(existing) if existing is not None else None,
                after=configuration_evidence(configuration),
            )
async def read_page(self,
*,
snapshot: LedgerSnapshot,
consumer_id: str | None,
delivery_config_ids: Sequence[str] | None,
media_buy_ids: Sequence[str] | None,
offset: int,
limit: int,
changes_after_sequence: int | None,
feed_purposes: Sequence[str] | None = None,
period_start: datetime | None = None,
period_end: datetime | None = None) ‑> LedgerPage
Expand source code
async def read_page(
    self,
    *,
    snapshot: LedgerSnapshot,
    consumer_id: str | None,
    delivery_config_ids: Sequence[str] | None,
    media_buy_ids: Sequence[str] | None,
    offset: int,
    limit: int,
    changes_after_sequence: int | None,
    feed_purposes: Sequence[str] | None = None,
    period_start: datetime | None = None,
    period_end: datetime | None = None,
) -> LedgerPage:
    if consumer_id not in {None, snapshot.consumer_id}:
        raise LedgerConflictError(
            "CURSOR_SNAPSHOT_MISMATCH", "snapshot belongs to another caller"
        )
    consumer_id = snapshot.consumer_id
    lower = changes_after_sequence or 0
    selected = [
        item
        for item in self._changes
        if item[1] == snapshot.account_id
        and self._change_owners[item[0]] == consumer_id
        and lower < item[0] <= snapshot.max_sequence
    ]
    config_filter = set(delivery_config_ids) if delivery_config_ids else None
    media_buy_filter = set(media_buy_ids) if media_buy_ids else None

    records: list[tuple[int, LedgerRecordKind, Any]] = []
    for sequence, _account, kind, record_id, _committed in sorted(selected):
        record = self._resolve(kind, record_id, snapshot.account_id, consumer_id)
        if record is None or record.account_id != snapshot.account_id:
            continue
        if not self._in_scope(
            kind,
            record,
            config_filter,
            media_buy_filter,
            consumer_id,
            feed_purposes,
            period_start,
            period_end,
        ):
            continue
        records.append((sequence, kind, record))

    window = records[offset : offset + limit]
    has_more = offset + limit < len(records)
    return LedgerPage(
        obligations=tuple(item[2] for item in window if item[1] == "obligation"),
        revisions=tuple(item[2] for item in window if item[1] == "revision"),
        adjustments=tuple(item[2] for item in window if item[1] == "adjustment"),
        consumer_statuses=tuple(item[2] for item in window if item[1] == "consumer_status"),
        total_count=len(records),
        has_more=has_more,
        cursor=(
            encode_cursor({"snapshot": snapshot.snapshot_id, "offset": offset + limit})
            if has_more
            else None
        ),
    )
async def read_revision_rows(self,
*,
account_id: str,
reporting_revision_id: str,
cursor: str | None = None,
limit: int = 500) ‑> ReportingRowPage
Expand source code
async def read_revision_rows(
    self,
    *,
    account_id: str,
    reporting_revision_id: str,
    cursor: str | None = None,
    limit: int = 500,
) -> ReportingRowPage:
    revision = await self.get_revision(
        account_id=account_id, reporting_revision_id=reporting_revision_id
    )
    if revision is None:
        raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account")
    rows = self._rows.get(reporting_revision_id, ())
    offset = revision_row_offset(cursor, reporting_revision_id, limit)
    window = rows[offset : offset + limit]
    has_more = offset + limit < len(rows)
    return ReportingRowPage(
        reporting_revision_id=reporting_revision_id,
        rows=tuple(deepcopy(row) for row in window),
        total_count=len(rows),
        has_more=has_more,
        cursor=(
            encode_cursor(
                {"ownership": 2, "revision": reporting_revision_id, "offset": offset + limit}
            )
            if has_more
            else None
        ),
    )
async def read_status_snapshot(self, *, caller: ReportingCaller) ‑> ReportingStatusSnapshot
Expand source code
async def read_status_snapshot(self, *, caller: ReportingCaller) -> ReportingStatusSnapshot:
    from adcp.reporting.ledger.status_snapshot import settle_memory_snapshot

    async with self._mutation():
        from adcp.reporting.ledger.delivery_models import ReportingDeliveryPrincipal
        from adcp.reporting.materializer.capture import private_snapshot

        return private_snapshot(
            settle_memory_snapshot(self, caller.account_id),
            ReportingDeliveryPrincipal(caller.account_id, caller.consumer_id),
        )
async def record_consumer_status(self, status: ConsumerStatusRecord) ‑> tuple[ConsumerStatusRecord, bool]
Expand source code
async def record_consumer_status(
    self, status: ConsumerStatusRecord
) -> tuple[ConsumerStatusRecord, bool]:
    from adcp.reporting.evidence import consumer_reference

    consumer_reference(status.consumer_id)
    async with self._mutation():
        identity = _consumer_status_identity(status)
        existing = self._replay(status)
        if existing is not None:
            return existing, False
        from adcp.reporting.ledger.status_snapshot import (
            memory_snapshot,
            validate_status_evidence,
        )

        validate_status_evidence(status, memory_snapshot(self, status.account_id))
        leaf = self._current_status_leaf(status.chain_key)
        if status.supersedes_reporting_status_id:
            if (
                leaf is None
                or leaf.reporting_status_id != status.supersedes_reporting_status_id
            ):
                raise LedgerConflictError(
                    "STATUS_SUPERSEDES_STALE",
                    "supersedes_reporting_status_id must name this chain's current leaf; "
                    "a stale pointer would let a successful retry erase a recorded outage",
                )
            self._statuses[(leaf.account_id, leaf.consumer_id, leaf.reporting_status_id)] = (
                replace(leaf, superseded=True)
            )
        elif leaf is not None:
            raise LedgerConflictError(
                "STATUS_SUPERSEDES_REQUIRED",
                "this chain already has a current statement; a new statement must "
                "explicitly supersede it",
            )
        key = (status.account_id, status.consumer_id, status.reporting_status_id)
        self._statuses[key] = status
        self._status_identity[key] = identity
        self._append(
            status.account_id,
            "consumer_status",
            status.reporting_status_id,
            consumer_id=status.consumer_id,
        )
        from adcp.reporting.ledger.status_snapshot import settle_memory_snapshot

        settle_memory_snapshot(self, status.account_id)
        if self._notification_state is not None:
            self._dirty_status(
                ReportingStatusScope(
                    status.account_id,
                    status.generation_key,
                    status.reporting_obligation_id,
                    status.consumer_id,
                ),
                "consumer_status",
                (
                    ReportingStatusEvidence("consumer_status", leaf.reporting_status_id)
                    if leaf is not None
                    else None
                ),
                ReportingStatusEvidence(
                    "consumer_status",
                    status.reporting_status_id,
                    supersedes_id=status.supersedes_reporting_status_id,
                ),
            )
        return status, True
async def record_consumer_status_with_lifecycle(self, status: ConsumerStatusRecord) ‑> tuple[ConsumerStatusRecord, bool]
Expand source code
async def record_consumer_status_with_lifecycle(
    self, status: ConsumerStatusRecord
) -> tuple[ConsumerStatusRecord, bool]:
    return await self.record_consumer_status(status)
async def record_restatement_checkpoint(self,
checkpoint: RestatementCheckpoint) ‑> RestatementCheckpoint
Expand source code
async def record_restatement_checkpoint(
    self, checkpoint: RestatementCheckpoint
) -> RestatementCheckpoint:
    async with self._mutation():
        obligation = self._obligations.get(checkpoint.reporting_obligation_id)
        if obligation is None or obligation.account_id != checkpoint.account_id:
            raise LedgerConflictError(
                "OBLIGATION_NOT_FOUND",
                "a restatement checkpoint must attach to an obligation for this account",
            )
        existing = self._restatement_checkpoints.get(checkpoint.reporting_obligation_id)
        if existing is not None:
            if checkpoint.next_observation < existing.next_observation:
                return existing
            if checkpoint.next_observation == existing.next_observation and _utc(
                checkpoint.checked_at
            ) <= _utc(existing.checked_at):
                return existing
        self._restatement_checkpoints[checkpoint.reporting_obligation_id] = checkpoint
        return checkpoint
async def record_retry_schedule(self,
entry: RetryScheduleEntry) ‑> RetryScheduleEntry
Expand source code
async def record_retry_schedule(self, entry: RetryScheduleEntry) -> RetryScheduleEntry:
    async with self._mutation():
        existing = self._retry_schedules.get(entry.scope_key)
        if existing is not None:
            assert entry.recorded_at is not None and existing.recorded_at is not None
            if _utc(entry.recorded_at) < _utc(existing.recorded_at):
                return existing
            if entry.attempt > 0 and not entry.blocked:
                if existing.blocked and not entry.replayed:
                    return existing
                if entry.attempt < existing.attempt and not entry.replayed:
                    return existing
                if not entry.replayed and not existing.blocked:
                    entry = replace(
                        entry,
                        retry_not_before=max(
                            entry.retry_not_before, existing.retry_not_before, key=_utc
                        ),
                    )
        stored = replace(entry, replayed=False)
        self._retry_schedules[entry.scope_key] = stored
        return stored
async def release_period_close(self,
lease: LeasedConfiguration,
*,
worker_id: str) ‑> None
Expand source code
async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) -> None:
    async with self._lock:
        key = lease.generation_key
        held = self._leases.get(key)
        if held == (worker_id, _utc(lease.lease_expires_at)):
            del self._leases[key]
async def reserve_provisional_acquisition(self, acquisition: ProvisionalAcquisition) ‑> ProvisionalAcquisition
Expand source code
async def reserve_provisional_acquisition(
    self, acquisition: ProvisionalAcquisition
) -> ProvisionalAcquisition:
    # Canonical round-trip also detaches all request collections.
    acquisition = ProvisionalAcquisition.from_wire(acquisition.to_wire())
    key = (acquisition.account_id, acquisition.obligation_id, acquisition.ordinal)
    async with self._mutation():
        existing = self._provisional_acquisitions.get(key)
        if existing is not None:
            return existing
        obligation = self._obligations.get(acquisition.obligation_id)
        if obligation is None or obligation.account_id != acquisition.account_id:
            raise LedgerConflictError("OBLIGATION_NOT_FOUND", "unknown observation obligation")
        if not acquisition.binds(obligation):
            raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition generation differs")
        if any(
            item.account_id == acquisition.account_id
            and item.execution_key == acquisition.execution_key
            for item in self._provisional_acquisitions.values()
        ):
            raise LedgerConflictError(
                "OBSERVATION_CONFLICT", "execution key is already reserved"
            )
        checkpoint = self._restatement_checkpoints.get(acquisition.obligation_id)
        expected = (
            checkpoint.next_observation
            if checkpoint
            else len(
                await self.list_revisions(
                    account_id=acquisition.account_id,
                    reporting_obligation_id=acquisition.obligation_id,
                )
            )
        )
        if acquisition.ordinal != expected:
            raise LedgerConflictError("OBSERVATION_CONFLICT", "observation ordinal changed")
        self._provisional_acquisitions[key] = acquisition
        return acquisition
async def resolve_consumer_status_replay(self, status: ConsumerStatusRecord) ‑> ConsumerStatusRecord | None
Expand source code
async def resolve_consumer_status_replay(
    self, status: ConsumerStatusRecord
) -> ConsumerStatusRecord | None:
    async with self._lock:
        return self._replay(status)
async def retire_issue(self,
*,
issue_key: str,
account_id: str,
at: datetime,
status_scope: ReportingStatusScope | None = None) ‑> ReportingIssueLifecycle | None
Expand source code
async def retire_issue(
    self,
    *,
    issue_key: str,
    account_id: str,
    at: datetime,
    status_scope: ReportingStatusScope | None = None,
) -> ReportingIssueLifecycle | None:
    async with self._mutation():
        key = (account_id, issue_key)
        live = self._issues.get(key)
        if live is None or not live.live or not issue_is_retirable(live.issue_state):
            # Convergent: nothing live, or already waived. A waived issue
            # is retired from the projection by agreement, and overwriting
            # that readable act with `resolved` is an edge the forward-only
            # lifecycle forbids -- enforced here rather than left to each
            # caller to remember.
            return None
        check_issue_state_transition(live.issue_state, "resolved")
        retired = replace(live, issue_state="resolved", retired_at=_utc(at))
        self._resolve_issue_scope(
            account_id=account_id,
            consumer_id=live.consumer_id,
            issue_id=live.issue_id,
            status_scope=status_scope,
        )
        self._issues[key] = retired
        self._dirty_issue(retired, status_scope, live)
        return retired
async def set_issue_state(self,
*,
issue_key: str,
account_id: str,
state: "Literal['acknowledged', 'waived']",
at: datetime,
external_ref: str | None = None,
status_scope: ReportingStatusScope | None = None) ‑> ReportingIssueLifecycle
Expand source code
async def set_issue_state(
    self,
    *,
    issue_key: str,
    account_id: str,
    state: Literal["acknowledged", "waived"],
    at: datetime,
    external_ref: str | None = None,
    status_scope: ReportingStatusScope | None = None,
) -> ReportingIssueLifecycle:
    async with self._mutation():
        if state not in {"acknowledged", "waived"}:
            raise LedgerConflictError(
                "ISSUE_STATE_NOT_OPERATOR_SETTABLE",
                f"issue_state {state!r} is not settable by an operator. 'resolved' is "
                "reachable only when the condition actually clears -- the projection "
                "retires it -- because a seller must not retire a mismatch out of a "
                "degraded projection while the statement that caused it is still the "
                "consumer's current leaf",
            )
        key = (account_id, issue_key)
        live = self._issues.get(key)
        if live is None or not live.live:
            raise LedgerConflictError(
                "ISSUE_NOT_OPEN",
                f"no open issue {issue_key!r} for this account; a retired issue cannot be "
                "reopened, and a recurrence gets a new occurrence",
            )
        check_issue_state_transition(live.issue_state, state)
        if state == "waived" and live.issue_state != "waived":
            from adcp.reporting.ledger.status_projection import bind_mismatch_waiver
            from adcp.reporting.ledger.status_snapshot import memory_snapshot

            live = bind_mismatch_waiver(memory_snapshot(self, account_id), live)
        updated = replace(
            live,
            issue_state=state,
            external_ref=external_ref or live.external_ref,
            # Set on the way into a retired state and never cleared.
            retired_at=(
                (
                    live.retired_at
                    if live.waived_reporting_status_id is not None
                    and live.retired_at is not None
                    else _utc(at)
                )
                if state == "waived"
                else live.retired_at
            ),
        )
        self._resolve_issue_scope(
            account_id=account_id,
            consumer_id=live.consumer_id,
            issue_id=live.issue_id,
            status_scope=status_scope,
        )
        self._issues[key] = updated
        # Derive the no-op from the resulting record rather than predicting
        # it: an idempotent re-acknowledge changes nothing and enqueues
        # nothing, while anything that does move retained evidence stays
        # reconstructable for the projector.
        self._dirty_issue(updated, status_scope, live)
        return updated
async def set_revision_readable(self, *, account_id: str, reporting_revision_id: str, readable: bool) ‑> None
Expand source code
async def set_revision_readable(
    self, *, account_id: str, reporting_revision_id: str, readable: bool
) -> None:
    async with self._mutation():
        existing = self._revisions.get(reporting_revision_id)
        if existing is None or existing.account_id != account_id:
            raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account")
        if existing.readable == readable:
            return
        self._revisions[reporting_revision_id] = replace(existing, readable=readable)
        if self._notification_state is not None:
            obligation = self._obligations[existing.reporting_obligation_id]
            self._dirty_status(
                ReportingStatusScope.for_obligation(obligation),
                "readability",
                ReportingStatusEvidence(
                    "revision", reporting_revision_id, readable=existing.readable
                ),
                ReportingStatusEvidence("revision", reporting_revision_id, readable=readable),
            )
async def transaction(self) ‑> AsyncIterator[InMemoryReportingLedgerStore]
Expand source code
@asynccontextmanager
async def transaction(self) -> AsyncIterator[InMemoryReportingLedgerStore]:
    """Group several source mutations into one atomic, replayable boundary."""
    async with self._mutation():
        yield self

Group several source mutations into one atomic, replayable boundary.

class LeasedConfiguration (account_id: str,
consumer_id: str,
delivery_config_id: str,
delivery_config_version: int,
lease_expires_at: datetime)
Expand source code
@dataclass(frozen=True)
class LeasedConfiguration:
    """One configuration generation leased for period-close work.

    Carries its account so the producer never has to invert the lease back into
    a tenant -- an inversion that is easy to get subtly wrong when two accounts
    reuse a ``delivery_config_id``.
    """

    account_id: str
    consumer_id: str
    delivery_config_id: str
    delivery_config_version: int
    lease_expires_at: datetime

    @property
    def generation_key(self) -> ReportingConfigurationGenerationKey:
        return ReportingConfigurationGenerationKey(
            account_id=self.account_id,
            consumer_id=self.consumer_id,
            delivery_config_id=self.delivery_config_id,
            delivery_config_version=self.delivery_config_version,
        )

One configuration generation leased for period-close work.

Carries its account so the producer never has to invert the lease back into a tenant – an inversion that is easy to get subtly wrong when two accounts reuse a delivery_config_id.

Instance variables

var account_id : str
var consumer_id : str
var delivery_config_id : str
var delivery_config_version : int
prop generation_key : ReportingConfigurationGenerationKey
Expand source code
@property
def generation_key(self) -> ReportingConfigurationGenerationKey:
    return ReportingConfigurationGenerationKey(
        account_id=self.account_id,
        consumer_id=self.consumer_id,
        delivery_config_id=self.delivery_config_id,
        delivery_config_version=self.delivery_config_version,
    )
var lease_expires_at : datetime.datetime
class LedgerConflictError (code: str, message: str)
Expand source code
class LedgerConflictError(RuntimeError):
    """A write would violate an invariant that makes the ledger evidence.

    Carries a stable ``code`` so a handler can map it to a wire error without
    matching on prose.
    """

    def __init__(self, code: str, message: str) -> None:
        super().__init__(message)
        self.code = code

A write would violate an invariant that makes the ledger evidence.

Carries a stable code so a handler can map it to a wire error without matching on prose.

Ancestors

  • builtins.RuntimeError
  • builtins.Exception
  • builtins.BaseException
class LedgerPage (obligations: tuple[ReportingObligationRecord, ...],
revisions: tuple[ReportingRevisionRecord, ...],
adjustments: tuple[ReportingAdjustmentRecord, ...],
consumer_statuses: tuple[ConsumerStatusRecord, ...],
total_count: int,
has_more: bool,
cursor: str | None)
Expand source code
@dataclass(frozen=True)
class LedgerPage:
    """One page of the flat union of ledger records.

    Pagination is over the *flat union* of obligations, revisions, adjustments
    and consumer statuses rather than nesting history under each obligation.
    Nesting looks friendlier and is unbounded: one obligation with ten thousand
    snapshot restatements would make a single "page" unpageable.
    """

    obligations: tuple[ReportingObligationRecord, ...]
    revisions: tuple[ReportingRevisionRecord, ...]
    adjustments: tuple[ReportingAdjustmentRecord, ...]
    consumer_statuses: tuple[ConsumerStatusRecord, ...]
    total_count: int
    has_more: bool
    cursor: str | None

One page of the flat union of ledger records.

Pagination is over the flat union of obligations, revisions, adjustments and consumer statuses rather than nesting history under each obligation. Nesting looks friendlier and is unbounded: one obligation with ten thousand snapshot restatements would make a single "page" unpageable.

Instance variables

var adjustments : tuple[ReportingAdjustmentRecord, ...]
var consumer_statuses : tuple[ConsumerStatusRecord, ...]
var cursor : str | None
var has_more : bool
var obligations : tuple[ReportingObligationRecord, ...]
var revisions : tuple[ReportingRevisionRecord, ...]
var total_count : int
class ReportingLedgerStore (*args, **kwargs)
Expand source code
@runtime_checkable
class ReportingLedgerStore(Protocol):
    """Durable home for obligations, revisions, adjustments, and statuses."""

    async def create_schema(self) -> None:
        """Idempotently create or upgrade this store's schema. Safe on every boot."""
        ...

    # -- configurations --------------------------------------------------

    async def put_configuration(self, configuration: ReportingConfiguration) -> None:
        """Record an accepted configuration generation.

        A generation is immutable: re-putting one with changed content is a
        conflict, because a buyer that retained it derives expectations from it.
        """
        ...

    async def list_configurations(
        self, *, caller: ReportingCaller, delivery_config_ids: Sequence[str] | None = None
    ) -> tuple[ReportingConfiguration, ...]: ...

    async def list_all_configurations(self) -> tuple[ReportingConfiguration, ...]:
        """Enumerate retained generations for trusted service startup recovery.

        This administrative API is never exposed through buyer task handlers;
        those continue to use account-scoped reads after authorization.
        """
        raise NotImplementedError

    # -- obligations -----------------------------------------------------

    async def commit_obligation(
        self, obligation: ReportingObligationRecord
    ) -> ReportingObligationRecord:
        """Commit an obligation, or return the existing one for its period.

        Idempotent by ``(account, config generation, period)``, not by id: two
        workers racing a period close must converge on one obligation.
        New records require an explicit, validated currency. A legacy record
        with unknown currency can be read/replayed but cannot be filled here.
        """
        ...

    async def get_obligation(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> ReportingObligationRecord | None: ...

    async def find_obligation(
        self,
        *,
        account_id: str,
        consumer_id: str,
        delivery_config_id: str,
        delivery_config_version: int,
        period_start: datetime,
        period_end: datetime,
    ) -> ReportingObligationRecord | None: ...

    # -- revisions -------------------------------------------------------

    async def commit_revision(
        self,
        revision: ReportingRevisionRecord,
        rows: Sequence[dict[str, Any]],
    ) -> ReportingRevisionRecord:
        """Commit an immutable revision and its frozen rows.

        Enforces terminal officials, exact supersession, and idempotent replay
        of an identical ``reporting_revision_id``.
        """
        ...

    async def list_revisions(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> tuple[ReportingRevisionRecord, ...]: ...

    async def get_revision(
        self, *, account_id: str, reporting_revision_id: str
    ) -> ReportingRevisionRecord | None: ...

    async def read_revision_rows(
        self,
        *,
        account_id: str,
        reporting_revision_id: str,
        cursor: str | None = None,
        limit: int = 500,
    ) -> ReportingRowPage:
        """Walk one revision's frozen rows in stable order."""
        ...

    async def set_revision_readable(
        self, *, account_id: str, reporting_revision_id: str, readable: bool
    ) -> None:
        """Record that a revision's content is (no longer) readable.

        Retention expiry and storage loss are real; representing them is how an
        obligation becomes honestly ``action_required`` instead of staying
        ``complete`` over evidence nobody can read.
        """
        ...

    # -- adjustments -----------------------------------------------------

    async def commit_adjustment(
        self, adjustment: ReportingAdjustmentRecord
    ) -> ReportingAdjustmentRecord: ...

    async def list_adjustments(
        self, *, account_id: str, reporting_revision_ids: Sequence[str]
    ) -> tuple[ReportingAdjustmentRecord, ...]: ...

    # -- consumer status (preview) ---------------------------------------

    async def record_consumer_status(
        self, status: ConsumerStatusRecord
    ) -> tuple[ConsumerStatusRecord, bool]:
        """Append a consumer status statement, superseding the chain's leaf.

        Returns ``(record, recorded)`` where ``recorded`` is ``False`` for an
        exact idempotent replay.  Supersession is atomic: naming a stale or
        missing leaf must raise rather than fork the chain, so a successful
        retry cannot erase a recorded outage.
        """
        ...

    async def resolve_consumer_status_replay(
        self, status: ConsumerStatusRecord
    ) -> ConsumerStatusRecord | None:
        """Return the stored row when this exact statement was already recorded.

        Exists so the ingest can answer "is this an exact retry?" *before* it
        validates anything time-dependent. ``record_consumer_status`` already
        makes that check, but it runs last, and some validation the ingest does
        first -- notably "``content_mismatch`` must name the revision the seller
        currently requires" -- has an answer that changes as the seller
        publishes. An exact retry sent after a restatement would then be
        rejected for a reason that did not apply when it was first accepted,
        and the spec's ``result_mapping`` requires ``unchanged``.

        Returns ``None`` when the id is unknown, so the caller proceeds to
        validate and record. Raises :class:`LedgerConflictError` with
        ``STATUS_IDENTITY_CONFLICT`` when the id exists with different content:
        reuse with changed content is still a conflict, and answering
        ``unchanged`` there would let a buyer rewrite a recorded statement.
        """
        ...

    async def list_consumer_statuses(
        self,
        *,
        account_id: str,
        consumer_id: str,
        reporting_obligation_ids: Sequence[str] | None = None,
    ) -> tuple[ConsumerStatusRecord, ...]: ...

    # -- issue lifecycle -------------------------------------------------

    async def ensure_issue_opened(
        self,
        *,
        issue_key: str,
        account_id: str,
        consumer_id: str | None,
        observed_at: datetime,
    ) -> ReportingIssueLifecycle:
        """Return this condition's live occurrence, opening one if needed.

        Idempotent: the second caller gets the first caller's ``opened_at``.
        That is the whole point -- AdCP 3.2.0-rc.3 anchors the escalation clock
        to ``opened_at``, so a re-emission that advanced it would let an
        unattended mismatch stay below ``action_required`` indefinitely.

        This is the one write the read path performs. ``get_reporting_status``
        is where a consumer mismatch is *first observed*, and a stable
        first-observation timestamp cannot be derived from immutable evidence
        alone. Implementations must make it safe under concurrent reads; two
        readers of one condition must converge on one row rather than open two
        occurrences.

        A condition retired as ``resolved`` or ``waived`` opens a *new*
        occurrence with the next ``generation``, a new ``issue_id``, and a new
        ``opened_at``, which is exactly the spec's "a recurrence after
        retirement receives a new issue_id".
        """
        ...

    async def set_issue_state(
        self,
        *,
        issue_key: str,
        account_id: str,
        state: Literal["acknowledged", "waived"],
        at: datetime,
        external_ref: str | None = None,
    ) -> ReportingIssueLifecycle:
        """Move a live occurrence forward, or attach an ``external_ref``.

        Deliberately cannot set ``resolved``. The spec forbids retiring a
        ``CONSUMER_STATUS_MISMATCH`` out of a degraded projection while the
        statement that caused it is still the consumer's current leaf, so
        resolution is not an operator action -- it is what the projection does
        when the condition actually clears (see :meth:`retire_issue`).
        Otherwise a seller could unilaterally erase a buyer-attributed
        disagreement, which is the one outcome this separately attributed loop
        exists to prevent.

        ``waived`` requires explicit off-protocol agreement by this consumer
        and seller for the exact issue, causing statement and diagnosed
        conflict. The adopter must retain its private consent audit before
        calling this method; ``external_ref`` is inert correlation, not proof.
        Rc.6 removes that issue and restores underlying seller health, without
        changing the consumer statement. Later statements or different
        conflicts are evaluated independently, never covered by this waiver.
        """
        ...

    async def retire_issue(
        self, *, issue_key: str, account_id: str, at: datetime
    ) -> ReportingIssueLifecycle | None:
        """Mark this condition's live occurrence ``resolved``.

        Called by the projection when the condition no longer holds. Returns
        ``None`` when there was nothing live to retire, so a repeated
        projection is convergent rather than an error.
        """
        ...

    async def get_issue(self, *, issue_key: str, account_id: str) -> ReportingIssueLifecycle | None:
        """The live occurrence of this condition, or ``None``."""
        ...

    # -- snapshots and pagination ----------------------------------------

    async def open_snapshot(
        self, *, caller: ReportingCaller, filters_fingerprint: str
    ) -> LedgerSnapshot:
        """Take a consistent read boundary for one account and filter set."""
        ...

    async def read_page(
        self,
        *,
        snapshot: LedgerSnapshot,
        consumer_id: str | None,
        delivery_config_ids: Sequence[str] | None,
        media_buy_ids: Sequence[str] | None,
        offset: int,
        limit: int,
        changes_after_sequence: int | None,
    ) -> LedgerPage: ...

    # -- worker leasing --------------------------------------------------

    async def lease_period_close(
        self, *, worker_id: str, now: datetime, lease_seconds: float
    ) -> LeasedConfiguration | None:
        """Lease one configuration generation for period-close work.

        Leasing rather than locking so a worker that dies mid-close releases its
        work by expiry instead of wedging the period forever, and so two workers
        cannot both close the same period.  Implementations should prefer the
        generation whose work is most overdue.
        """
        ...

    async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) -> None: ...

Durable home for obligations, revisions, adjustments, and statuses.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def commit_adjustment(self, adjustment: ReportingAdjustmentRecord) ‑> ReportingAdjustmentRecord
Expand source code
async def commit_adjustment(
    self, adjustment: ReportingAdjustmentRecord
) -> ReportingAdjustmentRecord: ...
async def commit_obligation(self, obligation: ReportingObligationRecord) ‑> ReportingObligationRecord
Expand source code
async def commit_obligation(
    self, obligation: ReportingObligationRecord
) -> ReportingObligationRecord:
    """Commit an obligation, or return the existing one for its period.

    Idempotent by ``(account, config generation, period)``, not by id: two
    workers racing a period close must converge on one obligation.
    New records require an explicit, validated currency. A legacy record
    with unknown currency can be read/replayed but cannot be filled here.
    """
    ...

Commit an obligation, or return the existing one for its period.

Idempotent by (account, config generation, period), not by id: two workers racing a period close must converge on one obligation. New records require an explicit, validated currency. A legacy record with unknown currency can be read/replayed but cannot be filled here.

async def commit_revision(self, revision: ReportingRevisionRecord, rows: Sequence[dict[str, Any]]) ‑> ReportingRevisionRecord
Expand source code
async def commit_revision(
    self,
    revision: ReportingRevisionRecord,
    rows: Sequence[dict[str, Any]],
) -> ReportingRevisionRecord:
    """Commit an immutable revision and its frozen rows.

    Enforces terminal officials, exact supersession, and idempotent replay
    of an identical ``reporting_revision_id``.
    """
    ...

Commit an immutable revision and its frozen rows.

Enforces terminal officials, exact supersession, and idempotent replay of an identical reporting_revision_id.

async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    """Idempotently create or upgrade this store's schema. Safe on every boot."""
    ...

Idempotently create or upgrade this store's schema. Safe on every boot.

async def ensure_issue_opened(self,
*,
issue_key: str,
account_id: str,
consumer_id: str | None,
observed_at: datetime) ‑> ReportingIssueLifecycle
Expand source code
async def ensure_issue_opened(
    self,
    *,
    issue_key: str,
    account_id: str,
    consumer_id: str | None,
    observed_at: datetime,
) -> ReportingIssueLifecycle:
    """Return this condition's live occurrence, opening one if needed.

    Idempotent: the second caller gets the first caller's ``opened_at``.
    That is the whole point -- AdCP 3.2.0-rc.3 anchors the escalation clock
    to ``opened_at``, so a re-emission that advanced it would let an
    unattended mismatch stay below ``action_required`` indefinitely.

    This is the one write the read path performs. ``get_reporting_status``
    is where a consumer mismatch is *first observed*, and a stable
    first-observation timestamp cannot be derived from immutable evidence
    alone. Implementations must make it safe under concurrent reads; two
    readers of one condition must converge on one row rather than open two
    occurrences.

    A condition retired as ``resolved`` or ``waived`` opens a *new*
    occurrence with the next ``generation``, a new ``issue_id``, and a new
    ``opened_at``, which is exactly the spec's "a recurrence after
    retirement receives a new issue_id".
    """
    ...

Return this condition's live occurrence, opening one if needed.

Idempotent: the second caller gets the first caller's opened_at. That is the whole point – AdCP 3.2.0-rc.3 anchors the escalation clock to opened_at, so a re-emission that advanced it would let an unattended mismatch stay below action_required indefinitely.

This is the one write the read path performs. get_reporting_status is where a consumer mismatch is first observed, and a stable first-observation timestamp cannot be derived from immutable evidence alone. Implementations must make it safe under concurrent reads; two readers of one condition must converge on one row rather than open two occurrences.

A condition retired as resolved or waived opens a new occurrence with the next generation, a new issue_id, and a new opened_at, which is exactly the spec's "a recurrence after retirement receives a new issue_id".

async def find_obligation(self,
*,
account_id: str,
consumer_id: str,
delivery_config_id: str,
delivery_config_version: int,
period_start: datetime,
period_end: datetime) ‑> ReportingObligationRecord | None
Expand source code
async def find_obligation(
    self,
    *,
    account_id: str,
    consumer_id: str,
    delivery_config_id: str,
    delivery_config_version: int,
    period_start: datetime,
    period_end: datetime,
) -> ReportingObligationRecord | None: ...
async def get_issue(self, *, issue_key: str, account_id: str) ‑> ReportingIssueLifecycle | None
Expand source code
async def get_issue(self, *, issue_key: str, account_id: str) -> ReportingIssueLifecycle | None:
    """The live occurrence of this condition, or ``None``."""
    ...

The live occurrence of this condition, or None.

async def get_obligation(self, *, account_id: str, reporting_obligation_id: str) ‑> ReportingObligationRecord | None
Expand source code
async def get_obligation(
    self, *, account_id: str, reporting_obligation_id: str
) -> ReportingObligationRecord | None: ...
async def get_revision(self, *, account_id: str, reporting_revision_id: str) ‑> ReportingRevisionRecord | None
Expand source code
async def get_revision(
    self, *, account_id: str, reporting_revision_id: str
) -> ReportingRevisionRecord | None: ...
async def lease_period_close(self, *, worker_id: str, now: datetime, lease_seconds: float) ‑> LeasedConfiguration | None
Expand source code
async def lease_period_close(
    self, *, worker_id: str, now: datetime, lease_seconds: float
) -> LeasedConfiguration | None:
    """Lease one configuration generation for period-close work.

    Leasing rather than locking so a worker that dies mid-close releases its
    work by expiry instead of wedging the period forever, and so two workers
    cannot both close the same period.  Implementations should prefer the
    generation whose work is most overdue.
    """
    ...

Lease one configuration generation for period-close work.

Leasing rather than locking so a worker that dies mid-close releases its work by expiry instead of wedging the period forever, and so two workers cannot both close the same period. Implementations should prefer the generation whose work is most overdue.

async def list_adjustments(self, *, account_id: str, reporting_revision_ids: Sequence[str]) ‑> tuple[ReportingAdjustmentRecord, ...]
Expand source code
async def list_adjustments(
    self, *, account_id: str, reporting_revision_ids: Sequence[str]
) -> tuple[ReportingAdjustmentRecord, ...]: ...
async def list_all_configurations(self) ‑> tuple[ReportingConfiguration, ...]
Expand source code
async def list_all_configurations(self) -> tuple[ReportingConfiguration, ...]:
    """Enumerate retained generations for trusted service startup recovery.

    This administrative API is never exposed through buyer task handlers;
    those continue to use account-scoped reads after authorization.
    """
    raise NotImplementedError

Enumerate retained generations for trusted service startup recovery.

This administrative API is never exposed through buyer task handlers; those continue to use account-scoped reads after authorization.

async def list_configurations(self,
*,
caller: ReportingCaller,
delivery_config_ids: Sequence[str] | None = None) ‑> tuple[ReportingConfiguration, ...]
Expand source code
async def list_configurations(
    self, *, caller: ReportingCaller, delivery_config_ids: Sequence[str] | None = None
) -> tuple[ReportingConfiguration, ...]: ...
async def list_consumer_statuses(self,
*,
account_id: str,
consumer_id: str,
reporting_obligation_ids: Sequence[str] | None = None) ‑> tuple[ConsumerStatusRecord, ...]
Expand source code
async def list_consumer_statuses(
    self,
    *,
    account_id: str,
    consumer_id: str,
    reporting_obligation_ids: Sequence[str] | None = None,
) -> tuple[ConsumerStatusRecord, ...]: ...
async def list_revisions(self, *, account_id: str, reporting_obligation_id: str) ‑> tuple[ReportingRevisionRecord, ...]
Expand source code
async def list_revisions(
    self, *, account_id: str, reporting_obligation_id: str
) -> tuple[ReportingRevisionRecord, ...]: ...
async def open_snapshot(self, *, caller: ReportingCaller, filters_fingerprint: str) ‑> LedgerSnapshot
Expand source code
async def open_snapshot(
    self, *, caller: ReportingCaller, filters_fingerprint: str
) -> LedgerSnapshot:
    """Take a consistent read boundary for one account and filter set."""
    ...

Take a consistent read boundary for one account and filter set.

async def put_configuration(self, configuration: ReportingConfiguration) ‑> None
Expand source code
async def put_configuration(self, configuration: ReportingConfiguration) -> None:
    """Record an accepted configuration generation.

    A generation is immutable: re-putting one with changed content is a
    conflict, because a buyer that retained it derives expectations from it.
    """
    ...

Record an accepted configuration generation.

A generation is immutable: re-putting one with changed content is a conflict, because a buyer that retained it derives expectations from it.

async def read_page(self,
*,
snapshot: LedgerSnapshot,
consumer_id: str | None,
delivery_config_ids: Sequence[str] | None,
media_buy_ids: Sequence[str] | None,
offset: int,
limit: int,
changes_after_sequence: int | None) ‑> LedgerPage
Expand source code
async def read_page(
    self,
    *,
    snapshot: LedgerSnapshot,
    consumer_id: str | None,
    delivery_config_ids: Sequence[str] | None,
    media_buy_ids: Sequence[str] | None,
    offset: int,
    limit: int,
    changes_after_sequence: int | None,
) -> LedgerPage: ...
async def read_revision_rows(self,
*,
account_id: str,
reporting_revision_id: str,
cursor: str | None = None,
limit: int = 500) ‑> ReportingRowPage
Expand source code
async def read_revision_rows(
    self,
    *,
    account_id: str,
    reporting_revision_id: str,
    cursor: str | None = None,
    limit: int = 500,
) -> ReportingRowPage:
    """Walk one revision's frozen rows in stable order."""
    ...

Walk one revision's frozen rows in stable order.

async def record_consumer_status(self, status: ConsumerStatusRecord) ‑> tuple[ConsumerStatusRecord, bool]
Expand source code
async def record_consumer_status(
    self, status: ConsumerStatusRecord
) -> tuple[ConsumerStatusRecord, bool]:
    """Append a consumer status statement, superseding the chain's leaf.

    Returns ``(record, recorded)`` where ``recorded`` is ``False`` for an
    exact idempotent replay.  Supersession is atomic: naming a stale or
    missing leaf must raise rather than fork the chain, so a successful
    retry cannot erase a recorded outage.
    """
    ...

Append a consumer status statement, superseding the chain's leaf.

Returns (record, recorded) where recorded is False for an exact idempotent replay. Supersession is atomic: naming a stale or missing leaf must raise rather than fork the chain, so a successful retry cannot erase a recorded outage.

async def release_period_close(self,
lease: LeasedConfiguration,
*,
worker_id: str) ‑> None
Expand source code
async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) -> None: ...
async def resolve_consumer_status_replay(self, status: ConsumerStatusRecord) ‑> ConsumerStatusRecord | None
Expand source code
async def resolve_consumer_status_replay(
    self, status: ConsumerStatusRecord
) -> ConsumerStatusRecord | None:
    """Return the stored row when this exact statement was already recorded.

    Exists so the ingest can answer "is this an exact retry?" *before* it
    validates anything time-dependent. ``record_consumer_status`` already
    makes that check, but it runs last, and some validation the ingest does
    first -- notably "``content_mismatch`` must name the revision the seller
    currently requires" -- has an answer that changes as the seller
    publishes. An exact retry sent after a restatement would then be
    rejected for a reason that did not apply when it was first accepted,
    and the spec's ``result_mapping`` requires ``unchanged``.

    Returns ``None`` when the id is unknown, so the caller proceeds to
    validate and record. Raises :class:`LedgerConflictError` with
    ``STATUS_IDENTITY_CONFLICT`` when the id exists with different content:
    reuse with changed content is still a conflict, and answering
    ``unchanged`` there would let a buyer rewrite a recorded statement.
    """
    ...

Return the stored row when this exact statement was already recorded.

Exists so the ingest can answer "is this an exact retry?" before it validates anything time-dependent. record_consumer_status already makes that check, but it runs last, and some validation the ingest does first – notably "content_mismatch must name the revision the seller currently requires" – has an answer that changes as the seller publishes. An exact retry sent after a restatement would then be rejected for a reason that did not apply when it was first accepted, and the spec's result_mapping requires unchanged.

Returns None when the id is unknown, so the caller proceeds to validate and record. Raises :class:LedgerConflictError with STATUS_IDENTITY_CONFLICT when the id exists with different content: reuse with changed content is still a conflict, and answering unchanged there would let a buyer rewrite a recorded statement.

async def retire_issue(self, *, issue_key: str, account_id: str, at: datetime) ‑> ReportingIssueLifecycle | None
Expand source code
async def retire_issue(
    self, *, issue_key: str, account_id: str, at: datetime
) -> ReportingIssueLifecycle | None:
    """Mark this condition's live occurrence ``resolved``.

    Called by the projection when the condition no longer holds. Returns
    ``None`` when there was nothing live to retire, so a repeated
    projection is convergent rather than an error.
    """
    ...

Mark this condition's live occurrence resolved.

Called by the projection when the condition no longer holds. Returns None when there was nothing live to retire, so a repeated projection is convergent rather than an error.

async def set_issue_state(self,
*,
issue_key: str,
account_id: str,
state: "Literal['acknowledged', 'waived']",
at: datetime,
external_ref: str | None = None) ‑> ReportingIssueLifecycle
Expand source code
async def set_issue_state(
    self,
    *,
    issue_key: str,
    account_id: str,
    state: Literal["acknowledged", "waived"],
    at: datetime,
    external_ref: str | None = None,
) -> ReportingIssueLifecycle:
    """Move a live occurrence forward, or attach an ``external_ref``.

    Deliberately cannot set ``resolved``. The spec forbids retiring a
    ``CONSUMER_STATUS_MISMATCH`` out of a degraded projection while the
    statement that caused it is still the consumer's current leaf, so
    resolution is not an operator action -- it is what the projection does
    when the condition actually clears (see :meth:`retire_issue`).
    Otherwise a seller could unilaterally erase a buyer-attributed
    disagreement, which is the one outcome this separately attributed loop
    exists to prevent.

    ``waived`` requires explicit off-protocol agreement by this consumer
    and seller for the exact issue, causing statement and diagnosed
    conflict. The adopter must retain its private consent audit before
    calling this method; ``external_ref`` is inert correlation, not proof.
    Rc.6 removes that issue and restores underlying seller health, without
    changing the consumer statement. Later statements or different
    conflicts are evaluated independently, never covered by this waiver.
    """
    ...

Move a live occurrence forward, or attach an external_ref.

Deliberately cannot set resolved. The spec forbids retiring a CONSUMER_STATUS_MISMATCH out of a degraded projection while the statement that caused it is still the consumer's current leaf, so resolution is not an operator action – it is what the projection does when the condition actually clears (see :meth:retire_issue). Otherwise a seller could unilaterally erase a buyer-attributed disagreement, which is the one outcome this separately attributed loop exists to prevent.

waived requires explicit off-protocol agreement by this consumer and seller for the exact issue, causing statement and diagnosed conflict. The adopter must retain its private consent audit before calling this method; external_ref is inert correlation, not proof. Rc.6 removes that issue and restores underlying seller health, without changing the consumer statement. Later statements or different conflicts are evaluated independently, never covered by this waiver.

async def set_revision_readable(self, *, account_id: str, reporting_revision_id: str, readable: bool) ‑> None
Expand source code
async def set_revision_readable(
    self, *, account_id: str, reporting_revision_id: str, readable: bool
) -> None:
    """Record that a revision's content is (no longer) readable.

    Retention expiry and storage loss are real; representing them is how an
    obligation becomes honestly ``action_required`` instead of staying
    ``complete`` over evidence nobody can read.
    """
    ...

Record that a revision's content is (no longer) readable.

Retention expiry and storage loss are real; representing them is how an obligation becomes honestly action_required instead of staying complete over evidence nobody can read.

class ReportingRowPage (rows: tuple[dict[str, Any], ...],
total_count: int,
has_more: bool,
cursor: str | None,
reporting_revision_id: str | None = None)
Expand source code
@dataclass(frozen=True)
class ReportingRowPage:
    """One page of an exact revision read.

    Every page repeats identical revision metadata and binding; a consumer
    hashes the concatenated rows in cursor order only after exhausting the
    frozen walk.  Do not substitute a fresh date-range pull for this read -- it
    may have changed since the immutable revision was published.
    """

    rows: tuple[dict[str, Any], ...]
    total_count: int
    has_more: bool
    cursor: str | None
    reporting_revision_id: str | None = None

One page of an exact revision read.

Every page repeats identical revision metadata and binding; a consumer hashes the concatenated rows in cursor order only after exhausting the frozen walk. Do not substitute a fresh date-range pull for this read – it may have changed since the immutable revision was published.

Instance variables

var cursor : str | None
var has_more : bool
var reporting_revision_id : str | None
var rows : tuple[dict[str, typing.Any], ...]
var total_count : int
class RestatementCheckpoint (account_id: str,
reporting_obligation_id: str,
checked_at: datetime,
next_observation: int,
provisional_until: datetime | None = None)
Expand source code
@dataclass(frozen=True)
class RestatementCheckpoint:
    """Durable scheduling state for successful source observations.

    ``next_observation`` advances even when the source content is unchanged.
    Without that durable ordinal the next scheduled refresh would replay the
    same sealed source execution forever instead of making a new observation.
    """

    account_id: str
    reporting_obligation_id: str
    checked_at: datetime
    next_observation: int
    provisional_until: datetime | None = None

    def __post_init__(self) -> None:
        if self.next_observation < 1:
            raise ValueError("next_observation must be at least one")
        _utc(self.checked_at)
        if self.provisional_until is not None:
            _utc(self.provisional_until)

Durable scheduling state for successful source observations.

next_observation advances even when the source content is unchanged. Without that durable ordinal the next scheduled refresh would replay the same sealed source execution forever instead of making a new observation.

Instance variables

var account_id : str
var checked_at : datetime.datetime
var next_observation : int
var provisional_until : datetime.datetime | None
var reporting_obligation_id : str
class RestatementCheckpointStore (*args, **kwargs)
Expand source code
@runtime_checkable
class RestatementCheckpointStore(Protocol):
    """Optional store extension used only by source settling policies."""

    async def get_restatement_checkpoint(
        self, *, account_id: str, reporting_obligation_id: str
    ) -> RestatementCheckpoint | None: ...

    async def record_restatement_checkpoint(
        self, checkpoint: RestatementCheckpoint
    ) -> RestatementCheckpoint:
        """Persist the next source observation ordinal after a successful read."""
        ...

Optional store extension used only by source settling policies.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def get_restatement_checkpoint(self, *, account_id: str, reporting_obligation_id: str) ‑> RestatementCheckpoint | None
Expand source code
async def get_restatement_checkpoint(
    self, *, account_id: str, reporting_obligation_id: str
) -> RestatementCheckpoint | None: ...
async def record_restatement_checkpoint(self,
checkpoint: RestatementCheckpoint) ‑> RestatementCheckpoint
Expand source code
async def record_restatement_checkpoint(
    self, checkpoint: RestatementCheckpoint
) -> RestatementCheckpoint:
    """Persist the next source observation ordinal after a successful read."""
    ...

Persist the next source observation ordinal after a successful read.

class RetryScheduleEntry (scope_key: str,
retry_not_before: datetime,
attempt: int,
blocked: bool = False,
recorded_at: datetime | None = None,
replayed: bool = False)
Expand source code
@dataclass(frozen=True)
class RetryScheduleEntry:
    """The next permitted source read for one slice, account, or source scope.

    ``attempt`` counts consecutive failures. Zero clears the schedule after a
    successful read. Scope keys are generated by the producer, not callers.
    """

    scope_key: str
    retry_not_before: datetime
    attempt: int
    blocked: bool = False
    recorded_at: datetime | None = None
    replayed: bool = field(default=False, compare=False, repr=False)

    def __post_init__(self) -> None:
        if not self.scope_key or self.attempt < 0:
            raise ValueError("retry schedule requires a scope key and nonnegative attempt")
        _utc(self.retry_not_before)
        if self.recorded_at is None:
            object.__setattr__(self, "recorded_at", self.retry_not_before)
        assert self.recorded_at is not None
        _utc(self.recorded_at)

The next permitted source read for one slice, account, or source scope.

attempt counts consecutive failures. Zero clears the schedule after a successful read. Scope keys are generated by the producer, not callers.

Instance variables

var attempt : int
var blocked : bool
var recorded_at : datetime.datetime | None
var replayed : bool
var retry_not_before : datetime.datetime
var scope_key : str
class RetryScheduleStore (*args, **kwargs)
Expand source code
@runtime_checkable
class RetryScheduleStore(Protocol):
    """Optional shared retry state for producer source acquisitions."""

    async def get_retry_schedule(self, *, scope_key: str) -> RetryScheduleEntry | None:
        pass

    async def record_retry_schedule(self, entry: RetryScheduleEntry) -> RetryScheduleEntry:
        pass

Optional shared retry state for producer source acquisitions.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def get_retry_schedule(self, *, scope_key: str) ‑> RetryScheduleEntry | None
Expand source code
async def get_retry_schedule(self, *, scope_key: str) -> RetryScheduleEntry | None:
    pass
async def record_retry_schedule(self,
entry: RetryScheduleEntry) ‑> RetryScheduleEntry
Expand source code
async def record_retry_schedule(self, entry: RetryScheduleEntry) -> RetryScheduleEntry:
    pass