Module adcp.reporting.projection

Versioned private reporting projection, independent of optional delivery workers.

Sub-modules

adcp.reporting.projection.capture

Closed immutable v2 inputs preserve legacy C/B2.1/B2.2 document shapes.

adcp.reporting.projection.history

Replay retained private boundaries without creating active notifications …

adcp.reporting.projection.memory

Memory conformance participant with the same unconditional rollback boundary.

adcp.reporting.projection.notifications

Isolated v2 status queues using the original SDK fanout/delivery transactions.

adcp.reporting.projection.pg

Connection-bound v2 capture over the reviewed feed and status participants.

adcp.reporting.projection.schema

B2.4's additive manifest never widens an earlier feature's required objects.

adcp.reporting.projection.wire

Tier-correct exact views over one captured, caller-private financial history.

Classes

class InMemoryReportingProjectionOutbox (store: InMemoryReportingLedgerStore)
Expand source code
class InMemoryReportingProjectionOutbox(InMemoryReportingStatusOutbox):
    _store: InMemoryReportingProjectionStore

    @property
    def _state(self) -> NotificationState:
        state = self._store._projection_outbox
        if state is None:
            raise ReportingNotificationError("notifications_disabled")
        return state

View over an opted-in ledger. Recreating this view preserves its state.

Ancestors

Inherited members

class InMemoryReportingProjectionStore (**kwargs: Any)
Expand source code
class InMemoryReportingProjectionStore(InMemoryReportingFeedStore):
    """Reference semantics, never a production durability claim."""

    _projection_read_policy: dict[str, Any]
    _projection_read_escalation: ReportingDeliveryEscalation | None

    def __init__(self, **kwargs: Any) -> None:
        super().__init__(**kwargs)
        self._projection_accounts: dict[str, _ProjectionAccount] = {}
        self._projection_outbox = (
            NotificationState() if self._notification_state is not None else None
        )

    def _capture_projection(self, account_id: str) -> ReportingProjectionInput:
        return capture_memory_input(
            memory_snapshot(self, account_id), tuple(getattr(self, "_delivery_records", ()))
        )

    @asynccontextmanager
    async def _mutation(self) -> AsyncIterator[None]:
        nested = _MEMORY_TRANSACTION.get() == (id(self), asyncio.current_task())
        async with super()._mutation():
            yield
            if not nested and _REPLAY.get() != id(self):
                for account_id, state in self._projection_accounts.items():
                    captured = self._capture_projection(account_id)
                    identity = source_identity(captured)
                    if identity != state.source:
                        state.inputs.append(captured)
                        state.source = identity

    def _feed_projection_options(self, caller: ReportingDeliveryPrincipal) -> dict[str, Any]:
        state = self._projection_accounts.get(caller.account_id)
        if state is None or not state.ready:
            return {}
        return {
            "representation_version": 2,
            "revision_ownership": state.policy["ownership_enabled"],
            "activated_consumer_status_enabled": state.policy["consumer_status_enabled"],
        }

    async def read_projection_input(
        self, *, caller: ReportingDeliveryPrincipal
    ) -> ReportingProjectionInput:
        async with self._lock:
            state = self._projection_accounts.get(caller.account_id)
            if state is None or not state.ready:
                raise ReportingNotificationError("status_projection_activation_required")
            return state.inputs[-1].at(self._clock())

    async def read_tier_status(
        self,
        request: dict[str, Any],
        *,
        caller: ReportingDeliveryPrincipal,
        consumer_status_enabled: bool = False,
    ) -> dict[str, Any]:
        from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
        from adcp.reporting.projection.wire import render_tier_status

        async with self._lock:
            state = self._projection_accounts.get(caller.account_id)
            value = state.inputs[-1].at(self._clock()) if state and state.ready else None
            policy = state.policy if state and state.ready else None
        if value is None or policy is None:
            return await ReportingStatusHandler(
                self, consumer_status_enabled=consumer_status_enabled
            ).handle(request, caller=ReportingStatusCaller(caller.account_id, caller.consumer_id))
        return render_tier_status(self, request, value, caller, policy, consumer_status_enabled)

Reference semantics, never a production durability claim.

Ancestors

Subclasses

Methods

async def read_projection_input(self, *, caller: ReportingDeliveryPrincipal) ‑> ReportingProjectionInput
Expand source code
async def read_projection_input(
    self, *, caller: ReportingDeliveryPrincipal
) -> ReportingProjectionInput:
    async with self._lock:
        state = self._projection_accounts.get(caller.account_id)
        if state is None or not state.ready:
            raise ReportingNotificationError("status_projection_activation_required")
        return state.inputs[-1].at(self._clock())
async def read_tier_status(self,
request: dict[str, Any],
*,
caller: ReportingDeliveryPrincipal,
consumer_status_enabled: bool = False) ‑> dict[str, typing.Any]
Expand source code
async def read_tier_status(
    self,
    request: dict[str, Any],
    *,
    caller: ReportingDeliveryPrincipal,
    consumer_status_enabled: bool = False,
) -> dict[str, Any]:
    from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
    from adcp.reporting.projection.wire import render_tier_status

    async with self._lock:
        state = self._projection_accounts.get(caller.account_id)
        value = state.inputs[-1].at(self._clock()) if state and state.ready else None
        policy = state.policy if state and state.ready else None
    if value is None or policy is None:
        return await ReportingStatusHandler(
            self, consumer_status_enabled=consumer_status_enabled
        ).handle(request, caller=ReportingStatusCaller(caller.account_id, caller.consumer_id))
    return render_tier_status(self, request, value, caller, policy, consumer_status_enabled)

Inherited members

