Module adcp.reporting.feed.pg

Connection-bound frozen reads over the unchanged receipt/materializer store.

Classes

class PgReportingFeedStore (*,
pool: AsyncConnectionPool,
clock: Callable[[], datetime] | None = None,
notifications: bool = False)
Expand source code
class PgReportingFeedStore(PgReportingReceiptStore):
    @storage_errors
    async def create_schema(self) -> None:
        async with self._connection() as connection, connection.transaction():
            await self._create_schema_on(connection)
            root = files("adcp.reporting.ledger")
            for name in (
                "reporting_materializer.sql",
                "reporting_receipt_ingestion.sql",
                "reporting_feed.sql",
            ):
                await connection.execute(root.joinpath(name).read_text())

    @storage_errors
    async def reporting_feed_ready(self) -> bool:
        async with self._connection() as connection:
            await validate_feed_schema(connection, notifications=self._notifications_enabled)
        return True

    async def _feed_snapshot_on(
        self, connection: Any, snapshot_id: str, caller: ReportingDeliveryPrincipal
    ) -> StoredFeedSnapshot | None:
        row = await (
            await connection.execute(
                "SELECT document, content_sha256, signing_key FROM reporting_feed_snapshots"
                " WHERE account_id=%s AND consumer_id=%s AND snapshot_id=%s",
                (caller.account_id, caller.consumer_id, snapshot_id),
            )
        ).fetchone()
        if row is None:
            return None
        if hashlib.sha256(row[0].encode()).hexdigest() != row[1] or len(row[2]) != 32:
            raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT")
        snapshot = decode_snapshot(json.loads(row[0]))
        if snapshot.caller != caller or snapshot.snapshot_id != snapshot_id:
            raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT")
        return StoredFeedSnapshot(snapshot, bytes(row[2]))

    async def _save_feed_snapshot_on(self, connection: Any, stored: _PreparedFeed) -> None:
        snapshot = stored.snapshot
        await connection.execute(
            "INSERT INTO reporting_feed_snapshots"
            " (account_id,consumer_id,snapshot_id,as_of,representation_version,ownership_mode,"
            " document,content_sha256,signing_key) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s)",
            (
                snapshot.caller.account_id,
                snapshot.caller.consumer_id,
                snapshot.snapshot_id,
                snapshot.as_of,
                snapshot.representation_version,
                snapshot.ownership_mode,
                stored.document,
                stored.content_sha256,
                stored.signing_key,
            ),
        )

    async def _capture_feed_on(
        self,
        connection: Any,
        caller: ReportingDeliveryPrincipal,
    ) -> _CapturedFeed:
        # Database time and both histories share the writer's account lock. No
        # public reader, second connection, lifecycle settlement, or queue write.
        as_of = await _now(connection)
        core = await read_snapshot_on(connection, account_id=caller.account_id, as_of=as_of)
        await self._validate_feed(connection, caller)
        changes = await self._changes(connection, caller)
        materializer = await (
            await connection.execute(
                "SELECT input,content_sha256=reporting_payload_sha256(input)"
                " FROM reporting_materializer_status_boundaries"
                " WHERE account_id=%s AND consumer_id=%s ORDER BY sequence",
                (caller.account_id, caller.consumer_id),
            )
        ).fetchall()
        receipts = await (
            await connection.execute(
                "SELECT input,content_sha256=reporting_receipt_ingestion_sha256(input)"
                " FROM reporting_receipt_ingestion_boundaries"
                " WHERE account_id=%s AND consumer_id=%s ORDER BY sequence",
                (caller.account_id, caller.consumer_id),
            )
        ).fetchall()
        if any(not r[1] for r in (*materializer, *receipts)):
            raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT")
        return _CapturedFeed(
            core,
            changes,
            tuple(r[0] for r in materializer),
            tuple(r[0] for r in receipts),
        )

    @storage_errors
    async def read_reporting_feed(
        self,
        request: dict[str, Any],
        *,
        caller: ReportingDeliveryPrincipal,
        consumer_status_enabled: bool = False,
        reauthorize: Callable[[], Awaitable[None]] | None = None,
    ) -> dict[str, Any]:
        parsed = FeedRequest.parse(request)
        bound = _BOUND_CONNECTION.get()
        if bound is not None and bound[:2] == (self._pool, asyncio.current_task()):
            # A savepoint cannot release the outer transaction's account lock
            # or establish that its uncommitted history is a public boundary.
            raise ReportingFeedError("REPORTING_FEED_TRANSACTION_UNAVAILABLE")
        async with self._connection() as connection, connection.transaction():
            await validate_feed_schema(connection, notifications=self._notifications_enabled)
            stored = None
            offset = 0
            # Restore before inspecting current projection/evidence. Authorization
            # may deny a caller; no missing/old token silently opens a new walk.
            if parsed.cursor is not None:
                stored = await self._feed_snapshot_on(
                    connection, token_position(parsed.cursor)[2], caller
                )
                if stored is None:
                    raise ReportingFeedError("INVALID_CHECKPOINT")
                offset = stored.check(parsed.cursor, "cursor", parsed, caller)
            after = (0, 0)
            if parsed.changes_after is not None:
                previous = await self._feed_snapshot_on(
                    connection, token_position(parsed.changes_after)[2], caller
                )
                if previous is None:
                    raise ReportingFeedError("INVALID_CHECKPOINT")
                previous.check(parsed.changes_after, "checkpoint", parsed, caller)
                after = previous.snapshot.through
                if stored is not None and stored.snapshot.after != after:
                    raise ReportingFeedError("INVALID_CHECKPOINT")
            if stored is None:
                await self._lock_account(connection, caller.account_id)
                captured = await self._capture_feed_on(connection, caller)
        # Both the transaction and pooled connection have been released. Slow
        # graph closure, canonicalization and token/page construction must not
        # serialize receipt/materializer writers, including a size-one pool.
        if stored is not None:
            page = await asyncio.to_thread(stored.page, offset, parsed.limit)
            if reauthorize is not None:
                await reauthorize()
            return inject_context(request, page)
        prepared = await asyncio.to_thread(
            _prepare_feed, captured, caller, parsed, after, consumer_status_enabled
        )
        page = inject_context(request, prepared.page)
        if reauthorize is not None:
            # The application ACL can use the same size-one pool: no connection
            # is borrowed while the callback rechecks this captured principal.
            await reauthorize()
        async with self._connection() as connection, connection.transaction():
            await validate_feed_schema(connection, notifications=self._notifications_enabled)
            # This independent immutable row needs no account writer lock.
            # Insert/commit failure publishes no page and cannot undo writers
            # that committed after the historical capture boundary.
            await self._save_feed_snapshot_on(connection, prepared)
        return page

    @storage_errors
    async def read_reporting_feed_snapshot(
        self, snapshot_id: str, *, caller: ReportingDeliveryPrincipal
    ) -> ReportingFeedSnapshot | None:
        async with self._connection() as connection, connection.transaction():
            await validate_feed_schema(connection, notifications=self._notifications_enabled)
            stored = await self._feed_snapshot_on(connection, snapshot_id, caller)
            return stored.snapshot if stored else None

