Module adcp.reporting.feed
Frozen authorized reporting feeds; PostgreSQL remains an optional lazy import.
Sub-modules
adcp.reporting.feed.errors-
Closed errors for authorized frozen reporting reads.
adcp.reporting.feed.memory-
Reference feed participant sharing the approved memory rollback boundary.
adcp.reporting.feed.pg-
Connection-bound frozen reads over the unchanged receipt/materializer store.
adcp.reporting.feed.projection-
Capture version 1 from committed inputs, with exact graph closure and no I/O.
adcp.reporting.feed.request-
Normalize the complete semantic request, independently of transport models.
adcp.reporting.feed.schema-
Independent feed manifest; old ingestion and materializer guards stay intact.
adcp.reporting.feed.snapshot-
Immutable wire membership and versioned private inputs for one complete walk …
adcp.reporting.feed.store-
Optional feed participant. All legacy required protocols remain unchanged.
Classes
class InMemoryReportingFeedStore (**kwargs: Any)-
Expand source code
class InMemoryReportingFeedStore(InMemoryReportingReceiptStore): """Persisted for this store's lifetime; production uses PgReportingFeedStore.""" _reporting_feed_snapshots: dict[str, tuple[bytes, str, bytes]] def _feed_projection_options(self, caller: ReportingDeliveryPrincipal) -> dict[str, Any]: return {} def _feed_snapshot( self, snapshot_id: str, caller: ReportingDeliveryPrincipal ) -> StoredFeedSnapshot | None: value = getattr(self, "_reporting_feed_snapshots", {}).get(snapshot_id) if value is None: return None document, digest, key = value if hashlib.sha256(document).hexdigest() != digest or len(key) != 32: raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT") snapshot = decode_snapshot(json.loads(document)) if snapshot.snapshot_id != snapshot_id: raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT") if snapshot.caller != caller: return None return StoredFeedSnapshot(snapshot, key) def _save_feed_snapshot(self, stored: StoredFeedSnapshot) -> None: document = canonical_json_utf8_v1(stored.snapshot.to_storage()) decode_snapshot(json.loads(document)) if not hasattr(self, "_reporting_feed_snapshots"): self._reporting_feed_snapshots = {} if stored.snapshot.snapshot_id in self._reporting_feed_snapshots: raise ReportingFeedError("REPORTING_FEED_STORAGE_UNAVAILABLE") self._reporting_feed_snapshots[stored.snapshot.snapshot_id] = ( document, hashlib.sha256(document).hexdigest(), stored.signing_key, ) def _capture_feed( self, caller: ReportingDeliveryPrincipal, request: FeedRequest, after: tuple[int, int], consumer_status_enabled: bool, ) -> ReportingFeedSnapshot: owned = tuple( (seq, who, record) for seq, who, record in getattr(self, "_delivery_records", ()) if who == caller or principal(record) == caller ) if any(who != caller or principal(record) != caller for _, who, record in owned): raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT") changes = tuple(ReportingReconciliationChange(seq, record) for seq, _, record in owned) return capture_feed( memory_snapshot(self, caller.account_id), changes, caller=caller, request=request, after=after, consumer_status_enabled=consumer_status_enabled, materializer_boundaries=tuple( b for b in getattr(self, "_materializer_boundaries", ()) if b.caller == caller ), receipt_boundaries=tuple( b for b in getattr(self, "_receipt_boundaries", ()) if b.caller == caller ), **self._feed_projection_options(caller), ) @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) try: async with self._mutation(): stored = None offset = 0 if parsed.cursor is not None: stored = self._feed_snapshot(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 = self._feed_snapshot(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: stored = StoredFeedSnapshot( self._capture_feed(caller, parsed, after, consumer_status_enabled), secrets.token_bytes(32), ) if reauthorize is not None: await reauthorize() self._save_feed_snapshot(stored) elif reauthorize is not None: await reauthorize() return inject_context(request, stored.page(offset, parsed.limit)) except LedgerConflictError: raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT") from None @storage_errors async def read_reporting_feed_snapshot( self, snapshot_id: str, *, caller: ReportingDeliveryPrincipal ) -> ReportingFeedSnapshot | None: async with self._lock: stored = self._feed_snapshot(snapshot_id, caller) return stored.snapshot if stored else None async def reporting_feed_ready(self) -> bool: return TruePersisted for this store's lifetime; production uses PgReportingFeedStore.
Ancestors
- InMemoryReportingReceiptStore
- InMemoryReportingMaterializerStore
- InMemoryReportingReconciliationStore
- InMemoryReportingLedgerStore
- 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) try: async with self._mutation(): stored = None offset = 0 if parsed.cursor is not None: stored = self._feed_snapshot(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 = self._feed_snapshot(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: stored = StoredFeedSnapshot( self._capture_feed(caller, parsed, after, consumer_status_enabled), secrets.token_bytes(32), ) if reauthorize is not None: await reauthorize() self._save_feed_snapshot(stored) elif reauthorize is not None: await reauthorize() return inject_context(request, stored.page(offset, parsed.limit)) except LedgerConflictError: raise ReportingFeedError("REPORTING_FEED_HISTORY_CORRUPT") from None 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._lock: stored = self._feed_snapshot(snapshot_id, caller) return stored.snapshot if stored else None async def reporting_feed_ready(self) ‑> bool-
Expand source code
async def reporting_feed_ready(self) -> bool: return True
Inherited members
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
class ReportingFeedError (code: FeedErrorCode)-
Expand source code
class ReportingFeedError(RuntimeError): def __init__(self, code: FeedErrorCode) -> None: self.code = code super().__init__(_MESSAGES[code])Unspecified run-time error.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class ReportingFeedRecord (kind: FeedKind, record_id: str, key: GlobalKey, wire_json: bytes)-
Expand source code
@dataclass(frozen=True) class ReportingFeedRecord: kind: FeedKind record_id: str key: GlobalKey wire_json: bytes = field(repr=False) def to_storage(self) -> dict[str, Any]: return { "kind": self.kind, "id": self.record_id, "key": list(self.key), "wire": json.loads(self.wire_json), }ReportingFeedRecord(kind: 'FeedKind', record_id: 'str', key: 'GlobalKey', wire_json: 'bytes')
Instance variables
var key : tuple[int, int, str, str]var kind : Literal['obligation', 'revision', 'adjustment', 'consumer_status', 'materialization', 'revision_receipt', 'adjustment_receipt']var record_id : strvar wire_json : bytes
Methods
def to_storage(self) ‑> dict[str, typing.Any]-
Expand source code
def to_storage(self) -> dict[str, Any]: return { "kind": self.kind, "id": self.record_id, "key": list(self.key), "wire": json.loads(self.wire_json), }
class ReportingFeedSnapshot (caller: ReportingDeliveryPrincipal,
snapshot_id: str,
as_of: datetime,
after: tuple[int, int],
through: tuple[int, int],
filters_json: bytes,
common_json: bytes,
records: tuple[ReportingFeedRecord, ...],
inputs_json: bytes,
representation_version: int = 1,
ownership_mode: "Literal['absent', 'bindings']" = 'absent')-
Expand source code
@dataclass(frozen=True) class ReportingFeedSnapshot: caller: ReportingDeliveryPrincipal snapshot_id: str as_of: datetime after: tuple[int, int] through: tuple[int, int] filters_json: bytes = field(repr=False) common_json: bytes = field(repr=False) records: tuple[ReportingFeedRecord, ...] = field(repr=False) inputs_json: bytes = field(repr=False) representation_version: int = 1 ownership_mode: Literal["absent", "bindings"] = "absent" @property def total_count(self) -> int: return len(self.records) @property def inputs(self) -> dict[str, Any]: """A detached PRIVATE copy. Never serialize this into a response/ext.""" return dict(json.loads(self.inputs_json)) def to_storage(self) -> dict[str, Any]: return { "version": 1, "representation_version": self.representation_version, "ownership_mode": self.ownership_mode, "account_id": self.caller.account_id, "consumer_id": self.caller.consumer_id, "snapshot_id": self.snapshot_id, "as_of": self.as_of.isoformat(), "after": list(self.after), "through": list(self.through), # Filters use ordinary JSON; keep their exact versioned bytes inside # the restricted canonical financial-history document as a string. "filters": self.filters_json.decode("ascii"), "common": json.loads(self.common_json), "total_count": self.total_count, "records": [r.to_storage() for r in self.records], "inputs": self.inputs, } @property def binding(self) -> str: # Includes the complete frozen projection and every private dependency, # not merely maxima followed by mutable reads. Principal lengths never # increase token size because the binding is a fixed SHA-256 digest. return _digest(self.to_storage())ReportingFeedSnapshot(caller: 'ReportingDeliveryPrincipal', snapshot_id: 'str', as_of: 'datetime', after: 'tuple[int, int]', through: 'tuple[int, int]', filters_json: 'bytes', common_json: 'bytes', records: 'tuple[ReportingFeedRecord, …]', inputs_json: 'bytes', representation_version: 'int' = 1, ownership_mode: "Literal['absent', 'bindings']" = 'absent')
Instance variables
var after : tuple[int, int]var as_of : datetime.datetimeprop binding : str-
Expand source code
@property def binding(self) -> str: # Includes the complete frozen projection and every private dependency, # not merely maxima followed by mutable reads. Principal lengths never # increase token size because the binding is a fixed SHA-256 digest. return _digest(self.to_storage()) var caller : ReportingDeliveryPrincipalvar common_json : bytesvar filters_json : bytesprop inputs : dict[str, Any]-
Expand source code
@property def inputs(self) -> dict[str, Any]: """A detached PRIVATE copy. Never serialize this into a response/ext.""" return dict(json.loads(self.inputs_json))A detached PRIVATE copy. Never serialize this into a response/ext.
var inputs_json : bytesvar ownership_mode : Literal['absent', 'bindings']var records : tuple[ReportingFeedRecord, ...]var representation_version : intvar snapshot_id : strvar through : tuple[int, int]prop total_count : int-
Expand source code
@property def total_count(self) -> int: return len(self.records)
Methods
def to_storage(self) ‑> dict[str, typing.Any]-
Expand source code
def to_storage(self) -> dict[str, Any]: return { "version": 1, "representation_version": self.representation_version, "ownership_mode": self.ownership_mode, "account_id": self.caller.account_id, "consumer_id": self.caller.consumer_id, "snapshot_id": self.snapshot_id, "as_of": self.as_of.isoformat(), "after": list(self.after), "through": list(self.through), # Filters use ordinary JSON; keep their exact versioned bytes inside # the restricted canonical financial-history document as a string. "filters": self.filters_json.decode("ascii"), "common": json.loads(self.common_json), "total_count": self.total_count, "records": [r.to_storage() for r in self.records], "inputs": self.inputs, }
class ReportingFeedStore (*args, **kwargs)-
Expand source code
@runtime_checkable class ReportingFeedStore(Protocol): """Trusted caller seam. The mounted handler reauthorizes every request. The request-local callback rechecks the same canonical principal before publishing a prepared snapshot or returning a continuation. It is never retained in a store, snapshot, or shared handler state. """ 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]: ... async def read_reporting_feed_snapshot( self, snapshot_id: str, *, caller: ReportingDeliveryPrincipal, ) -> ReportingFeedSnapshot | None: ... async def reporting_feed_ready(self) -> bool: ...Trusted caller seam. The mounted handler reauthorizes every request.
The request-local callback rechecks the same canonical principal before publishing a prepared snapshot or returning a continuation. It is never retained in a store, snapshot, or shared handler state.
Ancestors
- typing.Protocol
- typing.Generic
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
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]: ... async def read_reporting_feed_snapshot(self, snapshot_id: str, *, caller: ReportingDeliveryPrincipal) ‑> ReportingFeedSnapshot | None-
Expand source code
async def read_reporting_feed_snapshot( self, snapshot_id: str, *, caller: ReportingDeliveryPrincipal, ) -> ReportingFeedSnapshot | None: ... async def reporting_feed_ready(self) ‑> bool-
Expand source code
async def reporting_feed_ready(self) -> bool: ...