class InMemoryReportingStatusProjection (ledger: InMemoryReportingProjectionStore,
*,
consumer_status_enabled: bool = False,
revision_ownership: bool = False,
escalation: ReportingDeliveryEscalation | None = None)
Expand source code
class InMemoryReportingStatusProjection(InMemoryStatusNotificationStore):
    _projection_version = 2
    ledger: InMemoryReportingProjectionStore

    def __init__(
        self,
        ledger: InMemoryReportingProjectionStore,
        *,
        consumer_status_enabled: bool = False,
        revision_ownership: bool = False,
        escalation: ReportingDeliveryEscalation | None = None,
    ) -> None:
        if not isinstance(ledger, InMemoryReportingProjectionStore):
            raise ReportingNotificationError("status_projection_component_unready")
        if type(consumer_status_enabled) is not bool or type(revision_ownership) is not bool:
            raise ValueError("projection feature selections must be booleans")
        self.ledger, self.escalation = ledger, escalation
        self.consumer_status_enabled, self.revision_ownership = (
            consumer_status_enabled,
            revision_ownership,
        )
        if ledger._status_notification_state is None:
            ledger._status_notification_state = _StatusMemoryState()
        self.outbox = InMemoryReportingProjectionOutbox(ledger)
        ledger._projection_read_policy = self.policy
        ledger._projection_read_escalation = escalation

    def _enqueue_status(self, event: ReportingDomainEvent) -> None:
        self.outbox._state.enqueue(event)

    async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]:
        # A checkpoint contains a mutable JSON projection. Never lend that
        # dictionary to a reader or retain it in a caller's published result.
        return deepcopy(await super().checkpoints(account_id=account_id))

    @property
    def policy(self) -> dict[str, Any]:
        return {
            "version": 2,
            "escalation": escalation_identity(self.escalation),
            "consumer_status_enabled": self.consumer_status_enabled,
            "ownership_enabled": self.revision_ownership,
            "notifications_enabled": self.ledger._notification_state is not None,
        }

    @asynccontextmanager
    async def _transaction(self) -> AsyncIterator[None]:
        token = _REPLAY.set(id(self.ledger))
        try:
            async with self.ledger._mutation():
                yield
        finally:
            _REPLAY.reset(token)

    def _account(self, account_id: str) -> _ProjectionAccount:
        state = self.ledger._projection_accounts.get(account_id)
        if state is None:
            raise ReportingNotificationError("status_projection_activation_required")
        if state.policy != self.policy:
            raise ReportingNotificationError("status_projection_policy_conflict")
        return state

    def _cursor(self, account_id: str) -> int:
        self._account(account_id)
        return super()._cursor(account_id)

    async def baseline_ready(self, *, account_id: str) -> bool:
        async with self.ledger._lock:
            if account_id not in self.ledger._projection_accounts:
                return False
            return self._account(account_id).ready

    def _apply_value(
        self, value: ReportingProjectionInput, *, through: int, silent: bool = False
    ) -> int:
        return self._apply(
            settled_replay(value.core),
            through=through,
            reconciliation=value.reconciliation,
            consumer_status_enabled=self.consumer_status_enabled,
            enqueue=self.ledger._notification_state is not None and not silent,
        )

    async def activate(self, *, account_id: str) -> bool:
        started = await self._begin_activation(account_id=account_id)
        while True:
            async with self._transaction():
                state = self._account(account_id)
                if state.ready:
                    return started
                self._activation_step(account_id)

    async def _begin_activation(self, *, account_id: str) -> bool:
        async with self._transaction():
            if account_id in self.ledger._projection_accounts:
                self._account(account_id)
                return False
            old = self._state.accounts.get(account_id)
            notifications = self.ledger._notification_state
            through = (
                max(
                    (d.sequence for d in notifications.dirty if d.scope.account_id == account_id),
                    default=0,
                )
                if notifications
                else 0
            )
            if old is not None:
                if old[1] != escalation_identity(self.escalation):
                    raise ReportingNotificationError("status_policy_conflict")
                if old[0] != through or self._needs_rebuild(account_id):
                    raise ReportingNotificationError("status_projection_legacy_drain_required")
            checkpoints = [
                c for c in self._state.checkpoints.values() if c.scope.account_id == account_id
            ]
            now = self.ledger._clock()
            if any(
                (c.lease_expires_at and c.lease_expires_at > now)
                or (c.next_due_at and c.next_due_at <= now)
                for c in checkpoints
            ):
                raise ReportingNotificationError("status_projection_legacy_drain_required")
            value = self.ledger._capture_projection(account_id)
            floor = max((c.source_sequence for c in checkpoints), default=0)
            legacy = tuple(
                sorted(
                    (
                        b
                        for b in (
                            *self.ledger._materializer_boundaries,
                            *getattr(self.ledger, "_receipt_boundaries", ()),
                        )
                        if b.caller.account_id == account_id
                    ),
                    key=lambda b: b.account_sequence,
                )
            )
            expected = self.ledger._materializer_account_heads.get(account_id, 0)
            if len(legacy) != expected or any(
                b.account_sequence != n for n, b in enumerate(legacy, 1)
            ):
                raise ReportingNotificationError("status_projection_history_corrupt")
            self.ledger._projection_accounts[account_id] = _ProjectionAccount(
                self.policy,
                floor + expected,
                value,
                source_identity(value),
                inputs=[value],
                legacy=legacy,
                baselines={checkpoint_key(c): c for c in checkpoints},
            )
            self._state.accounts[account_id] = (through, escalation_identity(self.escalation))
            self._state.selector_accounts[account_id] = "complete"
            return True

    def _activation_step(self, account_id: str) -> StatusTurn:
        state = self._account(account_id)
        if state.legacy_cursor < len(state.legacy):
            boundary = state.legacy[state.legacy_cursor]
            for step in project_boundary(
                boundary,
                state.historical_checkpoints,
                baselines=state.baselines,
                source_sequence=state.checkpoint_floor
                - len(state.legacy)
                + boundary.account_sequence,
                escalation=self.escalation,
                consumer_status_enabled=self.consumer_status_enabled,
            ):
                state.historical_checkpoints[checkpoint_key(step.checkpoint)] = step.checkpoint
                state.historical_steps.append((boundary.account_sequence, step.document()))
                self._state.checkpoints[step.checkpoint.scope.checkpoint_key] = step.checkpoint
            state.legacy_cursor += 1
            return StatusTurn(True)
        self._apply_value(state.inputs[0], through=state.checkpoint_floor + 1, silent=True)
        state.cursor, state.ready = 1, True
        return StatusTurn(True)

    async def baseline(self, *, account_id: str) -> bool:
        return await self.activate(account_id=account_id)

    def _project_version(self, account_id: str) -> StatusTurn:
        state = self._account(account_id)
        if not state.ready:
            return self._activation_step(account_id)
        following = state.inputs[state.cursor] if state.cursor < len(state.inputs) else None
        due = min(
            (
                c.next_due_at
                for c in self._state.checkpoints.values()
                if c.scope.account_id == account_id and c.next_due_at is not None
            ),
            default=None,
        )
        cursor = state.cursor
        if (
            due is not None
            and due <= self.ledger._clock()
            and (following is None or due < following.core.as_of)
        ):
            value = state.current.at(due)
        elif following is not None:
            value, cursor = following, cursor + 1
        else:
            return StatusTurn(False)
        if value.core.as_of < state.current.core.as_of:
            raise ReportingNotificationError("status_projection_clock_regressed")
        value = with_projection_core(
            value, settled_replay(with_replay_lifecycles(value.core, state.current.core))
        )
        count = self._apply_value(value, through=state.checkpoint_floor + cursor)
        state.current, state.cursor = value, cursor
        return StatusTurn(True, count)

    async def project_one(self, *, account_id: str) -> StatusTurn:
        async with self._transaction():
            return self._project_version(account_id)

    async def rebuild_one(self) -> StatusTurn:
        async with self._transaction():
            for account_id, state in sorted(self.ledger._projection_accounts.items()):
                if state.cursor < len(state.inputs) and state.policy == self.policy:
                    return self._project_version(account_id)
            return StatusTurn(False)

    async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
        async with self._transaction():
            self._account(lease.scope.account_id)
            checkpoint = self._state.checkpoints.get(lease.scope.checkpoint_key)
            if (
                checkpoint is None
                or checkpoint.lease_token != lease.token
                or checkpoint.lease_expires_at != lease.expires_at
                or lease.expires_at <= self.ledger._clock()
            ):
                return StatusTurn(False)
            result = self._project_version(lease.scope.account_id)
            current = self._state.checkpoints[lease.scope.checkpoint_key]
            if current.lease_token != lease.token or lease.expires_at <= self.ledger._clock():
                raise ReportingNotificationError("status_lease_lost")
            self._state.checkpoints[lease.scope.checkpoint_key] = replace(
                current,
                lease_token=None,
                lease_expires_at=None,
            )
            return result

    async def sweep_one(self) -> StatusTurn:
        async with self.ledger._lock:
            at = self.ledger._clock()
            due = sorted(
                (c.next_due_at, c.scope.account_id)
                for c in self._state.checkpoints.values()
                if c.next_due_at is not None
                and c.next_due_at <= at
                and c.scope.account_id in self.ledger._projection_accounts
                and self.ledger._projection_accounts[c.scope.account_id].policy == self.policy
                and (c.lease_expires_at is None or c.lease_expires_at <= at)
            )
        if not due:
            return StatusTurn(False)
        lease = await self.claim_due(account_id=due[0][1])
        return await self.complete_due(lease) if lease is not None else StatusTurn(False)

