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 NoneOne 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
- PgReportingReceiptStore
- PgReportingMaterializerStore
- PgReportingReconciliationStore
- PgReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
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