Module adcp.reporting.ledger.pg
PostgreSQL-backed :class:~adcp.reporting.ledger.store.ReportingLedgerStore.
Durable counterpart to
:class:~adcp.reporting.ledger.store.InMemoryReportingLedgerStore, following
the same shape as :class:PgTaskRegistry: the caller supplies
an :class:psycopg_pool.AsyncConnectionPool, we never open, own, or close it,
and each operation takes a short-lived connection.
Quickstart
::
from psycopg_pool import AsyncConnectionPool
from adcp.reporting.ledger import PgReportingLedgerStore, ReportingProducer
async with AsyncConnectionPool(dsn, min_size=2, max_size=10) as pool:
store = PgReportingLedgerStore(pool=pool)
await store.create_schema() # idempotent; safe on every boot
producer = ReportingProducer(source=..., offerings=..., store=store)
await producer.run_worker()
Schema Bootstrap
:meth:create_schema creates or upgrades the schema transactionally, including
the account-qualified configuration primary key for beta.15 installations.
The raw DDL ships in :file:reporting_ledger.sql followed by
:file:reporting_ledger_account_generations.sql,
:file:reporting_ledger_obligation_currency.sql,
:file:reporting_ledger_reconciliation.sql, and
:file:reporting_notification_outbox.sql; run all five in one transaction
when using Alembic, Flyway, or psql. See :file:docs/reporting-ledger-migration.md
for deployment and compatibility notes.
Where The Invariants Actually Live
In the schema, not in Python. A ledger that enforces "one official revision per obligation" in application code is one deploy away from two workers racing past it. The partial unique indexes in the DDL are the enforcement:
- one obligation per
(account, config generation, period); - one official revision per obligation;
- one successor per superseded revision;
- one unsuperseded consumer-status leaf per logical chain.
Python turns the resulting integrity errors into typed
:class:~adcp.reporting.ledger.store.LedgerConflictError values so a handler can
map them to wire errors without matching on driver prose.
A note on the # noqa: S608
# nosec B608 comments
Several queries f-string a column list or a filter clause into the SQL.
Every
interpolated value is a module-level constant defined in this file – column
lists and literal clause fragments – and never a caller-supplied value.
Caller data is always bound as a %s parameter.
The alternative, repeating
twenty column names inline at each of nine call sites, is how column lists and
row mappers drift apart.
Change-feed ordering
Every immutable write appends to reporting_ledger_changes in the same
transaction, under a transaction-scoped advisory lock on the account.
That
lock is what makes sequence order equal commit order per account: without it a
transaction could take a low sequence, commit late, and be invisible to a
consumer that had already checkpointed past it.
A silently lost record is the
one failure mode a reporting ledger must not have.
Classes
class PgReportingLedgerStore (*,
pool: AsyncConnectionPool,
clock: Callable[[], datetime] | None = None,
notifications: bool = False)-
Expand source code
class PgReportingLedgerStore: """Durable reporting ledger over a caller-supplied connection pool.""" is_durable: ClassVar[bool] = True def __init__( self, *, pool: AsyncConnectionPool, clock: Callable[[], datetime] | None = None, notifications: bool = False, ) -> None: if not PG_AVAILABLE: raise ImportError(_INSTALL_HINT) self._pool = pool # Change timestamps and the snapshot boundary default to the *database* clock, # which is what makes two readers of one snapshot agree even across # application hosts with drifting clocks -- do not override it in # production for that reason. Overriding is for replay, backfill, and # tests that need to stand at a specific instant relative to seeded # evidence rather than wherever wall-clock time happens to fall. self._clock = clock self._notifications_enabled = notifications self._period_close_sample: tuple[int, datetime | None, str, str, str, int] | None = None @asynccontextmanager async def _connection(self) -> AsyncIterator[Any]: bound = _BOUND_CONNECTION.get() if bound is not None and bound[:2] == (self._pool, asyncio.current_task()): yield bound[2] else: async with self._pool.connection() as connection: yield connection @asynccontextmanager async def transaction(self) -> AsyncIterator[PgReportingLedgerStore]: """Group source operations into one committed status boundary. Pool ownership remains with the adopter. Nested turns are savepoints; callers acquire accounts in canonical order when touching several. """ async with self._connection() as connection, connection.transaction(): token = _BOUND_CONNECTION.set((self._pool, asyncio.current_task(), connection)) try: yield self finally: _BOUND_CONNECTION.reset(token) @asynccontextmanager async def _source_publication( self, account_id: str, *, seals: ReportingSealStore | None = None ) -> AsyncIterator[InlineSealPublisher | None]: from psycopg import Error from adcp.reporting.inline_storage import InlineStorageError, PgReportingSealStore if isinstance(seals, PgReportingSealStore) and seals._pool is self._pool: # The backend's audited READ COMMITTED transaction owns both the # account lock and seal write. A second pool checkout would deadlock # with a size-one pool, and would separate the publication boundary. failure = None try: async with seals._transaction() as connection: await self._lock_account(connection, account_id) owner = asyncio.current_task() async def publish(key: str, sealed: SealedSlice) -> SealedSlice: if asyncio.current_task() is not owner: raise InlineStorageError("INVALID_INPUT") prepared = seals._prepare_seal( account_id=account_id, source_execution_key=key, sealed=sealed ) return await seals._put_on(connection, prepared) yield publish except InlineStorageError as error: failure = error.code except Error: failure = "RESOURCE_UNAVAILABLE" if failure is not None: raise InlineStorageError(failure) return # Keep the live authorization check and ledger commit on the same # account transaction, including commits replayed after a restart. async with self.transaction(), self._connection() as connection: await self._lock_account(connection, account_id) yield None async def create_schema(self) -> None: """Create or upgrade the ledger atomically, serializing concurrent boots. The bootstrap takes a transaction-scoped schema lock before any DDL. Keep the migration in that transaction, including with an autocommit pool, so a second process cannot observe a partially upgraded schema. """ async with self._connection() as connection: async with connection.transaction(): await self._create_schema_on(connection) async def _create_schema_on(self, connection: Any) -> None: from adcp.reporting.migration import require_owned_schema_or_empty await require_owned_schema_or_empty(connection) for path in ( _DDL_PATH, _ACCOUNT_GENERATIONS_DDL_PATH, _CURRENCY_DDL_PATH, _RECONCILIATION_DDL_PATH, _NOTIFICATIONS_DDL_PATH, _ACTIVITY_DDL_PATH, Path(__file__).with_name("reporting_provisional_observations.sql"), Path(__file__).with_name("reporting_caller_ownership.sql"), ): await connection.execute(path.read_text()) await self._require_provisional_schema(connection) if self._notifications_enabled: from adcp.reporting.outbox._schema import validate_schema await validate_schema(connection) async def _require_provisional_schema(self, connection: Any) -> None: from adcp.reporting.outbox._schema import validate_provisional_schema await validate_provisional_schema(connection) async def _notification_now(self, connection: Any) -> datetime: from adcp.reporting.outbox.pg import database_now if self._clock is not None: # Preserve the explicit conformance clock across deferred boundary # capture. Production leaves this unset and captures database time # after the complete source transaction's writes, before commit. await connection.execute("SELECT set_config('adcp.status_clock_override', 'on', true)") return await database_now(connection, self._clock) async def _record_notification(self, connection: Any, event: ReportingDomainEvent) -> None: if self._notifications_enabled: from adcp.reporting.outbox.pg import enqueue_event await enqueue_event(connection, event) async def _dirty_status( self, connection: Any, scope: ReportingStatusScope, reason: DirtyReason, before: ReportingStatusEvidence | None = None, after: ReportingStatusEvidence | None = None, ) -> None: if self._notifications_enabled: from adcp.reporting.ledger.status_snapshot import settle_snapshot_on from adcp.reporting.outbox.pg import mark_dirty await settle_snapshot_on(self, connection, account_id=scope.account_id) await mark_dirty( connection, scope, reason, await self._notification_now(connection), before, after ) async def _dirty_issue( self, connection: Any, issue: ReportingIssueLifecycle, status_scope: ReportingStatusScope | None, before: ReportingIssueLifecycle | None = None, *, enqueue: bool = True, ) -> None: from dataclasses import asdict has_scope_storage = await self._issue_scope_storage_on(connection) row = ( await ( await connection.execute( "SELECT scope FROM reporting_issue_status_scopes" " WHERE account_id = %s AND issue_id = %s", (issue.account_id, issue.issue_id), ) ).fetchone() if has_scope_storage else None ) existing = decode_status_scope(row[0]) if row else None scope = ( status_scope or existing or ReportingStatusScope(issue.account_id, consumer_id=issue.consumer_id) ) if ( scope.account_id != issue.account_id or issue.consumer_id is not None and scope.consumer_id != issue.consumer_id ): raise ReportingNotificationError("invalid_status_scope") validate_scope_refinement(existing, scope) if scope.generation_key is not None: key = scope.generation_key configuration = await ( await connection.execute( "SELECT feed_purpose FROM reporting_configurations WHERE account_id = %s" " AND consumer_id = %s AND delivery_config_id = %s AND delivery_co" "nfig_version = %s", ( scope.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ), ) ).fetchone() if configuration is None or ( scope.feed_purpose is not None and configuration[0] != scope.feed_purpose ): raise ReportingNotificationError("invalid_status_scope") if scope.reporting_obligation_id is not None: obligation = await ( await connection.execute( "SELECT delivery_config_id, delivery_config_version, feed_purpose" " FROM reporting_obligations WHERE account_id = %s" " AND consumer_id = %s AND reporting_obligation_id = %s", (scope.account_id, scope.consumer_id, scope.reporting_obligation_id), ) ).fetchone() if ( obligation is None or ( scope.generation_key is not None and obligation[:2] != ( scope.generation_key.delivery_config_id, scope.generation_key.delivery_config_version, ) ) or (scope.feed_purpose is not None and obligation[2] != scope.feed_purpose) ): raise ReportingNotificationError("invalid_status_scope") if not has_scope_storage: # Default-off Core writers also support the pre-outbox schema. # Migration introduces scope persistence without rewriting history. return await connection.execute( "INSERT INTO reporting_issue_status_scopes (account_id, issue_id, scope)" " VALUES (%s,%s,%s::jsonb) ON CONFLICT (account_id, issue_id)" " DO UPDATE SET scope = EXCLUDED.scope", (issue.account_id, issue.issue_id, _json(asdict(scope))), ) if enqueue and (before != issue or existing != scope): await self._dirty_status( connection, scope, "issue", issue_evidence(before) if before else None, issue_evidence(issue), ) async def _issue_scope_storage_on(self, connection: Any) -> bool: row = await ( await connection.execute( "SELECT to_regclass(quote_ident(current_schema())" " || '.reporting_issue_status_scopes')" ) ).fetchone() present = row is not None and row[0] is not None if not present and self._notifications_enabled: raise ReportingNotificationError( "notification_schema_unready:missing:reporting_issue_status_scopes" ) return present # -- change feed ------------------------------------------------------ async def _append_change( self, connection: Any, account_id: str, kind: LedgerRecordKind, record_id: str, *, consumer_id: str, ) -> None: await connection.execute( "INSERT INTO reporting_ledger_changes" " (account_id, consumer_id, record_kind, record_id, committed_at)" " VALUES (%s, %s, %s, %s, COALESCE(%s, now()))" " ON CONFLICT (account_id, consumer_id, record_kind, record_id) DO NOTHING", ( account_id, consumer_id, kind, record_id, self._clock() if self._clock is not None else None, ), ) @staticmethod async def _bound_lock_waits(connection: Any) -> None: # set_config(..., true) is SET LOCAL: the cap ends with this transaction # (or savepoint rollback) and never lengthens an adopter's shorter limit. # SKIP LOCKED does not skip relation, FK, trigger, or advisory-lock waits. await connection.execute( "SELECT set_config('lock_timeout'," " LEAST(COALESCE(NULLIF(setting::integer,0),5000),5000)::text,true)" " FROM pg_settings WHERE name='lock_timeout'" ) def _worker_lock_errors(self) -> tuple[type[Exception], ...]: from psycopg.errors import DeadlockDetected, LockNotAvailable bound = _BOUND_CONNECTION.get() if bound is not None and bound[:2] == (self._pool, asyncio.current_task()): # The adopter owns the outer rollback and any locks acquired before # our savepoint. Only standalone worker turns may reacquire a lease. return () return (DeadlockDetected, LockNotAvailable) @staticmethod async def _lock_account(connection: Any, account_id: str) -> None: """Serialize this account's feed appends so seq order == commit order.""" await PgReportingLedgerStore._bound_lock_waits(connection) await connection.execute( "SELECT pg_advisory_xact_lock(hashtext(%s))", (f"adcp.reporting:{account_id}",) ) # Optional C lock order: account, boundary cursor, canonical typed scopes, # then source evidence. A/B schemas and custom stores need no C methods. present = await ( await connection.execute( "SELECT to_regclass(quote_ident(current_schema()) || '.reporting_status_accounts')" ) ).fetchone() if present is not None and present[0] is not None: await ( await connection.execute( "SELECT 1 FROM reporting_status_accounts WHERE account_id=%s FOR UPDATE", (account_id,), ) ).fetchall() await ( await connection.execute( "SELECT 1 FROM reporting_status_scope_checkpoints WHERE account_id=%s" " ORDER BY account_id, consumer_namespace, delivery_config_id, version," " scope_kind, obligation_namespace FOR UPDATE", (account_id,), ) ).fetchall() # -- configurations --------------------------------------------------- async def put_configuration(self, configuration: ReportingConfiguration) -> None: reject_reserved_authoritative_party(configuration) key = configuration.generation_key payload = _configuration_payload(configuration) digest = _fingerprint(payload) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, key.account_id) inserted_cursor = await connection.execute( "INSERT INTO reporting_configurations" " (delivery_config_id, delivery_config_version, account_id," " report_definition_id, reporting_profile, feed_purpose, required_finality," " account_timezone, schedule, media_buy_ids, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, definition," " authoritative_party, content_sha256, consumer_id, quarantined)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s::jsonb, %s, %s, %s, %s," " %s::jsonb, %s, %s, %s, %s)" " ON CONFLICT (account_id, consumer_id, delivery_config_id, delive" "ry_config_version) DO NOTHING" " RETURNING delivery_config_id", ( key.delivery_config_id, key.delivery_config_version, key.account_id, configuration.report_definition_id, configuration.reporting_profile, configuration.feed_purpose, configuration.required_finality, configuration.account_timezone, _json(payload["schedule"]), _json(sorted(configuration.media_buy_ids)), configuration.activated_at, configuration.deactivated_at, configuration.automated_recovery_window.total_seconds(), configuration.status_retention_days, _json(payload["definition"]) if payload["definition"] else None, configuration.authoritative_party, digest, configuration.consumer_id, configuration.quarantined, ), ) inserted = await inserted_cursor.fetchone() # Check *after* the insert. ON CONFLICT waits for a concurrent # winner; a fresh READ COMMITTED statement sees its retained # content. A pre-insert check followed by DO NOTHING could silently # accept a different immutable generation from a losing writer. row = await ( await connection.execute( "SELECT content_sha256, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, quarantined" " FROM reporting_configurations" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s FOR UPDATE", ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ), ) ).fetchone() if row is not None and row[5]: raise LedgerConflictError( "REPORTING_GENERATION_QUARANTINED", "legacy evidence is read-only" ) if row is None or row[0] != digest: 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", ) scope = ReportingStatusScope(configuration.account_id, configuration.generation_key) if inserted is not None: if self._notifications_enabled: await self._dirty_status( connection, scope, "configuration", after=configuration_evidence(configuration), ) return # rc.3 carries activation/deactivation and the recovery/retention # windows as lifecycle state over one immutable generation, so a # re-put that changes only those must apply -- otherwise a # deactivated feed keeps minting obligations -- and must co-commit # its status-dirty generation on this exact connection. An # unchanged re-put stays a no-op and enqueues nothing. retained = replace( configuration, activated_at=_utc(row[1]) if row[1] else None, deactivated_at=_utc(row[2]) if row[2] else None, automated_recovery_window=timedelta(seconds=float(row[3])), status_retention_days=row[4], ) if configuration_lifecycle(retained) == configuration_lifecycle(configuration): return await connection.execute( "UPDATE reporting_configurations SET activated_at = %s, deactivated_at = %s," " automated_recovery_seconds = %s, status_retention_days = %s" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s", ( configuration.activated_at, configuration.deactivated_at, configuration.automated_recovery_window.total_seconds(), configuration.status_retention_days, key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ), ) if self._notifications_enabled: await self._dirty_status( connection, scope, "configuration", before=configuration_evidence(retained), after=configuration_evidence(configuration), ) async def list_configurations( self, *, caller: ReportingCaller, delivery_config_ids: Sequence[str] | None = None ) -> tuple[ReportingConfiguration, ...]: async with self._connection() as connection: return await self._list_configurations_on( connection, account_id=caller.account_id, consumer_id=caller.consumer_id, delivery_config_ids=delivery_config_ids, ) async def list_all_configurations(self) -> tuple[ReportingConfiguration, ...]: """One ordered scan for trusted service recovery, including retired generations.""" async with self._connection() as connection: rows = await ( await connection.execute( "SELECT delivery_config_id, delivery_config_version, account_id," " report_definition_id, reporting_profile, feed_purpose, required_finality," " account_timezone, schedule, media_buy_ids, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, definition," " authoritative_party, consumer_id, quarantined" " FROM reporting_configurations WHERE NOT quarantined" " ORDER BY account_id, consumer_id, delivery_config_id, delivery_config_version" ) ).fetchall() return tuple(_configuration_from_row(row) for row in rows) @staticmethod async def _list_configurations_on( connection: Any, *, account_id: str, consumer_id: str, delivery_config_ids: Sequence[str] | None = None, ) -> tuple[ReportingConfiguration, ...]: clause = " AND delivery_config_id = ANY(%s)" if delivery_config_ids else "" params: list[Any] = [account_id, consumer_id] if delivery_config_ids: params.append(list(delivery_config_ids)) rows = await ( await connection.execute( "SELECT delivery_config_id, delivery_config_version, account_id," # noqa: S608 # nosec B608 " report_definition_id, reporting_profile, feed_purpose, required_finality," " account_timezone, schedule, media_buy_ids, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, definition," " authoritative_party, consumer_id, quarantined" " FROM reporting_configurations" f" WHERE account_id = %s AND consumer_id = %s{clause}" # noqa: S608 — clause is a literal " ORDER BY delivery_config_id, delivery_config_version", tuple(params), ) ).fetchall() return tuple(_configuration_from_row(row) for row in rows) @staticmethod async def _require_writable_generation( connection: Any, key: ReportingConfigurationGenerationKey ) -> None: row = await ( await connection.execute( "SELECT quarantined FROM reporting_configurations WHERE account_id=%s" " AND consumer_id=%s AND delivery_config_id=%s AND delivery_config_version=%s", ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ), ) ).fetchone() if row is None: raise LedgerConflictError( "UNKNOWN_CONFIGURATION_GENERATION", "generation is unavailable" ) if row[0]: raise LedgerConflictError( "REPORTING_GENERATION_QUARANTINED", "establish a new owned generation after operator reconciliation", ) # -- obligations ------------------------------------------------------ async def commit_obligation( self, obligation: ReportingObligationRecord ) -> ReportingObligationRecord: key = obligation.generation_key async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, key.account_id) await self._require_writable_generation(connection, key) existing_row = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND period_start = %s AND period_end = %s", ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, obligation.period.start, obligation.period.end, ), ) ).fetchone() if existing_row is not None: return _obligation_from_row(existing_row) require_frozen_currency(obligation.currency) try: inserted = await ( await connection.execute( "INSERT INTO reporting_obligations" " (reporting_obligation_id, account_id, delivery_config_id," " delivery_config_version, report_definition_id, reporting_profile," " feed_purpose, period_key, period_start, period_end, source_timezone," " expected_at, scope_resolved_at, automated_recovery_deadline_at," " required_finality, coverage_status, media_buy_ids, package_ids," " schedule, definition, created_at, currency, consumer_id)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s," " %s::jsonb, %s::jsonb, %s::jsonb, %s::jsonb, %s, %s, %s)" " ON CONFLICT (account_id, consumer_id, delivery_config_id, delive" "ry_config_version," " period_start, period_end) DO NOTHING" " RETURNING reporting_obligation_id", ( obligation.reporting_obligation_id, key.account_id, key.delivery_config_id, key.delivery_config_version, obligation.report_definition_id, obligation.reporting_profile, obligation.feed_purpose, obligation.period.period_key, obligation.period.start, obligation.period.end, obligation.period.source_timezone, obligation.period.expected_at, obligation.scope_resolved_at, obligation.automated_recovery_deadline_at, obligation.required_finality, obligation.coverage_status, _json(sorted(obligation.media_buy_ids)), _json(sorted(obligation.package_ids)), _json(_schedule_payload(obligation.schedule)), ( _json(_definition_payload(obligation.definition)) if obligation.definition else None ), obligation.created_at, obligation.currency, obligation.consumer_id, ), ) ).fetchone() except Exception as error: raise _translate_integrity_error(error) from error if inserted is not None: await self._append_change( connection, obligation.account_id, "obligation", obligation.reporting_obligation_id, consumer_id=obligation.consumer_id, ) if self._notifications_enabled: await self._dirty_status( connection, ReportingStatusScope.for_obligation(obligation), "obligation", after=ReportingStatusEvidence( "obligation", obligation.reporting_obligation_id ), ) return obligation # Another worker won the period close; converge on its obligation. existing = await self.find_obligation( account_id=obligation.account_id, consumer_id=obligation.consumer_id, delivery_config_id=obligation.delivery_config_id, delivery_config_version=obligation.delivery_config_version, period_start=obligation.period.start, period_end=obligation.period.end, ) if existing is None: # pragma: no cover - only under concurrent deletion raise LedgerConflictError( "OBLIGATION_LOST", "the obligation vanished between insert and read" ) return existing async def get_obligation( self, *, account_id: str, reporting_obligation_id: str ) -> ReportingObligationRecord | None: async with self._connection() as connection: row = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE account_id = %s AND reporting_obligation_id = %s", (account_id, reporting_obligation_id), ) ).fetchone() return _obligation_from_row(row) if row 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, ) async with self._connection() as connection: row = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND period_start = %s AND period_end = %s", ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, period_start, period_end, ), ) ).fetchone() return _obligation_from_row(row) if row else None # -- revisions -------------------------------------------------------- async def commit_revision( self, revision: ReportingRevisionRecord, rows: Sequence[dict[str, Any]] ) -> ReportingRevisionRecord: rows = tuple(deepcopy(row) for row in rows) 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) digest = _fingerprint(_revision_payload(revision)) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, revision.account_id) existing = await ( await connection.execute( f"SELECT {_REVISION_COLUMNS}, content_sha256 FROM reporting_revisions" # noqa: S608 # nosec B608 " WHERE reporting_revision_id = %s AND account_id = %s", (revision.reporting_revision_id, revision.account_id), ) ).fetchone() if existing is not None: if existing[-1] != digest: raise LedgerConflictError( "REVISION_IMMUTABLE", f"revision {revision.reporting_revision_id} already exists with " "different content; a restatement is a new revision", ) return _revision_from_row(existing[:-1]) obligation = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE reporting_obligation_id = %s AND account_id = %s", (revision.reporting_obligation_id, revision.account_id), ) ).fetchone() if obligation is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "a revision must attach to an obligation committed at the period close", ) await self._require_writable_generation( connection, _obligation_from_row(obligation).generation_key ) validate_revision_currency(_obligation_from_row(obligation), revision, rows) if revision.supersedes_reporting_revision_id: await self._require_current_leaf(connection, revision) write_error = None try: await connection.execute( "INSERT INTO reporting_revisions" " (reporting_revision_id, account_id, reporting_obligation_id, finality," " revision_content_sha256, row_count, control_totals, observed_at," " data_through, created_at, supersedes_reporting_revision_id," " finality_basis, finality_policy_id, finalized_at, readable," " readable_at_commit, source_publication_id, source_manifest_sha256," " content_sha256, canonical_content_digest, managed_control_totals)" " VALUES (%s, %s, %s, %s, %s, %s, %s::jsonb, %s, %s, %s, %s, %s, %s, %s," " %s, %s, %s, %s, %s, %s::jsonb, %s::jsonb)", ( revision.reporting_revision_id, revision.account_id, revision.reporting_obligation_id, revision.finality, revision.revision_content_sha256, revision.row_count, _json([[name, value] for name, value in revision.control_totals]), revision.observed_at, revision.data_through, revision.created_at, revision.supersedes_reporting_revision_id, revision.finality_basis, revision.finality_policy_id, revision.finalized_at, revision.readable, revision.readable_at_commit, revision.source_publication_id, revision.source_manifest_sha256, digest, ( _json(revision.canonical_content_digest.to_wire()) if revision.canonical_content_digest is not None else None ), ( _json([item.to_wire() for item in revision.managed_control_totals]) if revision.managed_control_totals is not None else None ), ), ) except Exception as error: # psycopg raises UniqueViolation subclasses write_error = _translate_integrity_error(error) if write_error is not None: raise write_error if rows: await connection.cursor().executemany( "INSERT INTO reporting_revision_rows" " (reporting_revision_id, ordinal, row_payload)" " SELECT reporting_revision_id, %s, %s::jsonb FROM reporting_revisions" " WHERE account_id = %s AND reporting_revision_id = %s", [ (ordinal, _json(row), revision.account_id, revision.reporting_revision_id) for ordinal, row in enumerate(rows) ], ) await self._append_change( connection, revision.account_id, "revision", revision.reporting_revision_id, consumer_id=_obligation_from_row(obligation).consumer_id, ) if self._notifications_enabled: from adcp.reporting.ledger.notification_events import revision_event await self._record_notification( connection, revision_event( revision, await self._notification_now(connection), consumer_id=_obligation_from_row(obligation).consumer_id, ), ) await self._dirty_status( connection, ReportingStatusScope.for_obligation(_obligation_from_row(obligation)), "revision", after=ReportingStatusEvidence( "revision", revision.reporting_revision_id, readable=revision.readable, supersedes_id=revision.supersedes_reporting_revision_id, ), ) return revision @staticmethod async def _require_current_leaf(connection: Any, revision: ReportingRevisionRecord) -> None: target = revision.supersedes_reporting_revision_id row = await ( await connection.execute( "SELECT r.reporting_revision_id," " EXISTS (SELECT 1 FROM reporting_revisions s" " WHERE s.account_id = r.account_id" " AND s.supersedes_reporting_revision_id = r.reporting_revision_id)" " FROM reporting_revisions r" " WHERE r.reporting_revision_id = %s AND r.account_id = %s" " AND r.reporting_obligation_id = %s", (target, revision.account_id, revision.reporting_obligation_id), ) ).fetchone() if row is None: raise LedgerConflictError( "SUPERSEDES_UNKNOWN", f"revision {target} is not part of this obligation's chain" ) if row[1]: 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, ...]: async with self._connection() as connection: rows = await ( await connection.execute( f"SELECT {_REVISION_COLUMNS} FROM reporting_revisions" # noqa: S608 # nosec B608 " WHERE account_id = %s AND reporting_obligation_id = %s" " ORDER BY created_at, reporting_revision_id", (account_id, reporting_obligation_id), ) ).fetchall() return tuple(_revision_from_row(row) for row in rows) async def get_provisional_acquisition( self, *, account_id: str, reporting_obligation_id: str, ordinal: int ) -> ProvisionalAcquisition | None: async with self._connection() as connection: row = await ( await connection.execute( "SELECT payload FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", (account_id, reporting_obligation_id, ordinal), ) ).fetchone() return ProvisionalAcquisition.from_wire(row[0]) if row is not None else None async def reserve_provisional_acquisition( self, acquisition: ProvisionalAcquisition ) -> ProvisionalAcquisition: async with self.transaction(), self._connection() as connection: await self._lock_account(connection, acquisition.account_id) await self._require_provisional_schema(connection) key = (acquisition.account_id, acquisition.obligation_id, acquisition.ordinal) existing = await ( await connection.execute( "SELECT payload FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", key, ) ).fetchone() if existing is not None: return ProvisionalAcquisition.from_wire(existing[0]) obligation = await self.get_obligation( account_id=acquisition.account_id, reporting_obligation_id=acquisition.obligation_id, ) if obligation is None: raise LedgerConflictError("OBLIGATION_NOT_FOUND", "unknown observation obligation") if not acquisition.binds(obligation): raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition generation differs") duplicate = await ( await connection.execute( "SELECT 1 FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND source_execution_key=%s", (acquisition.account_id, acquisition.execution_key), ) ).fetchone() if duplicate is not None: raise LedgerConflictError( "OBSERVATION_CONFLICT", "execution key is already reserved" ) checkpoint = await self.get_restatement_checkpoint( account_id=acquisition.account_id, reporting_obligation_id=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") await connection.execute( "INSERT INTO reporting_provisional_acquisitions" " (account_id,reporting_obligation_id,ordinal,source_execution_key,payload)" " VALUES (%s,%s,%s,%s,%s::jsonb)", (*key, acquisition.execution_key, _json(acquisition.to_wire())), ) return ProvisionalAcquisition.from_wire(acquisition.to_wire()) async def get_provisional_observation( self, *, account_id: str, reporting_obligation_id: str ) -> ProvisionalObservation | None: async with self._connection() as connection: await self._require_provisional_schema(connection) row = await ( await connection.execute( "SELECT payload FROM reporting_provisional_observations" " WHERE account_id=%s AND reporting_obligation_id=%s" " ORDER BY ordinal DESC LIMIT 1", (account_id, reporting_obligation_id), ) ).fetchone() return ProvisionalObservation.from_wire(row[0]) if row else 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.transaction(), self._connection() as connection: await self._lock_account(connection, acquisition.account_id) await self._require_provisional_schema(connection) existing = await ( await connection.execute( "SELECT reporting_revision_id,payload FROM reporting_provisional_observations" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", key, ) ).fetchone() if existing is not None: retained_observation = ProvisionalObservation.from_wire(existing[1]) if ( retained_observation.acquisition != acquisition or existing[0] != observation.revision_id ): raise LedgerConflictError("OBSERVATION_CONFLICT", "observation replay differs") retained_revision = await self.get_revision( account_id=acquisition.account_id, reporting_revision_id=existing[0] ) if retained_revision is None: raise LedgerConflictError( "HISTORY_UNAVAILABLE", "observation revision is missing" ) # Reuse the immutable publication replay checks while holding # the observation's account lock and transaction. return await self.commit_revision(revision, rows) reserved = await ( await connection.execute( "SELECT payload FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", key, ) ).fetchone() if reserved is None or ProvisionalAcquisition.from_wire(reserved[0]) != acquisition: raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition was not reserved") checkpoint = await self.get_restatement_checkpoint( account_id=acquisition.account_id, reporting_obligation_id=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, ) ) await connection.execute( "INSERT INTO reporting_provisional_observations" " (account_id,reporting_obligation_id,ordinal,reporting_revision_id,payload)" " VALUES (%s,%s,%s,%s,%s::jsonb)", (*key, revision.reporting_revision_id, _json(observation.to_wire())), ) return committed async def get_retry_schedule(self, *, scope_key: str) -> RetryScheduleEntry | None: async with self._connection() as connection: row = await ( await connection.execute( "SELECT retry_not_before, attempt, blocked, recorded_at" " FROM adcp_reporting_producer_retry_schedules" " WHERE scope_key = %s", (scope_key,), ) ).fetchone() return ( RetryScheduleEntry(scope_key, _utc(row[0]), int(row[1]), row[2], _utc(row[3])) if row else None ) async def record_retry_schedule(self, entry: RetryScheduleEntry) -> RetryScheduleEntry: async with self._connection() as connection: row = await ( await connection.execute( "INSERT INTO adcp_reporting_producer_retry_schedules" " (scope_key, retry_not_before, attempt, blocked, recorded_at)" " VALUES (%s, %s, %s, %s, %s)" " ON CONFLICT (scope_key) DO UPDATE SET" " retry_not_before = CASE WHEN EXCLUDED.attempt = 0 OR EXCLUDED.blocked" " OR adcp_reporting_producer_retry_schedules.blocked OR %s" " THEN EXCLUDED.retry_not_before ELSE GREATEST(" " adcp_reporting_producer_retry_schedules.retry_not_before," " EXCLUDED.retry_not_before) END," " attempt = EXCLUDED.attempt, blocked = EXCLUDED.blocked," " recorded_at = EXCLUDED.recorded_at" " WHERE EXCLUDED.recorded_at >=" " adcp_reporting_producer_retry_schedules.recorded_at" " AND (EXCLUDED.attempt = 0 OR EXCLUDED.blocked OR %s" " OR (NOT adcp_reporting_producer_retry_schedules.blocked AND" " adcp_reporting_producer_retry_schedules.attempt <= EXCLUDED.attempt))" " RETURNING retry_not_before, attempt, blocked, recorded_at", ( entry.scope_key, entry.retry_not_before, entry.attempt, entry.blocked, entry.recorded_at, entry.replayed, entry.replayed, ), ) ).fetchone() if row is None: stored = await self.get_retry_schedule(scope_key=entry.scope_key) assert stored is not None return stored return RetryScheduleEntry(entry.scope_key, _utc(row[0]), int(row[1]), row[2], _utc(row[3])) async def get_restatement_checkpoint( self, *, account_id: str, reporting_obligation_id: str ) -> RestatementCheckpoint | None: async with self._connection() as connection: row = await ( await connection.execute( "SELECT account_id, reporting_obligation_id, checked_at," " next_observation, provisional_until" " FROM reporting_restatement_checkpoints" " WHERE account_id = %s AND reporting_obligation_id = %s", (account_id, reporting_obligation_id), ) ).fetchone() if row is None: return None return RestatementCheckpoint( account_id=row[0], reporting_obligation_id=row[1], checked_at=_utc(row[2]), next_observation=int(row[3]), provisional_until=_utc(row[4]) if row[4] else None, ) async def record_restatement_checkpoint( self, checkpoint: RestatementCheckpoint ) -> RestatementCheckpoint: async with self._connection() as connection: obligation = await ( await connection.execute( "SELECT 1 FROM reporting_obligations" " WHERE reporting_obligation_id = %s AND account_id = %s", (checkpoint.reporting_obligation_id, checkpoint.account_id), ) ).fetchone() if obligation is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "a restatement checkpoint must attach to an obligation for this account", ) row = await ( await connection.execute( "INSERT INTO reporting_restatement_checkpoints" " (account_id, reporting_obligation_id, checked_at, next_observation," " provisional_until)" " VALUES (%s, %s, %s, %s, %s)" " ON CONFLICT (reporting_obligation_id) DO UPDATE SET" " checked_at = EXCLUDED.checked_at," " next_observation = EXCLUDED.next_observation," " provisional_until = EXCLUDED.provisional_until" " WHERE reporting_restatement_checkpoints.account_id = EXCLUDED.account_id" " AND (reporting_restatement_checkpoints.next_observation" " < EXCLUDED.next_observation" " OR (reporting_restatement_checkpoints.next_observation" " = EXCLUDED.next_observation" " AND reporting_restatement_checkpoints.checked_at" " < EXCLUDED.checked_at))" " RETURNING account_id, reporting_obligation_id, checked_at," " next_observation, provisional_until", ( checkpoint.account_id, checkpoint.reporting_obligation_id, checkpoint.checked_at, checkpoint.next_observation, checkpoint.provisional_until, ), ) ).fetchone() if row is None: stored = await self.get_restatement_checkpoint( account_id=checkpoint.account_id, reporting_obligation_id=checkpoint.reporting_obligation_id, ) if stored is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "a restatement checkpoint must attach to an obligation for this account", ) return stored return RestatementCheckpoint( account_id=row[0], reporting_obligation_id=row[1], checked_at=_utc(row[2]), next_observation=int(row[3]), provisional_until=_utc(row[4]) if row[4] else None, ) async def get_revision( self, *, account_id: str, reporting_revision_id: str ) -> ReportingRevisionRecord | None: async with self._connection() as connection: row = await ( await connection.execute( f"SELECT {_REVISION_COLUMNS} FROM reporting_revisions" # noqa: S608 # nosec B608 " WHERE account_id = %s AND reporting_revision_id = %s", (account_id, reporting_revision_id), ) ).fetchone() return _revision_from_row(row) if row else None async def read_revision_rows( self, *, account_id: str, reporting_revision_id: str, cursor: str | None = None, limit: int = 500, ) -> ReportingRowPage: from adcp.reporting.ledger.store import revision_row_offset offset = revision_row_offset(cursor, reporting_revision_id, limit) async with self._connection() as connection: owned = await ( await connection.execute( "SELECT row_count FROM reporting_revisions" " WHERE account_id = %s AND reporting_revision_id = %s", (account_id, reporting_revision_id), ) ).fetchone() if owned is None: raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account") rows = await ( await connection.execute( "SELECT row_payload FROM reporting_revision_rows" " WHERE reporting_revision_id = %s" " ORDER BY ordinal OFFSET %s LIMIT %s", (reporting_revision_id, offset, limit), ) ).fetchall() total = int(owned[0]) has_more = offset + limit < total return ReportingRowPage( reporting_revision_id=reporting_revision_id, rows=tuple(row[0] for row in rows), total_count=total, 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._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) existing = await ( await connection.execute( "SELECT readable, reporting_obligation_id FROM reporting_revisions" " WHERE account_id = %s AND reporting_revision_id = %s FOR UPDATE", (account_id, reporting_revision_id), ) ).fetchone() if existing is None: raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account") if existing[0] == readable: return await connection.execute( "UPDATE reporting_revisions SET readable = %s" " WHERE account_id = %s AND reporting_revision_id = %s", (readable, account_id, reporting_revision_id), ) if self._notifications_enabled: obligation = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # nosec B608 " WHERE account_id = %s AND reporting_obligation_id = %s", (account_id, existing[1]), ) ).fetchone() assert obligation is not None await self._dirty_status( connection, ReportingStatusScope.for_obligation(_obligation_from_row(obligation)), "readability", ReportingStatusEvidence( "revision", reporting_revision_id, readable=existing[0] ), ReportingStatusEvidence("revision", reporting_revision_id, readable=readable), ) # -- adjustments ------------------------------------------------------ async def commit_adjustment( self, adjustment: ReportingAdjustmentRecord ) -> ReportingAdjustmentRecord: async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, adjustment.account_id) existing = await ( await connection.execute( f"SELECT {_ADJUSTMENT_COLUMNS} FROM reporting_adjustments" # noqa: S608 # nosec B608 " WHERE reporting_adjustment_id = %s AND account_id = %s", (adjustment.reporting_adjustment_id, adjustment.account_id), ) ).fetchone() if existing is not None: stored = _adjustment_from_row(existing) if stored != adjustment: raise LedgerConflictError( "ADJUSTMENT_IMMUTABLE", "adjustment content is immutable" ) return stored revision = await ( await connection.execute( "SELECT finality, reporting_obligation_id FROM reporting_revisions" " WHERE reporting_revision_id = %s AND account_id = %s", (adjustment.adjusts_reporting_revision_id, adjustment.account_id), ) ).fetchone() if revision is None: raise LedgerConflictError( "REVISION_NOT_FOUND", "an adjustment must name a committed revision" ) if revision[0] != "official": raise LedgerConflictError( "ADJUSTMENT_REQUIRES_OFFICIAL", "adjustments correct an official revision; restate a snapshot with a " "superseding snapshot revision instead", ) obligation = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE reporting_obligation_id = %s AND account_id = %s", (revision[1], adjustment.account_id), ) ).fetchone() if obligation is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "the adjustment target has no obligation" ) await self._require_writable_generation( connection, _obligation_from_row(obligation).generation_key ) validate_adjustment_currency(_obligation_from_row(obligation), adjustment) inserted = await ( await connection.execute( "INSERT INTO reporting_adjustments" " (reporting_adjustment_id, account_id, adjusts_reporting_revision_id," " reason_code, reason_detail, accounting_period_start," " accounting_period_end, control_total_deltas, correction_observed_at," " created_at, managed_control_total_deltas)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s, %s, %s::jsonb)" " ON CONFLICT (reporting_adjustment_id) DO NOTHING" " RETURNING reporting_adjustment_id", ( adjustment.reporting_adjustment_id, adjustment.account_id, adjustment.adjusts_reporting_revision_id, adjustment.reason_code, adjustment.reason_detail, adjustment.accounting_period_start, adjustment.accounting_period_end, _json([[name, value] for name, value in adjustment.control_total_deltas]), adjustment.correction_observed_at, adjustment.created_at, ( _json( [item.to_wire() for item in adjustment.managed_control_total_deltas] ) if adjustment.managed_control_total_deltas is not None else None ), ), ) ).fetchone() if inserted is None: raise LedgerConflictError("ADJUSTMENT_UNAVAILABLE", "adjustment is unavailable") await self._append_change( connection, adjustment.account_id, "adjustment", adjustment.reporting_adjustment_id, consumer_id=_obligation_from_row(obligation).consumer_id, ) if self._notifications_enabled: from adcp.reporting.ledger.notification_events import adjustment_event await self._record_notification( connection, adjustment_event( adjustment, await self._notification_now(connection), consumer_id=_obligation_from_row(obligation).consumer_id, ), ) await self._dirty_status( connection, ReportingStatusScope.for_obligation(_obligation_from_row(obligation)), "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, ...]: if not reporting_revision_ids: return () async with self._connection() as connection: rows = await ( await connection.execute( f"SELECT {_ADJUSTMENT_COLUMNS} FROM reporting_adjustments" # noqa: S608 # nosec B608 " WHERE account_id = %s AND adjusts_reporting_revision_id = ANY(%s)" " ORDER BY created_at, reporting_adjustment_id", (account_id, list(reporting_revision_ids)), ) ).fetchall() return tuple(_adjustment_from_row(row) for row in rows) # -- consumer status (preview) ---------------------------------------- async def record_consumer_status( self, status: ConsumerStatusRecord ) -> tuple[ConsumerStatusRecord, bool]: from adcp.reporting.evidence import consumer_reference consumer_reference(status.consumer_id) key = status.generation_key digest = _fingerprint(_consumer_status_payload(status)) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, status.account_id) replay = await self._replay(connection, status, digest) if replay is not None: return replay, False from adcp.reporting.ledger.status_snapshot import ( read_snapshot_on, validate_status_evidence, ) validate_status_evidence( status, await read_snapshot_on( connection, account_id=status.account_id, clock=self._clock, include_issue_scopes=await self._issue_scope_storage_on(connection), ), ) leaf = await ( await connection.execute( "SELECT reporting_status_id FROM reporting_consumer_statuses" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND report_definition_id = %s" " AND period_start = %s AND period_end = %s AND superseded = FALSE", ( key.account_id, status.consumer_id, key.delivery_config_id, key.delivery_config_version, status.report_definition_id, status.period_start, status.period_end, ), ) ).fetchone() if status.supersedes_reporting_status_id: if leaf is None or leaf[0] != 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", ) await connection.execute( "UPDATE reporting_consumer_statuses SET superseded = TRUE" " WHERE account_id = %s AND consumer_id = %s AND reporting_status_id = %s", (key.account_id, status.consumer_id, status.supersedes_reporting_status_id), ) elif leaf is not None: raise LedgerConflictError( "STATUS_SUPERSEDES_REQUIRED", "this chain already has a current statement; a new statement must " "explicitly supersede it", ) try: await connection.execute( "INSERT INTO reporting_consumer_statuses" " (reporting_status_id, account_id, consumer_id, delivery_config_id," " delivery_config_version, report_definition_id, period_start, period_end," " period_source_timezone, consumer_status, status_as_of, recorded_at," " supersedes_reporting_status_id, reporting_obligation_id," " reporting_revision_id, observed_revision_content_sha256, failure_code," " mismatch_code, consumer_commit_ref, seller_ledger_snapshot_id," " seller_ledger_as_of, superseded, content_sha256)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s," " %s, %s, %s, %s, %s, FALSE, %s)", ( status.reporting_status_id, status.account_id, status.consumer_id, status.delivery_config_id, status.delivery_config_version, status.report_definition_id, status.period_start, status.period_end, status.period_source_timezone, status.consumer_status, status.status_as_of, status.recorded_at, status.supersedes_reporting_status_id, status.reporting_obligation_id, status.reporting_revision_id, status.observed_revision_content_sha256, status.failure_code, status.mismatch_code, status.consumer_commit_ref, status.seller_ledger_snapshot_id, status.seller_ledger_as_of, digest, ), ) except Exception as error: raise _translate_integrity_error(error) from error await self._append_change( connection, status.account_id, "consumer_status", status.reporting_status_id, consumer_id=status.consumer_id, ) from adcp.reporting.ledger.status_snapshot import settle_snapshot_on await settle_snapshot_on(self, connection, account_id=status.account_id) if self._notifications_enabled: await self._dirty_status( connection, ReportingStatusScope( status.account_id, status.generation_key, status.reporting_obligation_id, status.consumer_id, ), "consumer_status", ReportingStatusEvidence("consumer_status", leaf[0]) if leaf 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]: """Optional participant: the record and lifecycle share the source transaction.""" return await self.record_consumer_status(status) async def read_status_snapshot(self, *, caller: ReportingCaller) -> ReportingStatusSnapshot: from adcp.reporting.ledger.delivery_models import ReportingDeliveryPrincipal from adcp.reporting.ledger.status_snapshot import settle_snapshot_on from adcp.reporting.materializer.capture import private_snapshot async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, caller.account_id) return private_snapshot( await settle_snapshot_on(self, connection, account_id=caller.account_id), ReportingDeliveryPrincipal(caller.account_id, caller.consumer_id), ) @staticmethod async def _get_consumer_status( connection: Any, status: ConsumerStatusRecord ) -> ConsumerStatusRecord | None: row = await ( await connection.execute( # Bare columns: this query has no `s` alias to qualify against. f"SELECT {_STATUS_COLUMNS_BARE} FROM reporting_consumer_statuses" # noqa: S608 # nosec B608 " WHERE reporting_status_id = %s AND account_id = %s AND consumer_id = %s", (status.reporting_status_id, status.account_id, status.consumer_id), ) ).fetchone() return _status_from_row(row) if row else None async def resolve_consumer_status_replay( self, status: ConsumerStatusRecord ) -> ConsumerStatusRecord | None: async with self._connection() as connection: return await self._replay( connection, status, _fingerprint(_consumer_status_payload(status)) ) async def _replay( self, connection: Any, status: ConsumerStatusRecord, digest: str ) -> ConsumerStatusRecord | None: existing = await ( await connection.execute( "SELECT content_sha256 FROM reporting_consumer_statuses" " WHERE reporting_status_id = %s AND account_id = %s AND consumer_id = %s", (status.reporting_status_id, status.account_id, status.consumer_id), ) ).fetchone() if existing is None: return None if existing[0] != digest: raise LedgerConflictError( "STATUS_IDENTITY_CONFLICT", f"reporting_status_id {status.reporting_status_id} was already " "recorded with different content", ) stored = await self._get_consumer_status(connection, status) assert stored is not None return stored async def list_consumer_statuses( self, *, account_id: str, consumer_id: str, reporting_obligation_ids: Sequence[str] | None = None, ) -> tuple[ConsumerStatusRecord, ...]: # Attaches by obligation id *or* by the exact logical period key, so a # chain filed before the obligation existed is not lost, forked, or # reset when the seller later repairs the missing obligation. clause = "" params: list[Any] = [account_id, consumer_id] if reporting_obligation_ids is not None: clause = ( " AND EXISTS (" " SELECT 1 FROM reporting_obligations o" " WHERE o.reporting_obligation_id = ANY(%s)" " AND o.account_id = s.account_id" " AND (o.reporting_obligation_id = s.reporting_obligation_id OR (" " o.delivery_config_id = s.delivery_config_id" " AND o.delivery_config_version = s.delivery_config_version" " AND o.report_definition_id = s.report_definition_id" " AND o.period_start = s.period_start" " AND o.period_end = s.period_end)))" ) params.append(list(reporting_obligation_ids)) async with self._connection() as connection: rows = await ( await connection.execute( f"SELECT {_STATUS_COLUMNS} FROM reporting_consumer_statuses s" # noqa: S608 # nosec B608 f" WHERE s.account_id = %s AND s.consumer_id = %s{clause}" " ORDER BY s.recorded_at, s.reporting_status_id", tuple(params), ) ).fetchall() return tuple(_status_from_row(row) for row in rows) # -- 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._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) return await self._ensure_issue_opened_on( connection, issue_key=issue_key, account_id=account_id, consumer_id=consumer_id, observed_at=observed_at, status_scope=status_scope, ) async def _ensure_issue_opened_on( self, connection: Any, *, issue_key: str, account_id: str, consumer_id: str | None, observed_at: datetime, status_scope: ReportingStatusScope | None = None, enqueue: bool = True, ) -> ReportingIssueLifecycle: live = await self._live_issue(connection, issue_key, account_id) if live is not None: if live.consumer_id != consumer_id: raise ReportingNotificationError("invalid_status_scope") await self._dirty_issue(connection, live, status_scope, live, enqueue=enqueue) return live row = await ( await connection.execute( "SELECT generation, consumer_id, issue_id FROM reporting_issue_lifecycle" " WHERE account_id = %s AND issue_key = %s ORDER BY generation DESC LIMIT 1", (account_id, issue_key), ) ).fetchone() if row is not None: if row[1] != consumer_id: raise ReportingNotificationError("invalid_status_scope") previous = ( await ( await connection.execute( "SELECT scope FROM reporting_issue_status_scopes" " WHERE account_id=%s AND issue_id=%s", (account_id, row[2]), ) ).fetchone() if await self._issue_scope_storage_on(connection) else None ) previous_scope = decode_status_scope(previous[0]) if previous else None if status_scope is not None: validate_scope_refinement(previous_scope, status_scope) else: status_scope = previous_scope generation = int(row[0] if row else 0) + 1 issue_id = issue_id_for_occurrence(issue_key, generation) await connection.execute( "INSERT INTO reporting_issue_lifecycle" " (issue_key, account_id, generation, issue_id, consumer_id, opened_at, issue_state)" " VALUES (%s, %s, %s, %s, %s, %s, 'open')", (issue_key, account_id, generation, issue_id, consumer_id, _utc(observed_at)), ) record = ReportingIssueLifecycle( issue_key=issue_key, issue_id=issue_id, account_id=account_id, consumer_id=consumer_id, opened_at=_utc(observed_at), generation=generation, ) await self._dirty_issue(connection, record, status_scope, enqueue=enqueue) 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: 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", ) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) live = await self._live_issue(connection, issue_key, account_id) if live is None: 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 read_snapshot_on snapshot = await read_snapshot_on( connection, account_id=account_id, as_of=_utc(at), include_issue_scopes=await self._issue_scope_storage_on(connection), ) live = bind_mismatch_waiver(snapshot, live) await _save_waiver_binding(connection, live) waived_at = ( live.retired_at if live.waived_reporting_status_id is not None and live.retired_at is not None else _utc(at) ) await connection.execute( "UPDATE reporting_issue_lifecycle" " SET issue_state = %s," " external_ref = COALESCE(%s, external_ref)," " retired_at = CASE WHEN %s = 'waived' THEN %s ELSE retired_at END" " WHERE account_id = %s AND issue_key = %s AND generation = %s", ( state, external_ref, state, waived_at, account_id, issue_key, live.generation, ), ) refreshed = await self._issue_row(connection, issue_key, account_id, live.generation) assert refreshed is not None # Derive the no-op from the resulting row 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. await self._dirty_issue(connection, refreshed, status_scope, live) return refreshed async def retire_issue( self, *, issue_key: str, account_id: str, at: datetime, status_scope: ReportingStatusScope | None = None, ) -> ReportingIssueLifecycle | None: async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) return await self._retire_issue_on( connection, issue_key=issue_key, account_id=account_id, at=at, status_scope=status_scope, ) async def _retire_issue_on( self, connection: Any, *, issue_key: str, account_id: str, at: datetime, status_scope: ReportingStatusScope | None = None, enqueue: bool = True, ) -> ReportingIssueLifecycle | None: live = await self._live_issue(connection, issue_key, account_id) if live is None or not issue_is_retirable(live.issue_state): return None if live.issue_state != "waived": check_issue_state_transition(live.issue_state, "resolved") await connection.execute( "UPDATE reporting_issue_lifecycle SET issue_state = 'resolved', retired_at = %s" " WHERE account_id = %s AND issue_key = %s AND generation = %s", (_utc(at), account_id, issue_key, live.generation), ) retired = await self._issue_row(connection, issue_key, account_id, live.generation) assert retired is not None await self._dirty_issue(connection, retired, status_scope, live, enqueue=enqueue) return retired async def get_issue(self, *, issue_key: str, account_id: str) -> ReportingIssueLifecycle | None: async with self._connection() as connection: return await self._live_issue(connection, issue_key, account_id) @staticmethod async def _live_issue( connection: Any, issue_key: str, account_id: str ) -> ReportingIssueLifecycle | None: row = await ( await connection.execute( f"SELECT {_ISSUE_COLUMNS} FROM reporting_issue_lifecycle" # noqa: S608 # nosec B608 " WHERE account_id = %s AND issue_key = %s" " AND issue_state IN ('open', 'acknowledged', 'waived')", (account_id, issue_key), ) ).fetchone() return await _with_waiver_binding(connection, _issue_from_row(row)) if row else None @staticmethod async def _issue_row( connection: Any, issue_key: str, account_id: str, generation: int ) -> ReportingIssueLifecycle | None: row = await ( await connection.execute( f"SELECT {_ISSUE_COLUMNS} FROM reporting_issue_lifecycle" # noqa: S608 # nosec B608 " WHERE account_id = %s AND issue_key = %s AND generation = %s", (account_id, issue_key, generation), ) ).fetchone() return await _with_waiver_binding(connection, _issue_from_row(row)) if row else None # -- snapshots -------------------------------------------------------- async def open_snapshot( self, *, caller: ReportingCaller, filters_fingerprint: str ) -> LedgerSnapshot: account_id = caller.account_id async with self._connection() as connection: row = await ( await connection.execute( "SELECT COALESCE(MAX(seq), 0), now() FROM reporting_ledger_changes" " WHERE account_id=%s AND consumer_id=%s AND record_kind IN" " ('obligation','revision','adjustment','consumer_status')", (account_id, caller.consumer_id), ) ).fetchone() assert row is not None max_sequence = int(row[0]) as_of = _utc(self._clock()) if self._clock is not None else _utc(row[1]) 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 async with self._connection() as connection: rows = await ( await connection.execute( "SELECT seq, record_kind, record_id FROM reporting_ledger_changes" " WHERE account_id = %s AND consumer_id=%s AND seq > %s AND seq <= %s" " ORDER BY seq", (snapshot.account_id, consumer_id, lower, snapshot.max_sequence), ) ).fetchall() selected: list[tuple[int, LedgerRecordKind, Any]] = [] for _seq, kind, record_id in rows: record = await self._resolve( connection, snapshot.account_id, kind, record_id, consumer_id=consumer_id ) if record is None: continue if not await self._in_scope( connection, kind, record, delivery_config_ids, media_buy_ids, consumer_id, feed_purposes, period_start, period_end, ): continue selected.append((_seq, kind, record)) window = selected[offset : offset + limit] has_more = offset + limit < len(selected) 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(selected), has_more=has_more, cursor=( encode_cursor({"snapshot": snapshot.snapshot_id, "offset": offset + limit}) if has_more else None ), ) async def _resolve( self, connection: Any, account_id: str, kind: str, record_id: str, *, consumer_id: str | None = None, ) -> Any | None: if kind not in _RESOLVERS: return None # Higher-tier records are never projected by Core. table, columns, key, builder = _RESOLVERS[kind] consumer_filter = " AND consumer_id=%s" if kind == "consumer_status" else "" parameters = ( (account_id, record_id, consumer_id) if consumer_filter else (account_id, record_id) ) row = await ( await connection.execute( f"SELECT {columns} FROM {table} WHERE account_id = %s AND {key} = %s" # noqa: S608 # nosec B608 f"{consumer_filter}", # nosec B608 parameters, ) ).fetchone() return builder(row) if row else None async def _in_scope( self, connection: Any, kind: str, record: Any, delivery_config_ids: Sequence[str] | None, media_buy_ids: Sequence[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 never disclosed. if consumer_id is None or record.consumer_id != consumer_id: return False owner = record start, end = record.period_start, record.period_end else: owner = await self._obligation_for(connection, kind, record) if owner is None or owner.consumer_id != consumer_id: return False start, end = owner.period.start, owner.period.end if media_buy_ids and not set(media_buy_ids).intersection(owner.media_buy_ids): return False configurations = await self._list_configurations_on( connection, account_id=record.account_id, consumer_id=owner.consumer_id, delivery_config_ids=[owner.delivery_config_id], ) configuration = next( (c for c in configurations if c.generation_key == owner.generation_key), None ) if configuration is None: return False return configuration_selected( configuration, delivery_config_ids=delivery_config_ids or (), media_buy_ids=media_buy_ids or (), feed_purposes=feed_purposes or (), ) and period_selected(start, end, period_start, period_end) async def _obligation_for( self, connection: Any, kind: str, record: Any ) -> ReportingObligationRecord | None: if kind == "obligation": found: ReportingObligationRecord = record return found if kind == "revision": return await self._resolve( connection, record.account_id, "obligation", record.reporting_obligation_id ) revision = await self._resolve( connection, record.account_id, "revision", record.adjusts_reporting_revision_id ) if revision is None: return None return await self._resolve( connection, revision.account_id, "obligation", revision.reporting_obligation_id ) # -- leasing ---------------------------------------------------------- async def _period_close_generations_on( self, connection: Any, identity: tuple[str, str, str, int] ) -> list[tuple[str, str, int]]: # A base-tier store may share a schema installed by another participant. # Probe on this transaction without caching schema presence or validating # every catalog object on the lease path. Call only after the account lock. present = await ( await connection.execute( "SELECT to_regclass('reporting_materializer_candidates') IS NOT NULL" ) ).fetchone() if not present[0]: return [] rows = await ( await connection.execute( "SELECT consumer_id,reporting_obligation_id,generation" " FROM reporting_materializer_candidates WHERE account_id=%s" " AND consumer_id=%s AND delivery_config_id=%s AND delivery_config_version=%s", identity, ) ).fetchall() return [(row[0], row[1], row[2]) for row in rows] async def _restore_period_close_generations_on( self, connection: Any, identity: tuple[str, str, str, int], snapshot: list[tuple[str, str, int]], ) -> None: # These SDK writes change only lease fields. All supported source writers # hold the same account lock, so no source change can intervene. Restore # only the captured generation plus the known single trigger increment; # never decrement blindly if the trigger was disabled or did not fire. # Unexpected generation changes remain fenced for schema/target validation. # Wakeups and due_at/reason retain the production tier's existing behavior. if not snapshot: return await connection.execute( "UPDATE reporting_materializer_candidates c SET generation=s.generation" " FROM unnest(%s::text[],%s::text[],%s::bigint[])" " AS s(consumer_id,reporting_obligation_id,generation)" " WHERE c.account_id=%s AND c.consumer_id=%s AND c.delivery_config_id=%s" " AND c.delivery_config_version=%s AND c.consumer_id=s.consumer_id" " AND c.reporting_obligation_id=s.reporting_obligation_id" " AND c.generation=s.generation+1", ( [row[0] for row in snapshot], [row[1] for row in snapshot], [row[2] for row in snapshot], *identity, ), ) async def lease_period_close( self, *, worker_id: str, now: datetime, lease_seconds: float ) -> LeasedConfiguration | None: from psycopg.errors import LockNotAvailable from psycopg.pq import TransactionStatus moment = _utc(now) expires = moment + timedelta(seconds=lease_seconds) after = self._period_close_sample wait_until = None following = None result = None async with self._connection() as connection: # Never add a blocking account edge while a caller transaction may # already own locks. Standalone turns may queue only before the first # account acquisition, including acquisitions that find a stale row. can_wait = connection.info.transaction_status == TransactionStatus.IDLE async with connection.transaction(): await self._bound_lock_waits(connection) # A materializer installed by any participant adds an AFTER UPDATE # trigger taking the account lock. Sample without row locks, then take # that account lock before locking a configuration. The base store can # share that schema, so the ordering belongs here, not only on a tier. for _ in range(2): continuation = ( " AND (COALESCE(t.lease_turn,0)," " COALESCE(c.lease_expires_at,'-infinity'::timestamptz)," " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_confi" "g_version)" " > (%s,COALESCE(%s::timestamptz,'-infinity'::timestamptz),%s,%s,%s,%s)" if after is not None else "" ) query = ( "SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.delivery" "_config_version," # nosec B608 " COALESCE(t.lease_turn,0),c.lease_expires_at FROM" " reporting_configurations c" " LEFT JOIN adcp_reporting_configuration_lease_turns t" " ON (t.account_id,t.consumer_id,t.delivery_config_id,t.delivery_c" "onfig_version)" " = (c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_co" "nfig_version)" " WHERE NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lea" "se_expires_at<=%s)" + continuation + " ORDER BY COALESCE(t.lease_turn,0),c.lease_expires_at NULLS FIRST," " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_confi" "g_version LIMIT 32" ) rows = await ( await connection.execute(query, (moment, *(after or ()))) ).fetchall() if rows or after is None: break after = None for candidate in rows: row = candidate[:4] following = (candidate[4], candidate[5], *row) locked = await ( await connection.execute( "SELECT pg_try_advisory_xact_lock(hashtext('adcp.reporting:' || %s))", (row[0],), ) ).fetchone() if not locked[0]: if not can_wait: continue # A try-lock alone can starve behind a continuous queue of # ordinary writers, even though each writer commits promptly. # Join that queue briefly, sharing one wait budget per turn. # The savepoint rolls back a timed-out wait and its SET LOCAL; # successful acquisition restores the caller's lock timeout. if wait_until is None: wait_until = monotonic() + _LEASE_ACCOUNT_WAIT_SECONDS remaining_ms = int((wait_until - monotonic()) * 1000) if remaining_ms <= 0: continue try: async with connection.transaction(): previous, configured_ms = await ( await connection.execute( "SELECT current_setting('lock_timeout'),setting::integer" " FROM pg_settings WHERE name='lock_timeout'" ) ).fetchone() limit_ms = min(remaining_ms, configured_ms or remaining_ms) await connection.execute( "SELECT set_config('lock_timeout',%s,true)", (f"{limit_ms}ms",) ) await connection.execute( "SELECT pg_advisory_xact_lock(hashtext('adcp.reporting:' ||" " %s))", (row[0],), ) await connection.execute( "SELECT set_config('lock_timeout',%s,true)", (previous,) ) except LockNotAvailable: continue can_wait = False snapshot = await self._period_close_generations_on(connection, tuple(row)) acquired = await ( await connection.execute( "UPDATE reporting_configurations SET lease_worker_id = %s," " lease_expires_at = %s" " WHERE (account_id,consumer_id,delivery_config_id,delivery_config" "_version) = (" " SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.deliver" "y_config_version" " FROM reporting_configurations c WHERE c.account_id=%s" " AND c.consumer_id=%s AND c.delivery_config_id=%s AND c.delivery_" "config_version=%s" " AND NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lease" "_expires_at<=%s)" " FOR UPDATE OF c SKIP LOCKED) RETURNING account_id", (worker_id, expires, *row, moment), ) ).fetchone() if acquired is None: continue await self._restore_period_close_generations_on( connection, tuple(row), snapshot ) await connection.execute( "INSERT INTO adcp_reporting_configuration_lease_turns" " (account_id,consumer_id,delivery_config_id,delivery_config_versi" "on,lease_turn)" " VALUES(%s,%s,%s,%s,nextval('adcp_reporting_configuration_lease_t" "urn_seq'))" " ON CONFLICT(account_id,consumer_id,delivery_config_id,delivery_c" "onfig_version)" " DO UPDATE SET" " lease_turn=nextval('adcp_reporting_configuration_lease_turn_seq')", tuple(row), ) result = LeasedConfiguration(row[0], row[1], row[2], row[3], expires) break # Only sampling uses this hint: leases and fairness ranks stay transactional. # Continue past a busy prefix on the next turn; a successful turn returns to # durable fairness order. Empty tails wrap at most once without row locks. self._period_close_sample = following if result is None else None return result async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) -> None: key = lease.generation_key identity = ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, key.account_id) snapshot = await self._period_close_generations_on(connection, identity) released = await ( await connection.execute( "UPDATE reporting_configurations SET lease_worker_id = NULL," " lease_expires_at = NULL" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND lease_worker_id = %s" " AND lease_expires_at = %s RETURNING account_id", (*identity, worker_id, lease.lease_expires_at), ) ).fetchone() if released is not None: await self._restore_period_close_generations_on(connection, identity, snapshot)Durable reporting ledger over a caller-supplied connection pool.
Subclasses
Class variables
var is_durable : ClassVar[bool]
Methods
async def commit_adjustment(self, adjustment: ReportingAdjustmentRecord) ‑> ReportingAdjustmentRecord-
Expand source code
async def commit_adjustment( self, adjustment: ReportingAdjustmentRecord ) -> ReportingAdjustmentRecord: async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, adjustment.account_id) existing = await ( await connection.execute( f"SELECT {_ADJUSTMENT_COLUMNS} FROM reporting_adjustments" # noqa: S608 # nosec B608 " WHERE reporting_adjustment_id = %s AND account_id = %s", (adjustment.reporting_adjustment_id, adjustment.account_id), ) ).fetchone() if existing is not None: stored = _adjustment_from_row(existing) if stored != adjustment: raise LedgerConflictError( "ADJUSTMENT_IMMUTABLE", "adjustment content is immutable" ) return stored revision = await ( await connection.execute( "SELECT finality, reporting_obligation_id FROM reporting_revisions" " WHERE reporting_revision_id = %s AND account_id = %s", (adjustment.adjusts_reporting_revision_id, adjustment.account_id), ) ).fetchone() if revision is None: raise LedgerConflictError( "REVISION_NOT_FOUND", "an adjustment must name a committed revision" ) if revision[0] != "official": raise LedgerConflictError( "ADJUSTMENT_REQUIRES_OFFICIAL", "adjustments correct an official revision; restate a snapshot with a " "superseding snapshot revision instead", ) obligation = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE reporting_obligation_id = %s AND account_id = %s", (revision[1], adjustment.account_id), ) ).fetchone() if obligation is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "the adjustment target has no obligation" ) await self._require_writable_generation( connection, _obligation_from_row(obligation).generation_key ) validate_adjustment_currency(_obligation_from_row(obligation), adjustment) inserted = await ( await connection.execute( "INSERT INTO reporting_adjustments" " (reporting_adjustment_id, account_id, adjusts_reporting_revision_id," " reason_code, reason_detail, accounting_period_start," " accounting_period_end, control_total_deltas, correction_observed_at," " created_at, managed_control_total_deltas)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s, %s, %s::jsonb)" " ON CONFLICT (reporting_adjustment_id) DO NOTHING" " RETURNING reporting_adjustment_id", ( adjustment.reporting_adjustment_id, adjustment.account_id, adjustment.adjusts_reporting_revision_id, adjustment.reason_code, adjustment.reason_detail, adjustment.accounting_period_start, adjustment.accounting_period_end, _json([[name, value] for name, value in adjustment.control_total_deltas]), adjustment.correction_observed_at, adjustment.created_at, ( _json( [item.to_wire() for item in adjustment.managed_control_total_deltas] ) if adjustment.managed_control_total_deltas is not None else None ), ), ) ).fetchone() if inserted is None: raise LedgerConflictError("ADJUSTMENT_UNAVAILABLE", "adjustment is unavailable") await self._append_change( connection, adjustment.account_id, "adjustment", adjustment.reporting_adjustment_id, consumer_id=_obligation_from_row(obligation).consumer_id, ) if self._notifications_enabled: from adcp.reporting.ledger.notification_events import adjustment_event await self._record_notification( connection, adjustment_event( adjustment, await self._notification_now(connection), consumer_id=_obligation_from_row(obligation).consumer_id, ), ) await self._dirty_status( connection, ReportingStatusScope.for_obligation(_obligation_from_row(obligation)), "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: key = obligation.generation_key async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, key.account_id) await self._require_writable_generation(connection, key) existing_row = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND period_start = %s AND period_end = %s", ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, obligation.period.start, obligation.period.end, ), ) ).fetchone() if existing_row is not None: return _obligation_from_row(existing_row) require_frozen_currency(obligation.currency) try: inserted = await ( await connection.execute( "INSERT INTO reporting_obligations" " (reporting_obligation_id, account_id, delivery_config_id," " delivery_config_version, report_definition_id, reporting_profile," " feed_purpose, period_key, period_start, period_end, source_timezone," " expected_at, scope_resolved_at, automated_recovery_deadline_at," " required_finality, coverage_status, media_buy_ids, package_ids," " schedule, definition, created_at, currency, consumer_id)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s," " %s::jsonb, %s::jsonb, %s::jsonb, %s::jsonb, %s, %s, %s)" " ON CONFLICT (account_id, consumer_id, delivery_config_id, delive" "ry_config_version," " period_start, period_end) DO NOTHING" " RETURNING reporting_obligation_id", ( obligation.reporting_obligation_id, key.account_id, key.delivery_config_id, key.delivery_config_version, obligation.report_definition_id, obligation.reporting_profile, obligation.feed_purpose, obligation.period.period_key, obligation.period.start, obligation.period.end, obligation.period.source_timezone, obligation.period.expected_at, obligation.scope_resolved_at, obligation.automated_recovery_deadline_at, obligation.required_finality, obligation.coverage_status, _json(sorted(obligation.media_buy_ids)), _json(sorted(obligation.package_ids)), _json(_schedule_payload(obligation.schedule)), ( _json(_definition_payload(obligation.definition)) if obligation.definition else None ), obligation.created_at, obligation.currency, obligation.consumer_id, ), ) ).fetchone() except Exception as error: raise _translate_integrity_error(error) from error if inserted is not None: await self._append_change( connection, obligation.account_id, "obligation", obligation.reporting_obligation_id, consumer_id=obligation.consumer_id, ) if self._notifications_enabled: await self._dirty_status( connection, ReportingStatusScope.for_obligation(obligation), "obligation", after=ReportingStatusEvidence( "obligation", obligation.reporting_obligation_id ), ) return obligation # Another worker won the period close; converge on its obligation. existing = await self.find_obligation( account_id=obligation.account_id, consumer_id=obligation.consumer_id, delivery_config_id=obligation.delivery_config_id, delivery_config_version=obligation.delivery_config_version, period_start=obligation.period.start, period_end=obligation.period.end, ) if existing is None: # pragma: no cover - only under concurrent deletion raise LedgerConflictError( "OBLIGATION_LOST", "the obligation vanished between insert and read" ) return existing 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.transaction(), self._connection() as connection: await self._lock_account(connection, acquisition.account_id) await self._require_provisional_schema(connection) existing = await ( await connection.execute( "SELECT reporting_revision_id,payload FROM reporting_provisional_observations" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", key, ) ).fetchone() if existing is not None: retained_observation = ProvisionalObservation.from_wire(existing[1]) if ( retained_observation.acquisition != acquisition or existing[0] != observation.revision_id ): raise LedgerConflictError("OBSERVATION_CONFLICT", "observation replay differs") retained_revision = await self.get_revision( account_id=acquisition.account_id, reporting_revision_id=existing[0] ) if retained_revision is None: raise LedgerConflictError( "HISTORY_UNAVAILABLE", "observation revision is missing" ) # Reuse the immutable publication replay checks while holding # the observation's account lock and transaction. return await self.commit_revision(revision, rows) reserved = await ( await connection.execute( "SELECT payload FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", key, ) ).fetchone() if reserved is None or ProvisionalAcquisition.from_wire(reserved[0]) != acquisition: raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition was not reserved") checkpoint = await self.get_restatement_checkpoint( account_id=acquisition.account_id, reporting_obligation_id=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, ) ) await connection.execute( "INSERT INTO reporting_provisional_observations" " (account_id,reporting_obligation_id,ordinal,reporting_revision_id,payload)" " VALUES (%s,%s,%s,%s,%s::jsonb)", (*key, revision.reporting_revision_id, _json(observation.to_wire())), ) 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) 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) digest = _fingerprint(_revision_payload(revision)) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, revision.account_id) existing = await ( await connection.execute( f"SELECT {_REVISION_COLUMNS}, content_sha256 FROM reporting_revisions" # noqa: S608 # nosec B608 " WHERE reporting_revision_id = %s AND account_id = %s", (revision.reporting_revision_id, revision.account_id), ) ).fetchone() if existing is not None: if existing[-1] != digest: raise LedgerConflictError( "REVISION_IMMUTABLE", f"revision {revision.reporting_revision_id} already exists with " "different content; a restatement is a new revision", ) return _revision_from_row(existing[:-1]) obligation = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE reporting_obligation_id = %s AND account_id = %s", (revision.reporting_obligation_id, revision.account_id), ) ).fetchone() if obligation is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "a revision must attach to an obligation committed at the period close", ) await self._require_writable_generation( connection, _obligation_from_row(obligation).generation_key ) validate_revision_currency(_obligation_from_row(obligation), revision, rows) if revision.supersedes_reporting_revision_id: await self._require_current_leaf(connection, revision) write_error = None try: await connection.execute( "INSERT INTO reporting_revisions" " (reporting_revision_id, account_id, reporting_obligation_id, finality," " revision_content_sha256, row_count, control_totals, observed_at," " data_through, created_at, supersedes_reporting_revision_id," " finality_basis, finality_policy_id, finalized_at, readable," " readable_at_commit, source_publication_id, source_manifest_sha256," " content_sha256, canonical_content_digest, managed_control_totals)" " VALUES (%s, %s, %s, %s, %s, %s, %s::jsonb, %s, %s, %s, %s, %s, %s, %s," " %s, %s, %s, %s, %s, %s::jsonb, %s::jsonb)", ( revision.reporting_revision_id, revision.account_id, revision.reporting_obligation_id, revision.finality, revision.revision_content_sha256, revision.row_count, _json([[name, value] for name, value in revision.control_totals]), revision.observed_at, revision.data_through, revision.created_at, revision.supersedes_reporting_revision_id, revision.finality_basis, revision.finality_policy_id, revision.finalized_at, revision.readable, revision.readable_at_commit, revision.source_publication_id, revision.source_manifest_sha256, digest, ( _json(revision.canonical_content_digest.to_wire()) if revision.canonical_content_digest is not None else None ), ( _json([item.to_wire() for item in revision.managed_control_totals]) if revision.managed_control_totals is not None else None ), ), ) except Exception as error: # psycopg raises UniqueViolation subclasses write_error = _translate_integrity_error(error) if write_error is not None: raise write_error if rows: await connection.cursor().executemany( "INSERT INTO reporting_revision_rows" " (reporting_revision_id, ordinal, row_payload)" " SELECT reporting_revision_id, %s, %s::jsonb FROM reporting_revisions" " WHERE account_id = %s AND reporting_revision_id = %s", [ (ordinal, _json(row), revision.account_id, revision.reporting_revision_id) for ordinal, row in enumerate(rows) ], ) await self._append_change( connection, revision.account_id, "revision", revision.reporting_revision_id, consumer_id=_obligation_from_row(obligation).consumer_id, ) if self._notifications_enabled: from adcp.reporting.ledger.notification_events import revision_event await self._record_notification( connection, revision_event( revision, await self._notification_now(connection), consumer_id=_obligation_from_row(obligation).consumer_id, ), ) await self._dirty_status( connection, ReportingStatusScope.for_obligation(_obligation_from_row(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: """Create or upgrade the ledger atomically, serializing concurrent boots. The bootstrap takes a transaction-scoped schema lock before any DDL. Keep the migration in that transaction, including with an autocommit pool, so a second process cannot observe a partially upgraded schema. """ async with self._connection() as connection: async with connection.transaction(): await self._create_schema_on(connection)Create or upgrade the ledger atomically, serializing concurrent boots.
The bootstrap takes a transaction-scoped schema lock before any DDL. Keep the migration in that transaction, including with an autocommit pool, so a second process cannot observe a partially upgraded schema.
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._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) return await self._ensure_issue_opened_on( connection, issue_key=issue_key, account_id=account_id, consumer_id=consumer_id, observed_at=observed_at, status_scope=status_scope, ) 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, ) async with self._connection() as connection: row = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND period_start = %s AND period_end = %s", ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, period_start, period_end, ), ) ).fetchone() return _obligation_from_row(row) if row 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._connection() as connection: return await self._live_issue(connection, issue_key, account_id) 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 with self._connection() as connection: row = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # noqa: S608 # nosec B608 " WHERE account_id = %s AND reporting_obligation_id = %s", (account_id, reporting_obligation_id), ) ).fetchone() return _obligation_from_row(row) if row 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: async with self._connection() as connection: row = await ( await connection.execute( "SELECT payload FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", (account_id, reporting_obligation_id, ordinal), ) ).fetchone() return ProvisionalAcquisition.from_wire(row[0]) if row is not None else None 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: async with self._connection() as connection: await self._require_provisional_schema(connection) row = await ( await connection.execute( "SELECT payload FROM reporting_provisional_observations" " WHERE account_id=%s AND reporting_obligation_id=%s" " ORDER BY ordinal DESC LIMIT 1", (account_id, reporting_obligation_id), ) ).fetchone() return ProvisionalObservation.from_wire(row[0]) if row else 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: async with self._connection() as connection: row = await ( await connection.execute( "SELECT account_id, reporting_obligation_id, checked_at," " next_observation, provisional_until" " FROM reporting_restatement_checkpoints" " WHERE account_id = %s AND reporting_obligation_id = %s", (account_id, reporting_obligation_id), ) ).fetchone() if row is None: return None return RestatementCheckpoint( account_id=row[0], reporting_obligation_id=row[1], checked_at=_utc(row[2]), next_observation=int(row[3]), provisional_until=_utc(row[4]) if row[4] 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: async with self._connection() as connection: row = await ( await connection.execute( "SELECT retry_not_before, attempt, blocked, recorded_at" " FROM adcp_reporting_producer_retry_schedules" " WHERE scope_key = %s", (scope_key,), ) ).fetchone() return ( RetryScheduleEntry(scope_key, _utc(row[0]), int(row[1]), row[2], _utc(row[3])) if row else 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 with self._connection() as connection: row = await ( await connection.execute( f"SELECT {_REVISION_COLUMNS} FROM reporting_revisions" # noqa: S608 # nosec B608 " WHERE account_id = %s AND reporting_revision_id = %s", (account_id, reporting_revision_id), ) ).fetchone() return _revision_from_row(row) if row 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 psycopg.errors import LockNotAvailable from psycopg.pq import TransactionStatus moment = _utc(now) expires = moment + timedelta(seconds=lease_seconds) after = self._period_close_sample wait_until = None following = None result = None async with self._connection() as connection: # Never add a blocking account edge while a caller transaction may # already own locks. Standalone turns may queue only before the first # account acquisition, including acquisitions that find a stale row. can_wait = connection.info.transaction_status == TransactionStatus.IDLE async with connection.transaction(): await self._bound_lock_waits(connection) # A materializer installed by any participant adds an AFTER UPDATE # trigger taking the account lock. Sample without row locks, then take # that account lock before locking a configuration. The base store can # share that schema, so the ordering belongs here, not only on a tier. for _ in range(2): continuation = ( " AND (COALESCE(t.lease_turn,0)," " COALESCE(c.lease_expires_at,'-infinity'::timestamptz)," " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_confi" "g_version)" " > (%s,COALESCE(%s::timestamptz,'-infinity'::timestamptz),%s,%s,%s,%s)" if after is not None else "" ) query = ( "SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.delivery" "_config_version," # nosec B608 " COALESCE(t.lease_turn,0),c.lease_expires_at FROM" " reporting_configurations c" " LEFT JOIN adcp_reporting_configuration_lease_turns t" " ON (t.account_id,t.consumer_id,t.delivery_config_id,t.delivery_c" "onfig_version)" " = (c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_co" "nfig_version)" " WHERE NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lea" "se_expires_at<=%s)" + continuation + " ORDER BY COALESCE(t.lease_turn,0),c.lease_expires_at NULLS FIRST," " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_confi" "g_version LIMIT 32" ) rows = await ( await connection.execute(query, (moment, *(after or ()))) ).fetchall() if rows or after is None: break after = None for candidate in rows: row = candidate[:4] following = (candidate[4], candidate[5], *row) locked = await ( await connection.execute( "SELECT pg_try_advisory_xact_lock(hashtext('adcp.reporting:' || %s))", (row[0],), ) ).fetchone() if not locked[0]: if not can_wait: continue # A try-lock alone can starve behind a continuous queue of # ordinary writers, even though each writer commits promptly. # Join that queue briefly, sharing one wait budget per turn. # The savepoint rolls back a timed-out wait and its SET LOCAL; # successful acquisition restores the caller's lock timeout. if wait_until is None: wait_until = monotonic() + _LEASE_ACCOUNT_WAIT_SECONDS remaining_ms = int((wait_until - monotonic()) * 1000) if remaining_ms <= 0: continue try: async with connection.transaction(): previous, configured_ms = await ( await connection.execute( "SELECT current_setting('lock_timeout'),setting::integer" " FROM pg_settings WHERE name='lock_timeout'" ) ).fetchone() limit_ms = min(remaining_ms, configured_ms or remaining_ms) await connection.execute( "SELECT set_config('lock_timeout',%s,true)", (f"{limit_ms}ms",) ) await connection.execute( "SELECT pg_advisory_xact_lock(hashtext('adcp.reporting:' ||" " %s))", (row[0],), ) await connection.execute( "SELECT set_config('lock_timeout',%s,true)", (previous,) ) except LockNotAvailable: continue can_wait = False snapshot = await self._period_close_generations_on(connection, tuple(row)) acquired = await ( await connection.execute( "UPDATE reporting_configurations SET lease_worker_id = %s," " lease_expires_at = %s" " WHERE (account_id,consumer_id,delivery_config_id,delivery_config" "_version) = (" " SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.deliver" "y_config_version" " FROM reporting_configurations c WHERE c.account_id=%s" " AND c.consumer_id=%s AND c.delivery_config_id=%s AND c.delivery_" "config_version=%s" " AND NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lease" "_expires_at<=%s)" " FOR UPDATE OF c SKIP LOCKED) RETURNING account_id", (worker_id, expires, *row, moment), ) ).fetchone() if acquired is None: continue await self._restore_period_close_generations_on( connection, tuple(row), snapshot ) await connection.execute( "INSERT INTO adcp_reporting_configuration_lease_turns" " (account_id,consumer_id,delivery_config_id,delivery_config_versi" "on,lease_turn)" " VALUES(%s,%s,%s,%s,nextval('adcp_reporting_configuration_lease_t" "urn_seq'))" " ON CONFLICT(account_id,consumer_id,delivery_config_id,delivery_c" "onfig_version)" " DO UPDATE SET" " lease_turn=nextval('adcp_reporting_configuration_lease_turn_seq')", tuple(row), ) result = LeasedConfiguration(row[0], row[1], row[2], row[3], expires) break # Only sampling uses this hint: leases and fairness ranks stay transactional. # Continue past a busy prefix on the next turn; a successful turn returns to # durable fairness order. Empty tails wrap at most once without row locks. self._period_close_sample = following if result is None else None return result 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, ...]: if not reporting_revision_ids: return () async with self._connection() as connection: rows = await ( await connection.execute( f"SELECT {_ADJUSTMENT_COLUMNS} FROM reporting_adjustments" # noqa: S608 # nosec B608 " WHERE account_id = %s AND adjusts_reporting_revision_id = ANY(%s)" " ORDER BY created_at, reporting_adjustment_id", (account_id, list(reporting_revision_ids)), ) ).fetchall() return tuple(_adjustment_from_row(row) for row in rows) async def list_all_configurations(self) ‑> tuple[ReportingConfiguration, ...]-
Expand source code
async def list_all_configurations(self) -> tuple[ReportingConfiguration, ...]: """One ordered scan for trusted service recovery, including retired generations.""" async with self._connection() as connection: rows = await ( await connection.execute( "SELECT delivery_config_id, delivery_config_version, account_id," " report_definition_id, reporting_profile, feed_purpose, required_finality," " account_timezone, schedule, media_buy_ids, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, definition," " authoritative_party, consumer_id, quarantined" " FROM reporting_configurations WHERE NOT quarantined" " ORDER BY account_id, consumer_id, delivery_config_id, delivery_config_version" ) ).fetchall() return tuple(_configuration_from_row(row) for row in rows)One ordered scan for trusted service recovery, including retired generations.
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 with self._connection() as connection: return await self._list_configurations_on( connection, account_id=caller.account_id, consumer_id=caller.consumer_id, delivery_config_ids=delivery_config_ids, ) 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, ...]: # Attaches by obligation id *or* by the exact logical period key, so a # chain filed before the obligation existed is not lost, forked, or # reset when the seller later repairs the missing obligation. clause = "" params: list[Any] = [account_id, consumer_id] if reporting_obligation_ids is not None: clause = ( " AND EXISTS (" " SELECT 1 FROM reporting_obligations o" " WHERE o.reporting_obligation_id = ANY(%s)" " AND o.account_id = s.account_id" " AND (o.reporting_obligation_id = s.reporting_obligation_id OR (" " o.delivery_config_id = s.delivery_config_id" " AND o.delivery_config_version = s.delivery_config_version" " AND o.report_definition_id = s.report_definition_id" " AND o.period_start = s.period_start" " AND o.period_end = s.period_end)))" ) params.append(list(reporting_obligation_ids)) async with self._connection() as connection: rows = await ( await connection.execute( f"SELECT {_STATUS_COLUMNS} FROM reporting_consumer_statuses s" # noqa: S608 # nosec B608 f" WHERE s.account_id = %s AND s.consumer_id = %s{clause}" " ORDER BY s.recorded_at, s.reporting_status_id", tuple(params), ) ).fetchall() return tuple(_status_from_row(row) for row in rows) 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 with self._connection() as connection: rows = await ( await connection.execute( f"SELECT {_REVISION_COLUMNS} FROM reporting_revisions" # noqa: S608 # nosec B608 " WHERE account_id = %s AND reporting_obligation_id = %s" " ORDER BY created_at, reporting_revision_id", (account_id, reporting_obligation_id), ) ).fetchall() return tuple(_revision_from_row(row) for row in rows) 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._connection() as connection: row = await ( await connection.execute( "SELECT COALESCE(MAX(seq), 0), now() FROM reporting_ledger_changes" " WHERE account_id=%s AND consumer_id=%s AND record_kind IN" " ('obligation','revision','adjustment','consumer_status')", (account_id, caller.consumer_id), ) ).fetchone() assert row is not None max_sequence = int(row[0]) as_of = _utc(self._clock()) if self._clock is not None else _utc(row[1]) 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) key = configuration.generation_key payload = _configuration_payload(configuration) digest = _fingerprint(payload) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, key.account_id) inserted_cursor = await connection.execute( "INSERT INTO reporting_configurations" " (delivery_config_id, delivery_config_version, account_id," " report_definition_id, reporting_profile, feed_purpose, required_finality," " account_timezone, schedule, media_buy_ids, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, definition," " authoritative_party, content_sha256, consumer_id, quarantined)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s::jsonb, %s::jsonb, %s, %s, %s, %s," " %s::jsonb, %s, %s, %s, %s)" " ON CONFLICT (account_id, consumer_id, delivery_config_id, delive" "ry_config_version) DO NOTHING" " RETURNING delivery_config_id", ( key.delivery_config_id, key.delivery_config_version, key.account_id, configuration.report_definition_id, configuration.reporting_profile, configuration.feed_purpose, configuration.required_finality, configuration.account_timezone, _json(payload["schedule"]), _json(sorted(configuration.media_buy_ids)), configuration.activated_at, configuration.deactivated_at, configuration.automated_recovery_window.total_seconds(), configuration.status_retention_days, _json(payload["definition"]) if payload["definition"] else None, configuration.authoritative_party, digest, configuration.consumer_id, configuration.quarantined, ), ) inserted = await inserted_cursor.fetchone() # Check *after* the insert. ON CONFLICT waits for a concurrent # winner; a fresh READ COMMITTED statement sees its retained # content. A pre-insert check followed by DO NOTHING could silently # accept a different immutable generation from a losing writer. row = await ( await connection.execute( "SELECT content_sha256, activated_at, deactivated_at," " automated_recovery_seconds, status_retention_days, quarantined" " FROM reporting_configurations" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s FOR UPDATE", ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ), ) ).fetchone() if row is not None and row[5]: raise LedgerConflictError( "REPORTING_GENERATION_QUARANTINED", "legacy evidence is read-only" ) if row is None or row[0] != digest: 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", ) scope = ReportingStatusScope(configuration.account_id, configuration.generation_key) if inserted is not None: if self._notifications_enabled: await self._dirty_status( connection, scope, "configuration", after=configuration_evidence(configuration), ) return # rc.3 carries activation/deactivation and the recovery/retention # windows as lifecycle state over one immutable generation, so a # re-put that changes only those must apply -- otherwise a # deactivated feed keeps minting obligations -- and must co-commit # its status-dirty generation on this exact connection. An # unchanged re-put stays a no-op and enqueues nothing. retained = replace( configuration, activated_at=_utc(row[1]) if row[1] else None, deactivated_at=_utc(row[2]) if row[2] else None, automated_recovery_window=timedelta(seconds=float(row[3])), status_retention_days=row[4], ) if configuration_lifecycle(retained) == configuration_lifecycle(configuration): return await connection.execute( "UPDATE reporting_configurations SET activated_at = %s, deactivated_at = %s," " automated_recovery_seconds = %s, status_retention_days = %s" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s", ( configuration.activated_at, configuration.deactivated_at, configuration.automated_recovery_window.total_seconds(), configuration.status_retention_days, key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ), ) if self._notifications_enabled: await self._dirty_status( connection, scope, "configuration", before=configuration_evidence(retained), 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 async with self._connection() as connection: rows = await ( await connection.execute( "SELECT seq, record_kind, record_id FROM reporting_ledger_changes" " WHERE account_id = %s AND consumer_id=%s AND seq > %s AND seq <= %s" " ORDER BY seq", (snapshot.account_id, consumer_id, lower, snapshot.max_sequence), ) ).fetchall() selected: list[tuple[int, LedgerRecordKind, Any]] = [] for _seq, kind, record_id in rows: record = await self._resolve( connection, snapshot.account_id, kind, record_id, consumer_id=consumer_id ) if record is None: continue if not await self._in_scope( connection, kind, record, delivery_config_ids, media_buy_ids, consumer_id, feed_purposes, period_start, period_end, ): continue selected.append((_seq, kind, record)) window = selected[offset : offset + limit] has_more = offset + limit < len(selected) 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(selected), 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: from adcp.reporting.ledger.store import revision_row_offset offset = revision_row_offset(cursor, reporting_revision_id, limit) async with self._connection() as connection: owned = await ( await connection.execute( "SELECT row_count FROM reporting_revisions" " WHERE account_id = %s AND reporting_revision_id = %s", (account_id, reporting_revision_id), ) ).fetchone() if owned is None: raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account") rows = await ( await connection.execute( "SELECT row_payload FROM reporting_revision_rows" " WHERE reporting_revision_id = %s" " ORDER BY ordinal OFFSET %s LIMIT %s", (reporting_revision_id, offset, limit), ) ).fetchall() total = int(owned[0]) has_more = offset + limit < total return ReportingRowPage( reporting_revision_id=reporting_revision_id, rows=tuple(row[0] for row in rows), total_count=total, 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.delivery_models import ReportingDeliveryPrincipal from adcp.reporting.ledger.status_snapshot import settle_snapshot_on from adcp.reporting.materializer.capture import private_snapshot async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, caller.account_id) return private_snapshot( await settle_snapshot_on(self, connection, account_id=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) key = status.generation_key digest = _fingerprint(_consumer_status_payload(status)) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, status.account_id) replay = await self._replay(connection, status, digest) if replay is not None: return replay, False from adcp.reporting.ledger.status_snapshot import ( read_snapshot_on, validate_status_evidence, ) validate_status_evidence( status, await read_snapshot_on( connection, account_id=status.account_id, clock=self._clock, include_issue_scopes=await self._issue_scope_storage_on(connection), ), ) leaf = await ( await connection.execute( "SELECT reporting_status_id FROM reporting_consumer_statuses" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND report_definition_id = %s" " AND period_start = %s AND period_end = %s AND superseded = FALSE", ( key.account_id, status.consumer_id, key.delivery_config_id, key.delivery_config_version, status.report_definition_id, status.period_start, status.period_end, ), ) ).fetchone() if status.supersedes_reporting_status_id: if leaf is None or leaf[0] != 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", ) await connection.execute( "UPDATE reporting_consumer_statuses SET superseded = TRUE" " WHERE account_id = %s AND consumer_id = %s AND reporting_status_id = %s", (key.account_id, status.consumer_id, status.supersedes_reporting_status_id), ) elif leaf is not None: raise LedgerConflictError( "STATUS_SUPERSEDES_REQUIRED", "this chain already has a current statement; a new statement must " "explicitly supersede it", ) try: await connection.execute( "INSERT INTO reporting_consumer_statuses" " (reporting_status_id, account_id, consumer_id, delivery_config_id," " delivery_config_version, report_definition_id, period_start, period_end," " period_source_timezone, consumer_status, status_as_of, recorded_at," " supersedes_reporting_status_id, reporting_obligation_id," " reporting_revision_id, observed_revision_content_sha256, failure_code," " mismatch_code, consumer_commit_ref, seller_ledger_snapshot_id," " seller_ledger_as_of, superseded, content_sha256)" " VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s," " %s, %s, %s, %s, %s, FALSE, %s)", ( status.reporting_status_id, status.account_id, status.consumer_id, status.delivery_config_id, status.delivery_config_version, status.report_definition_id, status.period_start, status.period_end, status.period_source_timezone, status.consumer_status, status.status_as_of, status.recorded_at, status.supersedes_reporting_status_id, status.reporting_obligation_id, status.reporting_revision_id, status.observed_revision_content_sha256, status.failure_code, status.mismatch_code, status.consumer_commit_ref, status.seller_ledger_snapshot_id, status.seller_ledger_as_of, digest, ), ) except Exception as error: raise _translate_integrity_error(error) from error await self._append_change( connection, status.account_id, "consumer_status", status.reporting_status_id, consumer_id=status.consumer_id, ) from adcp.reporting.ledger.status_snapshot import settle_snapshot_on await settle_snapshot_on(self, connection, account_id=status.account_id) if self._notifications_enabled: await self._dirty_status( connection, ReportingStatusScope( status.account_id, status.generation_key, status.reporting_obligation_id, status.consumer_id, ), "consumer_status", ReportingStatusEvidence("consumer_status", leaf[0]) if leaf 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]: """Optional participant: the record and lifecycle share the source transaction.""" return await self.record_consumer_status(status)Optional participant: the record and lifecycle share the source transaction.
async def record_restatement_checkpoint(self, checkpoint: RestatementCheckpoint) ‑> RestatementCheckpoint-
Expand source code
async def record_restatement_checkpoint( self, checkpoint: RestatementCheckpoint ) -> RestatementCheckpoint: async with self._connection() as connection: obligation = await ( await connection.execute( "SELECT 1 FROM reporting_obligations" " WHERE reporting_obligation_id = %s AND account_id = %s", (checkpoint.reporting_obligation_id, checkpoint.account_id), ) ).fetchone() if obligation is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "a restatement checkpoint must attach to an obligation for this account", ) row = await ( await connection.execute( "INSERT INTO reporting_restatement_checkpoints" " (account_id, reporting_obligation_id, checked_at, next_observation," " provisional_until)" " VALUES (%s, %s, %s, %s, %s)" " ON CONFLICT (reporting_obligation_id) DO UPDATE SET" " checked_at = EXCLUDED.checked_at," " next_observation = EXCLUDED.next_observation," " provisional_until = EXCLUDED.provisional_until" " WHERE reporting_restatement_checkpoints.account_id = EXCLUDED.account_id" " AND (reporting_restatement_checkpoints.next_observation" " < EXCLUDED.next_observation" " OR (reporting_restatement_checkpoints.next_observation" " = EXCLUDED.next_observation" " AND reporting_restatement_checkpoints.checked_at" " < EXCLUDED.checked_at))" " RETURNING account_id, reporting_obligation_id, checked_at," " next_observation, provisional_until", ( checkpoint.account_id, checkpoint.reporting_obligation_id, checkpoint.checked_at, checkpoint.next_observation, checkpoint.provisional_until, ), ) ).fetchone() if row is None: stored = await self.get_restatement_checkpoint( account_id=checkpoint.account_id, reporting_obligation_id=checkpoint.reporting_obligation_id, ) if stored is None: raise LedgerConflictError( "OBLIGATION_NOT_FOUND", "a restatement checkpoint must attach to an obligation for this account", ) return stored return RestatementCheckpoint( account_id=row[0], reporting_obligation_id=row[1], checked_at=_utc(row[2]), next_observation=int(row[3]), provisional_until=_utc(row[4]) if row[4] else None, ) async def record_retry_schedule(self, entry: RetryScheduleEntry) ‑> RetryScheduleEntry-
Expand source code
async def record_retry_schedule(self, entry: RetryScheduleEntry) -> RetryScheduleEntry: async with self._connection() as connection: row = await ( await connection.execute( "INSERT INTO adcp_reporting_producer_retry_schedules" " (scope_key, retry_not_before, attempt, blocked, recorded_at)" " VALUES (%s, %s, %s, %s, %s)" " ON CONFLICT (scope_key) DO UPDATE SET" " retry_not_before = CASE WHEN EXCLUDED.attempt = 0 OR EXCLUDED.blocked" " OR adcp_reporting_producer_retry_schedules.blocked OR %s" " THEN EXCLUDED.retry_not_before ELSE GREATEST(" " adcp_reporting_producer_retry_schedules.retry_not_before," " EXCLUDED.retry_not_before) END," " attempt = EXCLUDED.attempt, blocked = EXCLUDED.blocked," " recorded_at = EXCLUDED.recorded_at" " WHERE EXCLUDED.recorded_at >=" " adcp_reporting_producer_retry_schedules.recorded_at" " AND (EXCLUDED.attempt = 0 OR EXCLUDED.blocked OR %s" " OR (NOT adcp_reporting_producer_retry_schedules.blocked AND" " adcp_reporting_producer_retry_schedules.attempt <= EXCLUDED.attempt))" " RETURNING retry_not_before, attempt, blocked, recorded_at", ( entry.scope_key, entry.retry_not_before, entry.attempt, entry.blocked, entry.recorded_at, entry.replayed, entry.replayed, ), ) ).fetchone() if row is None: stored = await self.get_retry_schedule(scope_key=entry.scope_key) assert stored is not None return stored return RetryScheduleEntry(entry.scope_key, _utc(row[0]), int(row[1]), row[2], _utc(row[3])) 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: key = lease.generation_key identity = ( key.account_id, key.consumer_id, key.delivery_config_id, key.delivery_config_version, ) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, key.account_id) snapshot = await self._period_close_generations_on(connection, identity) released = await ( await connection.execute( "UPDATE reporting_configurations SET lease_worker_id = NULL," " lease_expires_at = NULL" " WHERE account_id = %s AND consumer_id = %s AND delivery_config_id = %s" " AND delivery_config_version = %s AND lease_worker_id = %s" " AND lease_expires_at = %s RETURNING account_id", (*identity, worker_id, lease.lease_expires_at), ) ).fetchone() if released is not None: await self._restore_period_close_generations_on(connection, identity, snapshot) async def reserve_provisional_acquisition(self, acquisition: ProvisionalAcquisition) ‑> ProvisionalAcquisition-
Expand source code
async def reserve_provisional_acquisition( self, acquisition: ProvisionalAcquisition ) -> ProvisionalAcquisition: async with self.transaction(), self._connection() as connection: await self._lock_account(connection, acquisition.account_id) await self._require_provisional_schema(connection) key = (acquisition.account_id, acquisition.obligation_id, acquisition.ordinal) existing = await ( await connection.execute( "SELECT payload FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND reporting_obligation_id=%s AND ordinal=%s", key, ) ).fetchone() if existing is not None: return ProvisionalAcquisition.from_wire(existing[0]) obligation = await self.get_obligation( account_id=acquisition.account_id, reporting_obligation_id=acquisition.obligation_id, ) if obligation is None: raise LedgerConflictError("OBLIGATION_NOT_FOUND", "unknown observation obligation") if not acquisition.binds(obligation): raise LedgerConflictError("OBSERVATION_CONFLICT", "acquisition generation differs") duplicate = await ( await connection.execute( "SELECT 1 FROM reporting_provisional_acquisitions" " WHERE account_id=%s AND source_execution_key=%s", (acquisition.account_id, acquisition.execution_key), ) ).fetchone() if duplicate is not None: raise LedgerConflictError( "OBSERVATION_CONFLICT", "execution key is already reserved" ) checkpoint = await self.get_restatement_checkpoint( account_id=acquisition.account_id, reporting_obligation_id=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") await connection.execute( "INSERT INTO reporting_provisional_acquisitions" " (account_id,reporting_obligation_id,ordinal,source_execution_key,payload)" " VALUES (%s,%s,%s,%s,%s::jsonb)", (*key, acquisition.execution_key, _json(acquisition.to_wire())), ) return ProvisionalAcquisition.from_wire(acquisition.to_wire()) 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._connection() as connection: return await self._replay( connection, status, _fingerprint(_consumer_status_payload(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._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) return await self._retire_issue_on( connection, issue_key=issue_key, account_id=account_id, at=at, status_scope=status_scope, ) 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: 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", ) async with self._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) live = await self._live_issue(connection, issue_key, account_id) if live is None: 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 read_snapshot_on snapshot = await read_snapshot_on( connection, account_id=account_id, as_of=_utc(at), include_issue_scopes=await self._issue_scope_storage_on(connection), ) live = bind_mismatch_waiver(snapshot, live) await _save_waiver_binding(connection, live) waived_at = ( live.retired_at if live.waived_reporting_status_id is not None and live.retired_at is not None else _utc(at) ) await connection.execute( "UPDATE reporting_issue_lifecycle" " SET issue_state = %s," " external_ref = COALESCE(%s, external_ref)," " retired_at = CASE WHEN %s = 'waived' THEN %s ELSE retired_at END" " WHERE account_id = %s AND issue_key = %s AND generation = %s", ( state, external_ref, state, waived_at, account_id, issue_key, live.generation, ), ) refreshed = await self._issue_row(connection, issue_key, account_id, live.generation) assert refreshed is not None # Derive the no-op from the resulting row 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. await self._dirty_issue(connection, refreshed, status_scope, live) return refreshed 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._connection() as connection, connection.transaction(): await self._lock_account(connection, account_id) existing = await ( await connection.execute( "SELECT readable, reporting_obligation_id FROM reporting_revisions" " WHERE account_id = %s AND reporting_revision_id = %s FOR UPDATE", (account_id, reporting_revision_id), ) ).fetchone() if existing is None: raise LedgerConflictError("REVISION_NOT_FOUND", "no such revision for this account") if existing[0] == readable: return await connection.execute( "UPDATE reporting_revisions SET readable = %s" " WHERE account_id = %s AND reporting_revision_id = %s", (readable, account_id, reporting_revision_id), ) if self._notifications_enabled: obligation = await ( await connection.execute( f"SELECT {_OBLIGATION_COLUMNS} FROM reporting_obligations" # nosec B608 " WHERE account_id = %s AND reporting_obligation_id = %s", (account_id, existing[1]), ) ).fetchone() assert obligation is not None await self._dirty_status( connection, ReportingStatusScope.for_obligation(_obligation_from_row(obligation)), "readability", ReportingStatusEvidence( "revision", reporting_revision_id, readable=existing[0] ), ReportingStatusEvidence("revision", reporting_revision_id, readable=readable), ) async def transaction(self) ‑> AsyncIterator[PgReportingLedgerStore]-
Expand source code
@asynccontextmanager async def transaction(self) -> AsyncIterator[PgReportingLedgerStore]: """Group source operations into one committed status boundary. Pool ownership remains with the adopter. Nested turns are savepoints; callers acquire accounts in canonical order when touching several. """ async with self._connection() as connection, connection.transaction(): token = _BOUND_CONNECTION.set((self._pool, asyncio.current_task(), connection)) try: yield self finally: _BOUND_CONNECTION.reset(token)Group source operations into one committed status boundary.
Pool ownership remains with the adopter. Nested turns are savepoints; callers acquire accounts in canonical order when touching several.