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_idwith 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:
LedgerConflictErrorrather 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.waivedis not retirable: it is already out of the projection by agreement, and overwriting that readable act withresolvedis an edge the forward-only lifecycle forbids. Both stores consult this rather than each remembering the rule. -
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: consumerbefore 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_FEATUREbefore the generation becomes ready and before any obligation exists – and must not silently coerce it toseller.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.
clocksupplies 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 selfGroup 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 : strvar consumer_id : strvar delivery_config_id : strvar delivery_config_version : intprop 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 = codeA write would violate an invariant that makes the ledger evidence.
Carries a stable
codeso 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 | NoneOne 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 | Nonevar has_more : boolvar 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 toopened_at, so a re-emission that advanced it would let an unattended mismatch stay belowaction_requiredindefinitely.This is the one write the read path performs.
get_reporting_statusis 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
resolvedorwaivedopens a new occurrence with the nextgeneration, a newissue_id, and a newopened_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 NotImplementedErrorEnumerate 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)whererecordedisFalsefor 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_statusalready makes that check, but it runs last, and some validation the ingest does first – notably "content_mismatchmust 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'sresult_mappingrequiresunchanged.Returns
Nonewhen the id is unknown, so the caller proceeds to validate and record. Raises :class:LedgerConflictErrorwithSTATUS_IDENTITY_CONFLICTwhen the id exists with different content: reuse with changed content is still a conflict, and answeringunchangedthere 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
Nonewhen 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 aCONSUMER_STATUS_MISMATCHout 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.waivedrequires 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_refis 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_requiredinstead of stayingcompleteover 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 = NoneOne 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 | Nonevar has_more : boolvar reporting_revision_id : str | Nonevar 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_observationadvances 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 : strvar checked_at : datetime.datetimevar next_observation : intvar provisional_until : datetime.datetime | Nonevar 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.
attemptcounts consecutive failures. Zero clears the schedule after a successful read. Scope keys are generated by the producer, not callers.Instance variables
var attempt : intvar blocked : boolvar recorded_at : datetime.datetime | Nonevar replayed : boolvar retry_not_before : datetime.datetimevar 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: passOptional 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