Ancestors

Class variables

var ledger : InMemoryReportingProjectionStore

Instance variables

prop policy : dict[str, Any]
Expand source code
@property
def policy(self) -> dict[str, Any]:
    return {
        "version": 2,
        "escalation": escalation_identity(self.escalation),
        "consumer_status_enabled": self.consumer_status_enabled,
        "ownership_enabled": self.revision_ownership,
        "notifications_enabled": self.ledger._notification_state is not None,
    }

Methods

async def activate(self, *, account_id: str) ‑> bool
Expand source code
async def activate(self, *, account_id: str) -> bool:
    started = await self._begin_activation(account_id=account_id)
    while True:
        async with self._transaction():
            state = self._account(account_id)
            if state.ready:
                return started
            self._activation_step(account_id)
async def baseline(self, *, account_id: str) ‑> bool
Expand source code
async def baseline(self, *, account_id: str) -> bool:
    return await self.activate(account_id=account_id)
async def baseline_ready(self, *, account_id: str) ‑> bool
Expand source code
async def baseline_ready(self, *, account_id: str) -> bool:
    async with self.ledger._lock:
        if account_id not in self.ledger._projection_accounts:
            return False
        return self._account(account_id).ready
async def checkpoints(self, *, account_id: str) ‑> tuple[StatusCheckpoint, ...]
Expand source code
async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]:
    # A checkpoint contains a mutable JSON projection. Never lend that
    # dictionary to a reader or retain it in a caller's published result.
    return deepcopy(await super().checkpoints(account_id=account_id))
async def complete_due(self, lease: StatusDueLease) ‑> StatusTurn
Expand source code
async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
    async with self._transaction():
        self._account(lease.scope.account_id)
        checkpoint = self._state.checkpoints.get(lease.scope.checkpoint_key)
        if (
            checkpoint is None
            or checkpoint.lease_token != lease.token
            or checkpoint.lease_expires_at != lease.expires_at
            or lease.expires_at <= self.ledger._clock()
        ):
            return StatusTurn(False)
        result = self._project_version(lease.scope.account_id)
        current = self._state.checkpoints[lease.scope.checkpoint_key]
        if current.lease_token != lease.token or lease.expires_at <= self.ledger._clock():
            raise ReportingNotificationError("status_lease_lost")
        self._state.checkpoints[lease.scope.checkpoint_key] = replace(
            current,
            lease_token=None,
            lease_expires_at=None,
        )
        return result
