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 True

Persisted for this store's lifetime; production uses PgReportingFeedStore.

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)
    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 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

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 : str
var 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.datetime
prop 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 : ReportingDeliveryPrincipal
var common_json : bytes
var filters_json : bytes
prop 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 : bytes
var ownership_mode : Literal['absent', 'bindings']
var records : tuple[ReportingFeedRecord, ...]
var representation_version : int
var snapshot_id : str
var 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: ...