One composition: old public stores plus optional ingestion and capture.

No network I/O, session lock, or second pool participates in an ordinal. All actual receipt timestamps come from PostgreSQL, including when an old public conformance clock was supplied to the inherited constructor.

Ancestors

Subclasses

Methods

async def read_reporting_feed(self,
request: dict[str, Any],
*,
caller: ReportingDeliveryPrincipal,
consumer_status_enabled: bool = False,
reauthorize: Callable[[], Awaitable[None]] | None = None) ‑> dict[str, typing.Any]
Expand source code
@storage_errors
async def read_reporting_feed(
    self,
    request: dict[str, Any],
    *,
    caller: ReportingDeliveryPrincipal,
    consumer_status_enabled: bool = False,
    reauthorize: Callable[[], Awaitable[None]] | None = None,
) -> dict[str, Any]:
    parsed = FeedRequest.parse(request)
    bound = _BOUND_CONNECTION.get()
    if bound is not None and bound[:2] == (self._pool, asyncio.current_task()):
        # A savepoint cannot release the outer transaction's account lock
        # or establish that its uncommitted history is a public boundary.
        raise ReportingFeedError("REPORTING_FEED_TRANSACTION_UNAVAILABLE")
    async with self._connection() as connection, connection.transaction():
        await validate_feed_schema(connection, notifications=self._notifications_enabled)
        stored = None
        offset = 0
        # Restore before inspecting current projection/evidence. Authorization
        # may deny a caller; no missing/old token silently opens a new walk.
        if parsed.cursor is not None:
            stored = await self._feed_snapshot_on(
                connection, token_position(parsed.cursor)[2], caller
            )
            if stored is None:
                raise ReportingFeedError("INVALID_CHECKPOINT")
            offset = stored.check(parsed.cursor, "cursor", parsed, caller)
        after = (0, 0)
        if parsed.changes_after is not None:
            previous = await self._feed_snapshot_on(
                connection, token_position(parsed.changes_after)[2], caller
            )
            if previous is None:
                raise ReportingFeedError("INVALID_CHECKPOINT")
            previous.check(parsed.changes_after, "checkpoint", parsed, caller)
            after = previous.snapshot.through
            if stored is not None and stored.snapshot.after != after:
                raise ReportingFeedError("INVALID_CHECKPOINT")
        if stored is None:
            await self._lock_account(connection, caller.account_id)
            captured = await self._capture_feed_on(connection, caller)
    # Both the transaction and pooled connection have been released. Slow
    # graph closure, canonicalization and token/page construction must not
    # serialize receipt/materializer writers, including a size-one pool.
    if stored is not None:
        page = await asyncio.to_thread(stored.page, offset, parsed.limit)
        if reauthorize is not None:
            await reauthorize()
        return inject_context(request, page)
    prepared = await asyncio.to_thread(
        _prepare_feed, captured, caller, parsed, after, consumer_status_enabled
    )
    page = inject_context(request, prepared.page)
    if reauthorize is not None:
        # The application ACL can use the same size-one pool: no connection
        # is borrowed while the callback rechecks this captured principal.
        await reauthorize()
    async with self._connection() as connection, connection.transaction():
        await validate_feed_schema(connection, notifications=self._notifications_enabled)
        # This independent immutable row needs no account writer lock.
        # Insert/commit failure publishes no page and cannot undo writers
        # that committed after the historical capture boundary.
        await self._save_feed_snapshot_on(connection, prepared)
    return page
async def read_reporting_feed_snapshot(self, snapshot_id: str, *, caller: ReportingDeliveryPrincipal) ‑> ReportingFeedSnapshot | None
Expand source code
@storage_errors
async def read_reporting_feed_snapshot(
    self, snapshot_id: str, *, caller: ReportingDeliveryPrincipal
) -> ReportingFeedSnapshot | None:
    async with self._connection() as connection, connection.transaction():
        await validate_feed_schema(connection, notifications=self._notifications_enabled)
        stored = await self._feed_snapshot_on(connection, snapshot_id, caller)
        return stored.snapshot if stored else None
async def reporting_feed_ready(self) ‑> bool
Expand source code
@storage_errors
async def reporting_feed_ready(self) -> bool:
    async with self._connection() as connection:
        await validate_feed_schema(connection, notifications=self._notifications_enabled)
    return True

Inherited members