async def project_one(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def project_one(self, *, account_id: str) -> StatusTurn:
    async with self._transaction():
        return self._project_version(account_id)
async def rebuild_one(self) ‑> StatusTurn
Expand source code
async def rebuild_one(self) -> StatusTurn:
    async with self._transaction():
        for account_id, state in sorted(self.ledger._projection_accounts.items()):
            if state.cursor < len(state.inputs) and state.policy == self.policy:
                return self._project_version(account_id)
        return StatusTurn(False)
async def sweep_one(self) ‑> StatusTurn
Expand source code
async def sweep_one(self) -> StatusTurn:
    async with self.ledger._lock:
        at = self.ledger._clock()
        due = sorted(
            (c.next_due_at, c.scope.account_id)
            for c in self._state.checkpoints.values()
            if c.next_due_at is not None
            and c.next_due_at <= at
            and c.scope.account_id in self.ledger._projection_accounts
            and self.ledger._projection_accounts[c.scope.account_id].policy == self.policy
            and (c.lease_expires_at is None or c.lease_expires_at <= at)
        )
    if not due:
        return StatusTurn(False)
    lease = await self.claim_due(account_id=due[0][1])
    return await self.complete_due(lease) if lease is not None else StatusTurn(False)
class PgReportingProjectionOutbox (*, pool: AsyncConnectionPool, clock: Callable[[], datetime] | None = None)
Expand source code
class PgReportingProjectionOutbox(PgReportingStatusOutbox):
    @asynccontextmanager
    async def _connection(self) -> AsyncIterator[Any]:
        async with self._pool.connection() as connection:
            yield _ProjectionQueueConnection(connection)

    async def create_schema(self) -> None:
        from adcp.reporting.projection.pg import PgReportingProjectionStore

        await PgReportingProjectionStore(pool=self._pool, notifications=True).create_schema()

Separate C events, expansion, delivery and activity, with generic HTTP logic.

Ancestors

Methods

async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    from adcp.reporting.projection.pg import PgReportingProjectionStore

    await PgReportingProjectionStore(pool=self._pool, notifications=True).create_schema()

Inherited members

class PgReportingProjectionStore (*,
pool: AsyncConnectionPool,
clock: Callable[[], datetime] | None = None,
notifications: bool = False)
Expand source code
class PgReportingProjectionStore(PgReportingFeedStore):
    """Optional v2 store. Installation alone neither activates tiers nor sends readiness."""

    _projection_read_policy: dict[str, Any]
    _projection_read_escalation: ReportingDeliveryEscalation | None

    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",
                "reporting_status_notifications.sql",
                "reporting_status_selector_version.sql",
                "reporting_projection.sql",
                "reporting_projection_notifications.sql",
                "reporting_projection_feed.sql",
            ):
                await connection.execute(root.joinpath(name).read_text())

    async def _feed_snapshot_on(
        self, connection: Any, snapshot_id: str, caller: ReportingDeliveryPrincipal
    ) -> StoredFeedSnapshot | None:
        legacy = await super()._feed_snapshot_on(connection, snapshot_id, caller)
        if legacy is not None:
            return legacy
        return await super()._feed_snapshot_on(
            _ProjectionFeedConnection(connection), snapshot_id, caller
        )

    async def _save_feed_snapshot_on(self, connection: Any, stored: _PreparedFeed) -> None:
        if stored.snapshot.representation_version == 2:
            connection = _ProjectionFeedConnection(connection)
        await super()._save_feed_snapshot_on(connection, stored)

    async def _capture_feed_on(
        self, connection: Any, caller: ReportingDeliveryPrincipal
    ) -> _CapturedFeed:
        captured = await super()._capture_feed_on(connection, caller)
        row = await (
            await connection.execute(
                "SELECT ownership_enabled,consumer_status_enabled"
                " FROM reporting_projection_accounts WHERE account_id=%s"
                " AND current_input IS NOT NULL",
                (caller.account_id,),
            )
        ).fetchone()
        if row is None:
            return captured
        await validate_projection_schema(connection, notifications=self._notifications_enabled)
        return replace(
            captured,
            representation_version=2,
            revision_ownership=row[0],
            activated_consumer_status_enabled=row[1],
        )

    async def read_projection_input(
        self, *, caller: ReportingDeliveryPrincipal
    ) -> ReportingProjectionInput:
        async with self._connection() as connection, connection.transaction():
            await self._lock_account(connection, caller.account_id)
            await validate_projection_schema(connection, notifications=self._notifications_enabled)
            row = await (
                await connection.execute(
                    "SELECT input,content_sha256=reporting_receipt_ingestion_sha256(input)"
                    " FROM reporting_projection_inputs WHERE account_id=%s"
                    " AND EXISTS (SELECT 1 FROM reporting_projection_accounts a"
                    " WHERE a.account_id=reporting_projection_inputs.account_id"
                    " AND a.current_input IS NOT NULL)"
                    " ORDER BY sequence DESC LIMIT 1",
                    (caller.account_id,),
                )
            ).fetchone()
            if row is None:
                raise ReportingNotificationError("status_projection_activation_required")
            if not row[1]:
                raise ReportingNotificationError("status_projection_history_corrupt")
            value = decode_projection_input(row[0])
            return value.at(await database_now(connection, self._clock))

    async def read_tier_status(
        self,
        request: dict[str, Any],
        *,
        caller: ReportingDeliveryPrincipal,
        consumer_status_enabled: bool = False,
    ) -> dict[str, Any]:
        from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
        from adcp.reporting.projection.wire import render_tier_status

        async with self._connection() as connection:
            row = await (
                await connection.execute(
                    "SELECT policy FROM reporting_projection_accounts WHERE account_id=%s"
                    " AND current_input IS NOT NULL",
                    (caller.account_id,),
                )
            ).fetchone()
        if row is None:
            return await ReportingStatusHandler(
                self, consumer_status_enabled=consumer_status_enabled
            ).handle(request, caller=ReportingStatusCaller(caller.account_id, caller.consumer_id))
        value = await self.read_projection_input(caller=caller)
        return await asyncio.to_thread(
            render_tier_status, self, request, value, caller, row[0], consumer_status_enabled
        )

Optional v2 store. Installation alone neither activates tiers nor sends readiness.

Ancestors

Subclasses

Methods

async def read_projection_input(self, *, caller: ReportingDeliveryPrincipal) ‑> ReportingProjectionInput
Expand source code
async def read_projection_input(
    self, *, caller: ReportingDeliveryPrincipal
) -> ReportingProjectionInput:
    async with self._connection() as connection, connection.transaction():
        await self._lock_account(connection, caller.account_id)
        await validate_projection_schema(connection, notifications=self._notifications_enabled)
        row = await (
            await connection.execute(
                "SELECT input,content_sha256=reporting_receipt_ingestion_sha256(input)"
                " FROM reporting_projection_inputs WHERE account_id=%s"
                " AND EXISTS (SELECT 1 FROM reporting_projection_accounts a"
                " WHERE a.account_id=reporting_projection_inputs.account_id"
                " AND a.current_input IS NOT NULL)"
                " ORDER BY sequence DESC LIMIT 1",
                (caller.account_id,),
            )
        ).fetchone()
        if row is None:
            raise ReportingNotificationError("status_projection_activation_required")
        if not row[1]:
            raise ReportingNotificationError("status_projection_history_corrupt")
        value = decode_projection_input(row[0])
        return value.at(await database_now(connection, self._clock))
