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.