async def read_tier_status(self,
request: dict[str, Any],
*,
caller: ReportingDeliveryPrincipal,
consumer_status_enabled: bool = False) ‑> dict[str, typing.Any]
Expand source code
async def read_tier_status(
    self,
    request: dict[str, Any],
    *,
    caller: ReportingDeliveryPrincipal,
    consumer_status_enabled: bool = False,
) -> dict[str, Any]:
    from adcp.reporting.ledger.status import ReportingStatusCaller, ReportingStatusHandler
    from adcp.reporting.projection.wire import render_tier_status

    async with self._connection() as connection:
        row = await (
            await connection.execute(
                "SELECT policy FROM reporting_projection_accounts WHERE account_id=%s"
                " AND current_input IS NOT NULL",
                (caller.account_id,),
            )
        ).fetchone()
    if row is None:
        return await ReportingStatusHandler(
            self, consumer_status_enabled=consumer_status_enabled
        ).handle(request, caller=ReportingStatusCaller(caller.account_id, caller.consumer_id))
    value = await self.read_projection_input(caller=caller)
    return await asyncio.to_thread(
        render_tier_status, self, request, value, caller, row[0], consumer_status_enabled
    )

Inherited members

class PgReportingStatusProjection (ledger: PgReportingProjectionStore,
*,
consumer_status_enabled: bool = False,
revision_ownership: bool = False,
escalation: ReportingDeliveryEscalation | None = None)
Expand source code
class PgReportingStatusProjection(PgStatusNotificationStore):
    """v2 input/cursor lifecycle reusing C's pure projector, checkpoint and event transaction."""

    _projection_version = 2
    ledger: PgReportingProjectionStore

    def __init__(
        self,
        ledger: PgReportingProjectionStore,
        *,
        consumer_status_enabled: bool = False,
        revision_ownership: bool = False,
        escalation: ReportingDeliveryEscalation | None = None,
    ) -> None:
        if not isinstance(ledger, PgReportingProjectionStore):
            raise ReportingNotificationError("status_projection_component_unready")
        if type(consumer_status_enabled) is not bool or type(revision_ownership) is not bool:
            raise ValueError("projection feature selections must be booleans")
        self.ledger, self.escalation = ledger, escalation
        self.consumer_status_enabled, self.revision_ownership = (
            consumer_status_enabled,
            revision_ownership,
        )
        self.outbox = PgReportingProjectionOutbox(pool=ledger._pool, clock=ledger._clock)
        ledger._projection_read_policy = self.policy
        ledger._projection_read_escalation = escalation

    async def _enqueue_status_on(self, connection: Any, event: ReportingDomainEvent) -> None:
        await _enqueue_on(_ProjectionQueueConnection(connection), event)

    @property
    def policy(self) -> dict[str, Any]:
        return {
            "version": 2,
            "escalation": escalation_identity(self.escalation),
            "consumer_status_enabled": self.consumer_status_enabled,
            "ownership_enabled": self.revision_ownership,
            "notifications_enabled": self.ledger._notifications_enabled,
        }

    async def create_schema(self) -> None:
        await self.ledger.create_schema()

    @asynccontextmanager
    async def _transaction(self, account_id: str) -> AsyncIterator[Any]:
        async with self.ledger._connection() as connection, connection.transaction():
            await connection.execute(
                "SELECT set_config('adcp.reporting.projection_version','2',true)"
            )
            await connection.execute(
                "SELECT set_config('adcp.reporting.selector_semantics_version','2',true)"
            )
            await connection.execute(
                "SELECT set_config('adcp.reporting.projection_internal','on',true)"
            )
            await self.ledger._lock_account(connection, account_id)
            if self.ledger._clock is not None:
                await connection.execute(
                    "SELECT set_config('adcp.reporting.projection_clock',%s,true)",
                    (self.ledger._clock().isoformat(),),
                )
            yield connection

    async def _state_on(self, connection: Any, account_id: str) -> tuple[Any, ...]:
        row = await (
            await connection.execute(
                "SELECT policy,cursor,max_sequence,checkpoint_floor,current_input,current_as_of"
                " FROM reporting_projection_accounts WHERE account_id=%s FOR UPDATE",
                (account_id,),
            )
        ).fetchone()
        if row is None:
            raise ReportingNotificationError("status_projection_activation_required")
        if row[0] != self.policy:
            raise ReportingNotificationError("status_projection_policy_conflict")
        return tuple(row)

    async def _account_on(self, connection: Any, account_id: str) -> int:
        await self._state_on(connection, account_id)
        return await super()._account_on(connection, account_id)

    async def baseline_ready(self, *, account_id: str) -> bool:
        async with self._transaction(account_id) as connection:
            await validate_projection_schema(
                connection, notifications=self.ledger._notifications_enabled
            )
            row = await (
                await connection.execute(
                    "SELECT policy,current_input IS NOT NULL FROM reporting_projection_accounts"
                    " WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            if row is None:
                return False
            if row[0] != self.policy:
                raise ReportingNotificationError("status_projection_policy_conflict")
            return bool(row[1])

    async def activate(self, *, account_id: str) -> bool:
        """Fence, replay captured private history, then activate the frozen baseline.

        Each retained input advances in its own transaction. Cancellation or a
        process crash resumes at that exact cursor. The derived historical
        events stay in an immutable epoch-zero journal, never a delivery queue.
        Old C projectors/sweepers must first be stopped and drained.
        """
        started = await self._begin_activation(account_id=account_id)
        while True:
            async with self._transaction(account_id) as connection:
                state = await self._state_on(connection, account_id)
                if state[4] is not None:
                    return started
                await self._activation_step_on(connection, account_id)

    async def _begin_activation(self, *, account_id: str) -> bool:
        async with self._transaction(account_id) as connection:
            await validate_projection_schema(
                connection, notifications=self.ledger._notifications_enabled
            )
            existing = await (
                await connection.execute(
                    "SELECT policy FROM reporting_projection_accounts WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            if existing is not None:
                await self._state_on(connection, account_id)
                return False
            now = await database_now(connection, self.ledger._clock)
            old = await (
                await connection.execute(
                    "SELECT dirty_sequence,baseline_complete,policy FROM reporting_status_accounts"
                    " WHERE account_id=%s FOR UPDATE",
                    (account_id,),
                )
            ).fetchone()
            dirty = await (
                await connection.execute(
                    "SELECT coalesce(max_sequence,0) FROM reporting_status_dirty_heads"
                    " WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            through = int(dirty[0]) if dirty else 0
            if old is not None and old[1]:
                if old[2] != escalation_identity(self.escalation):
                    raise ReportingNotificationError("status_policy_conflict")
                if int(old[0]) != through or await self._needs_rebuild_on(connection, account_id):
                    raise ReportingNotificationError("status_projection_legacy_drain_required")
            rows = await (
                await connection.execute(
                    "SELECT source_sequence,lease_expires_at,next_due_at FROM"
                    " reporting_status_scope_checkpoints"
                    " WHERE account_id=%s ORDER BY"
                    " consumer_namespace,delivery_config_id,version,scope_kind,"
                    " obligation_namespace FOR UPDATE",
                    (account_id,),
                )
            ).fetchall()
            if any(
                (r[1] is not None and r[1] > now) or (r[2] is not None and r[2] <= now)
                for r in rows
            ):
                raise ReportingNotificationError("status_projection_legacy_drain_required")
            floor = max((int(r[0]) for r in rows), default=0)
            # The inherited SDK column list is fixed; the account is bound.
            old_checkpoints = [
                _legacy_checkpoint(r)
                for r in await (
                    await connection.execute(
                        f"SELECT {_LEGACY_CHECKPOINT} FROM reporting_status_scope_checkpoints"  # nosec B608
                        " WHERE account_id=%s",
                        (account_id,),
                    )
                ).fetchall()
            ]
            generation_floor = max((c.generation for c in old_checkpoints if c), default=0)
            legacy_head = await (
                await connection.execute(
                    "SELECT captured_sequence FROM reporting_materializer_accounts"
                    " WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            legacy_through = int(legacy_head[0]) if legacy_head else 0
            await connection.execute(
                "INSERT INTO"
                " reporting_status_accounts(account_id,policy,baseline_complete,dirty_sequence,"
                "baseline_highwater,baseline_at,selector_target_version,selector_transition)"
                " VALUES(%s,%s::jsonb,TRUE,%s,%s,%s,2,'complete')"
                " ON CONFLICT(account_id) DO NOTHING",
                (
                    account_id,
                    json.dumps(escalation_identity(self.escalation)),
                    through,
                    through,
                    now,
                ),
            )
            await connection.execute(
                "INSERT INTO"
                " reporting_projection_accounts(account_id,activated_at,notifications_enabled,"
                "consumer_status_enabled,ownership_enabled,policy,legacy_through,checkpoint_floor,"
                "legacy_capture_through,legacy_generation_floor)"
                " VALUES(%s,%s,%s,%s,%s,%s::jsonb,%s,%s,%s,%s)",
                (
                    account_id,
                    now,
                    self.ledger._notifications_enabled,
                    self.consumer_status_enabled,
                    self.revision_ownership,
                    json.dumps(self.policy),
                    through,
                    floor + legacy_through,
                    legacy_through,
                    generation_floor,
                ),
            )
            for checkpoint in old_checkpoints:
                if checkpoint is None:
                    continue
                document = json.dumps(checkpoint_document(checkpoint))
                await connection.execute(
                    "INSERT INTO reporting_projection_legacy_baselines VALUES"
                    " (%s,%s,%s::jsonb,reporting_receipt_ingestion_sha256(%s::jsonb))",
                    (account_id, checkpoint_key(checkpoint), document, document),
                )
            await connection.execute(
                "UPDATE reporting_status_scope_checkpoints SET projection_writer_floor=2,"
                "lease_token=NULL,lease_expires_at=NULL WHERE account_id=%s",
                (account_id,),
            )
            await connection.execute("SELECT reporting_projection_capture(%s)", (account_id,))
            return True

    async def _activation_step_on(self, connection: Any, account_id: str) -> StatusTurn:
        metadata = await (
            await connection.execute(
                "SELECT legacy_capture_through,legacy_capture_cursor,"
                "checkpoint_floor FROM reporting_projection_accounts WHERE account_id=%s",
                (account_id,),
            )
        ).fetchone()
        through, cursor, floor = map(int, metadata)
        if cursor < through:
            row = await (
                await connection.execute(
                    "SELECT kind,input,content_sha256=reporting_receipt_ingestion_sha256(input),"
                    " count(*) OVER (PARTITION BY account_sequence)"
                    " FROM (SELECT 'materializer' AS kind,account_sequence,input,content_sha256"
                    " FROM reporting_materializer_status_boundaries WHERE account_id=%s"
                    " AND account_sequence>%s AND account_sequence<=%s UNION ALL"
                    " SELECT 'receipt',account_sequence,input,content_sha256"
                    " FROM reporting_receipt_ingestion_boundaries WHERE account_id=%s"
                    " AND account_sequence>%s AND account_sequence<=%s) b"
                    " ORDER BY account_sequence LIMIT 1",
                    (account_id, cursor, through, account_id, cursor, through),
                )
            ).fetchone()
            if row is None or not row[2] or row[3] != 1:
                raise ReportingNotificationError("status_projection_history_corrupt")
            boundary = decode_boundary(row[0], row[1])
            if boundary.account_sequence != cursor + 1 or boundary.caller.account_id != account_id:
                raise ReportingNotificationError("status_projection_history_corrupt")
            previous_rows = await (
                await connection.execute(
                    "SELECT scope_key,checkpoint FROM reporting_projection_legacy_checkpoints"
                    " WHERE account_id=%s AND consumer_id=%s ORDER BY scope_key",
                    (account_id, boundary.caller.consumer_id),
                )
            ).fetchall()
            previous = {r[0]: decode_checkpoint(r[1]) for r in previous_rows}
            baseline_rows = await (
                await connection.execute(
                    "SELECT scope_key,checkpoint FROM reporting_projection_legacy_baselines"
                    " WHERE account_id=%s ORDER BY scope_key",
                    (account_id,),
                )
            ).fetchall()
            baselines = {r[0]: decode_checkpoint(r[1]) for r in baseline_rows}
            steps = project_boundary(
                boundary,
                previous,
                baselines=baselines,
                source_sequence=floor - through + boundary.account_sequence,
                escalation=self.escalation,
                consumer_status_enabled=self.consumer_status_enabled,
            )
            raw = json.dumps(row[1])
            await connection.execute(
                "INSERT INTO reporting_projection_legacy_inputs VALUES"
                " (%s,%s,%s,%s,%s::jsonb,reporting_receipt_ingestion_sha256(%s::jsonb))",
                (
                    account_id,
                    boundary.account_sequence,
                    boundary.caller.consumer_id,
                    row[0],
                    raw,
                    raw,
                ),
            )
            await self._lock_scopes_on(connection, boundary.core)
            for step in steps:
                key = checkpoint_key(step.checkpoint)
                document = json.dumps(step.document())
                await connection.execute(
                    "INSERT INTO reporting_projection_legacy_steps VALUES"
                    " (%s,%s,%s,%s::jsonb,reporting_receipt_ingestion_sha256(%s::jsonb))",
                    (account_id, boundary.account_sequence, key, document, document),
                )
                checkpoint = json.dumps(checkpoint_document(step.checkpoint))
                await connection.execute(
                    "INSERT INTO reporting_projection_legacy_checkpoints VALUES"
                    " (%s,%s,%s,%s::jsonb,reporting_receipt_ingestion_sha256(%s::jsonb))"
                    " ON CONFLICT(account_id,scope_key) DO UPDATE SET"
                    " checkpoint=EXCLUDED.checkpoint,content_sha256=EXCLUDED.content_sha256",
                    (account_id, boundary.caller.consumer_id, key, checkpoint, checkpoint),
                )
                # Advance the original guarded checkpoint once per retained
                # boundary. Only the immutable epoch-zero journal receives the
                # event; activation never releases it into an active queue.
                await self._write_on(connection, step.checkpoint)
            await connection.execute(
                "UPDATE reporting_projection_accounts SET legacy_capture_cursor=%s"
                " WHERE account_id=%s",
                (boundary.account_sequence, account_id),
            )
            return StatusTurn(True)
        row = await (
            await connection.execute(
                "SELECT input FROM reporting_projection_inputs WHERE account_id=%s AND sequence=1",
                (account_id,),
            )
        ).fetchone()
        value = decode_projection_input(row[0])
        await self._lock_scopes_on(connection, value.core)
        await self._apply_value_on(connection, value, through=floor + 1, silent=True)
        await connection.execute(
            "UPDATE reporting_projection_accounts SET"
            " cursor=1,current_input=%s::jsonb,current_as_of=%s WHERE account_id=%s",
            (value.document.decode(), value.core.as_of, account_id),
        )
        return StatusTurn(True)

    async def baseline(self, *, account_id: str) -> bool:
        return await self.activate(account_id=account_id)

    async def _apply_value_on(
        self,
        connection: Any,
        value: ReportingProjectionInput,
        *,
        through: int,
        silent: bool = False,
    ) -> int:
        return await self._apply_on(
            connection,
            settled_replay(value.core),
            through=through,
            reconciliation=value.reconciliation,
            consumer_status_enabled=self.consumer_status_enabled,
            enqueue=self.ledger._notifications_enabled and not silent,
        )

    async def _project_version_on(self, connection: Any, account_id: str) -> StatusTurn:
        state = await self._state_on(connection, account_id)
        if state[4] is None:
            return await self._activation_step_on(connection, account_id)
        cursor, floor = int(state[1]), int(state[3])
        row = await (
            await connection.execute(
                "SELECT sequence,input,content_sha256=reporting_receipt_ingestion_sha256(input)"
                " FROM reporting_projection_inputs WHERE account_id=%s AND sequence>%s"
                " ORDER BY sequence LIMIT 1",
                (account_id, cursor),
            )
        ).fetchone()
        next_value = None
        if row is not None:
            if row[0] != cursor + 1 or not row[2]:
                raise ReportingNotificationError("status_projection_history_corrupt")
            next_value = decode_projection_input(row[1])
        now = await database_now(connection, self.ledger._clock)
        deadline = await (
            await connection.execute(
                "SELECT min(next_due_at) FROM reporting_status_scope_checkpoints"
                " WHERE account_id=%s",
                (account_id,),
            )
        ).fetchone()
        due = deadline[0] if deadline else None
        if due is not None and due <= now and (next_value is None or due < next_value.core.as_of):
            value = decode_projection_input(state[4]).at(due)
            if due < state[5]:
                raise ReportingNotificationError("status_projection_clock_regressed")
        elif next_value is not None:
            value, cursor = next_value, int(row[0])
            if value.core.as_of < state[5]:
                raise ReportingNotificationError("status_projection_clock_regressed")
        else:
            return StatusTurn(False)
        prior = decode_projection_input(state[4])
        value = with_projection_core(
            value, settled_replay(with_replay_lifecycles(value.core, prior.core))
        )
        await persist_replay_lifecycles_on(self.ledger, connection, value.core)
        count = await self._apply_value_on(connection, value, through=floor + cursor)
        await connection.execute(
            "UPDATE reporting_projection_accounts SET"
            " cursor=%s,current_input=%s::jsonb,current_as_of=%s"
            " WHERE account_id=%s",
            (cursor, value.document.decode(), value.core.as_of, account_id),
        )
        return StatusTurn(True, count)

    async def project_one(self, *, account_id: str) -> StatusTurn:
        async with self._transaction(account_id) as connection:
            return await self._project_version_on(connection, account_id)

    async def rebuild_one(self) -> StatusTurn:
        async with self.ledger._connection() as connection:
            row = await (
                await connection.execute(
                    "SELECT account_id FROM reporting_projection_accounts WHERE cursor<max_sequence"
                    " AND policy=%s::jsonb ORDER BY account_id LIMIT 1",
                    (json.dumps(self.policy),),
                )
            ).fetchone()
        return await self.project_one(account_id=row[0]) if row else StatusTurn(False)

    async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
        async with self._transaction(lease.scope.account_id) as connection:
            await self._state_on(connection, lease.scope.account_id)
            if not await self._held_on(connection, lease):
                return StatusTurn(False)
            result = await self._project_version_on(connection, lease.scope.account_id)
            if not await self._ack_on(connection, lease):
                raise ReportingNotificationError("status_lease_lost")
            return result

    async def sweep_one(self) -> StatusTurn:
        async with self.ledger._connection() as connection:
            at = await database_now(connection, self.ledger._clock)
            row = await (
                await connection.execute(
                    "SELECT c.account_id FROM reporting_status_scope_checkpoints c"
                    " JOIN reporting_projection_accounts a ON a.account_id=c.account_id"
                    " WHERE c.next_due_at<=%s AND a.policy=%s::jsonb"
                    " AND (c.lease_expires_at IS NULL OR c.lease_expires_at<=%s)"
                    " ORDER BY c.next_due_at,c.account_id LIMIT 1",
                    (at, json.dumps(self.policy), at),
                )
            ).fetchone()
        if row is None:
            return StatusTurn(False)
        lease = await self.claim_due(account_id=row[0])
        return await self.complete_due(lease) if lease is not None else StatusTurn(False)

v2 input/cursor lifecycle reusing C's pure projector, checkpoint and event transaction.

Ancestors

Class variables

var ledger : PgReportingProjectionStore

Instance variables

prop policy : dict[str, Any]
Expand source code
@property
def policy(self) -> dict[str, Any]:
    return {
        "version": 2,
        "escalation": escalation_identity(self.escalation),
        "consumer_status_enabled": self.consumer_status_enabled,
        "ownership_enabled": self.revision_ownership,
        "notifications_enabled": self.ledger._notifications_enabled,
    }

Methods

async def activate(self, *, account_id: str) ‑> bool
Expand source code
async def activate(self, *, account_id: str) -> bool:
    """Fence, replay captured private history, then activate the frozen baseline.

    Each retained input advances in its own transaction. Cancellation or a
    process crash resumes at that exact cursor. The derived historical
    events stay in an immutable epoch-zero journal, never a delivery queue.
    Old C projectors/sweepers must first be stopped and drained.
    """
    started = await self._begin_activation(account_id=account_id)
    while True:
        async with self._transaction(account_id) as connection:
            state = await self._state_on(connection, account_id)
            if state[4] is not None:
                return started
            await self._activation_step_on(connection, account_id)

Fence, replay captured private history, then activate the frozen baseline.

Each retained input advances in its own transaction. Cancellation or a process crash resumes at that exact cursor. The derived historical events stay in an immutable epoch-zero journal, never a delivery queue. Old C projectors/sweepers must first be stopped and drained.

async def baseline(self, *, account_id: str) ‑> bool
Expand source code
async def baseline(self, *, account_id: str) -> bool:
    return await self.activate(account_id=account_id)
async def baseline_ready(self, *, account_id: str) ‑> bool
Expand source code
async def baseline_ready(self, *, account_id: str) -> bool:
    async with self._transaction(account_id) as connection:
        await validate_projection_schema(
            connection, notifications=self.ledger._notifications_enabled
        )
        row = await (
            await connection.execute(
                "SELECT policy,current_input IS NOT NULL FROM reporting_projection_accounts"
                " WHERE account_id=%s",
                (account_id,),
            )
        ).fetchone()
        if row is None:
            return False
        if row[0] != self.policy:
            raise ReportingNotificationError("status_projection_policy_conflict")
        return bool(row[1])
async def complete_due(self, lease: StatusDueLease) ‑> StatusTurn
Expand source code
async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
    async with self._transaction(lease.scope.account_id) as connection:
        await self._state_on(connection, lease.scope.account_id)
        if not await self._held_on(connection, lease):
            return StatusTurn(False)
        result = await self._project_version_on(connection, lease.scope.account_id)
        if not await self._ack_on(connection, lease):
            raise ReportingNotificationError("status_lease_lost")
        return result
async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    await self.ledger.create_schema()
async def project_one(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def project_one(self, *, account_id: str) -> StatusTurn:
    async with self._transaction(account_id) as connection:
        return await self._project_version_on(connection, account_id)
async def sweep_one(self) ‑> StatusTurn
Expand source code
async def sweep_one(self) -> StatusTurn:
    async with self.ledger._connection() as connection:
        at = await database_now(connection, self.ledger._clock)
        row = await (
            await connection.execute(
                "SELECT c.account_id FROM reporting_status_scope_checkpoints c"
                " JOIN reporting_projection_accounts a ON a.account_id=c.account_id"
                " WHERE c.next_due_at<=%s AND a.policy=%s::jsonb"
                " AND (c.lease_expires_at IS NULL OR c.lease_expires_at<=%s)"
                " ORDER BY c.next_due_at,c.account_id LIMIT 1",
                (at, json.dumps(self.policy), at),
            )
        ).fetchone()
    if row is None:
        return StatusTurn(False)
    lease = await self.claim_due(account_id=row[0])
    return await self.complete_due(lease) if lease is not None else StatusTurn(False)

Inherited members

class ReportingProjectionInput (core: ReportingStatusSnapshot,
reconciliation: tuple[ReportingDeliveryRecord, ...],
document: bytes)
Expand source code
@dataclass(frozen=True)
class ReportingProjectionInput:
    core: ReportingStatusSnapshot = field(repr=False)
    reconciliation: tuple[ReportingDeliveryRecord, ...] = field(repr=False)
    document: bytes = field(repr=False)

    def __post_init__(self) -> None:
        # Structural validation also protects direct internal construction:
        # frozen outer dataclasses alone do not make nested lists/dicts safe.
        _require_frozen(self.core)
        _require_frozen(self.reconciliation)
        if type(self.document) is not bytes:
            raise ReportingNotificationError("status_projection_history_corrupt")

    def __deepcopy__(self, memo: dict[int, Any]) -> ReportingProjectionInput:
        # Decoding closes every collection into tuples of frozen records. A
        # rollback copies the mutable account cursor, queues and input list;
        # the immutable historical values can safely remain shared. Copying
        # every earlier full-account snapshot on each new source mutation made
        # bounded production catch-up grow cubically in retained history.
        memo[id(self)] = self
        return self

    def at(self, at: datetime) -> ReportingProjectionInput:
        """Evaluate a deadline on the same captured history, never today's records."""
        return replace(self, core=replace(self.core, as_of=aware_utc(at)))

ReportingProjectionInput(core: 'ReportingStatusSnapshot', reconciliation: 'tuple[ReportingDeliveryRecord, …]', document: 'bytes')

Instance variables

var core : ReportingStatusSnapshot
var document : bytes
var reconciliation : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...]

Methods

def at(self, at: datetime) ‑> ReportingProjectionInput
Expand source code
def at(self, at: datetime) -> ReportingProjectionInput:
    """Evaluate a deadline on the same captured history, never today's records."""
    return replace(self, core=replace(self.core, as_of=aware_utc(at)))

Evaluate a deadline on the same captured history, never today's records.