Module adcp.reporting.outbox.status

Optional durable status lifecycle API and deterministic worker turns.

Importable without PostgreSQL. Stores own the transaction; neither the projector nor the sweeper performs HTTP or supplies a second persistence seam.

Functions

def advance_checkpoint(previous: StatusCheckpoint | None,
result: StatusProjectionResult,
*,
fired_at: datetime,
source_sequence: int,
baseline: bool = False) ‑> tuple[StatusCheckpoint, ReportingDomainEvent | None]
Expand source code
def advance_checkpoint(
    previous: StatusCheckpoint | None,
    result: StatusProjectionResult,
    *,
    fired_at: datetime,
    source_sequence: int,
    baseline: bool = False,
) -> tuple[StatusCheckpoint, ReportingDomainEvent | None]:
    """Shared mutation/sweep primitive. Caller already owns the typed scope lock."""
    if result.intents:
        raise ReportingNotificationError("status_lifecycle_pending")
    generation = previous.generation if previous is not None else 0
    event = None
    if (
        not baseline
        and result.publishable
        and (previous is None or previous.fingerprint != result.fingerprint)
    ):
        generation += 1
        event = ReportingDomainEvent(
            result.scope.account_id,
            str(uuid4()),
            fired_at,
            StatusChanged(
                scope=result.scope,
                health=result.health,
                fingerprint=result.fingerprint,
                checkpoint_generation=generation,
                previous_health=previous.snapshot["health"] if previous else None,
                issue_ids=result.issue_ids,
            ),
            cause_generation=generation,
        )
        event.body(subscriber_id="validation", idempotency_key="0" * 32)
    return (
        StatusCheckpoint(
            scope=result.scope,
            fingerprint=(
                result.fingerprint
                if result.publishable or previous is None
                else previous.fingerprint
            ),
            generation=generation,
            snapshot=result.canonical(),
            next_due_at=result.next_due_at,
            source_sequence=source_sequence,
            baseline=baseline or bool(previous and previous.baseline),
            publishable=result.publishable,
            lease_token=previous.lease_token if previous else None,
            lease_expires_at=previous.lease_expires_at if previous else None,
        ),
        event,
    )

Shared mutation/sweep primitive. Caller already owns the typed scope lock.

def escalation_identity(escalation: ReportingDeliveryEscalation | None) ‑> dict[str, typing.Any]
Expand source code
def escalation_identity(escalation: ReportingDeliveryEscalation | None) -> dict[str, Any]:
    return {
        **(escalation.to_wire() if escalation else {}),
        "selector_semantics_version": REPORTING_SELECTOR_VERSION,
    }
def settled_replay(snapshot: ReportingStatusSnapshot) ‑> ReportingStatusSnapshot
Expand source code
def settled_replay(snapshot: ReportingStatusSnapshot) -> ReportingStatusSnapshot:
    """Replay a captured committed boundary without consulting final current rows."""
    for _ in range(3):
        intents = lifecycle_intents(snapshot)
        if not intents:
            return snapshot
        snapshot = apply_intents_to_snapshot(snapshot, intents)
    raise ReportingNotificationError("status_lifecycle_did_not_converge")

Replay a captured committed boundary without consulting final current rows.

Classes

class ReportingStatusProjector (store: StatusNotificationStore)
Expand source code
@dataclass(frozen=True)
class ReportingStatusProjector:
    store: StatusNotificationStore

    async def run_once(self, *, account_id: str) -> StatusTurn:
        return await self.store.project_one(account_id=account_id)

    async def rebuild_once(self) -> StatusTurn:
        """Reproject one populated old scope's account, without an account list."""
        if not isinstance(self.store, StatusSelectorRebuildStore):
            raise ReportingNotificationError("status_selector_rebuild_unsupported")
        return await self.store.rebuild_one()

ReportingStatusProjector(store: 'StatusNotificationStore')

Instance variables

var store : StatusNotificationStore

Methods

async def rebuild_once(self) ‑> StatusTurn
Expand source code
async def rebuild_once(self) -> StatusTurn:
    """Reproject one populated old scope's account, without an account list."""
    if not isinstance(self.store, StatusSelectorRebuildStore):
        raise ReportingNotificationError("status_selector_rebuild_unsupported")
    return await self.store.rebuild_one()

Reproject one populated old scope's account, without an account list.

async def run_once(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def run_once(self, *, account_id: str) -> StatusTurn:
    return await self.store.project_one(account_id=account_id)
class ReportingStatusSweeper (store: StatusNotificationStore,
lease_seconds: float = 30)
Expand source code
@dataclass(frozen=True)
class ReportingStatusSweeper:
    store: StatusNotificationStore
    lease_seconds: float = 30

    async def run_once(self, *, account_id: str) -> StatusTurn:
        lease = await self.store.claim_due(account_id=account_id, lease_seconds=self.lease_seconds)
        if lease is None:
            return StatusTurn(False)
        return await self.store.complete_due(lease)

ReportingStatusSweeper(store: 'StatusNotificationStore', lease_seconds: 'float' = 30)

Instance variables

var lease_seconds : float
var store : StatusNotificationStore

Methods

async def run_once(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def run_once(self, *, account_id: str) -> StatusTurn:
    lease = await self.store.claim_due(account_id=account_id, lease_seconds=self.lease_seconds)
    if lease is None:
        return StatusTurn(False)
    return await self.store.complete_due(lease)
class StatusBoundary (transaction_id: str,
through: int,
dirty: tuple[ReportingStatusDirty, ...],
snapshot: ReportingStatusSnapshot)
Expand source code
@dataclass(frozen=True)
class StatusBoundary:
    transaction_id: str
    through: int
    dirty: tuple[ReportingStatusDirty, ...]
    snapshot: ReportingStatusSnapshot

StatusBoundary(transaction_id: 'str', through: 'int', dirty: 'tuple[ReportingStatusDirty, …]', snapshot: 'ReportingStatusSnapshot')

Instance variables

var dirty : tuple[ReportingStatusDirty, ...]
var snapshot : ReportingStatusSnapshot
var through : int
var transaction_id : str
class StatusCheckpoint (scope: ReportingStatusScope,
fingerprint: str,
generation: int,
snapshot: dict[str, Any],
next_due_at: datetime | None,
source_sequence: int,
baseline: bool,
publishable: bool,
lease_token: str | None = None,
lease_expires_at: datetime | None = None,
selector_semantics_version: int = 2,
selector_writer_floor: int = 2)
Expand source code
@dataclass(frozen=True)
class StatusCheckpoint:
    scope: ReportingStatusScope
    fingerprint: str
    generation: int
    snapshot: dict[str, Any]
    next_due_at: datetime | None
    source_sequence: int
    baseline: bool
    publishable: bool
    lease_token: str | None = field(default=None, repr=False)
    lease_expires_at: datetime | None = None
    selector_semantics_version: int = REPORTING_SELECTOR_VERSION
    selector_writer_floor: int = REPORTING_SELECTOR_VERSION

StatusCheckpoint(scope: 'ReportingStatusScope', fingerprint: 'str', generation: 'int', snapshot: 'dict[str, Any]', next_due_at: 'datetime | None', source_sequence: 'int', baseline: 'bool', publishable: 'bool', lease_token: 'str | None' = None, lease_expires_at: 'datetime | None' = None, selector_semantics_version: 'int' = 2, selector_writer_floor: 'int' = 2)

Instance variables

var baseline : bool
var fingerprint : str
var generation : int
var lease_expires_at : datetime.datetime | None
var lease_token : str | None
var next_due_at : datetime.datetime | None
var publishable : bool
var scope : ReportingStatusScope
var selector_semantics_version : int
var selector_writer_floor : int
var snapshot : dict[str, typing.Any]
var source_sequence : int
class StatusDueLease (scope: ReportingStatusScope, token: str, expires_at: datetime, due_at: datetime)
Expand source code
@dataclass(frozen=True)
class StatusDueLease:
    scope: ReportingStatusScope
    token: str = field(repr=False)
    expires_at: datetime
    due_at: datetime

StatusDueLease(scope: 'ReportingStatusScope', token: 'str', expires_at: 'datetime', due_at: 'datetime')

Instance variables

var due_at : datetime.datetime
var expires_at : datetime.datetime
var scope : ReportingStatusScope
var token : str
class StatusNotificationStore (*args, **kwargs)
Expand source code
class StatusNotificationStore(Protocol):
    """Opt-in C participant; no methods are added to the ledger/outbox protocols."""

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

    async def baseline(self, *, account_id: str) -> bool: ...

    async def baseline_ready(self, *, account_id: str) -> bool: ...

    async def project_one(self, *, account_id: str) -> StatusTurn: ...

    async def claim_due(
        self, *, account_id: str, lease_seconds: float
    ) -> StatusDueLease | None: ...

    async def complete_due(self, lease: StatusDueLease) -> StatusTurn: ...

    async def release_due(self, lease: StatusDueLease) -> bool: ...

    async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]: ...

Opt-in C participant; no methods are added to the ledger/outbox protocols.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def baseline(self, *, account_id: str) ‑> bool
Expand source code
async def baseline(self, *, account_id: str) -> bool: ...
async def baseline_ready(self, *, account_id: str) ‑> bool
Expand source code
async def baseline_ready(self, *, account_id: str) -> bool: ...
async def checkpoints(self, *, account_id: str) ‑> tuple[StatusCheckpoint, ...]
Expand source code
async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]: ...
async def claim_due(self, *, account_id: str, lease_seconds: float) ‑> StatusDueLease | None
Expand source code
async def claim_due(
    self, *, account_id: str, lease_seconds: float
) -> StatusDueLease | None: ...
async def complete_due(self,
lease: StatusDueLease) ‑> StatusTurn
Expand source code
async def complete_due(self, lease: StatusDueLease) -> StatusTurn: ...
async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None: ...
async def project_one(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def project_one(self, *, account_id: str) -> StatusTurn: ...
async def release_due(self,
lease: StatusDueLease) ‑> bool
Expand source code
async def release_due(self, lease: StatusDueLease) -> bool: ...
class StatusSelectorRebuildStore (*args, **kwargs)
Expand source code
@runtime_checkable
class StatusSelectorRebuildStore(Protocol):
    """Indexed, account-discovering C cutover seam; no materializer work queue."""

    async def rebuild_one(self) -> StatusTurn: ...

Indexed, account-discovering C cutover seam; no materializer work queue.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def rebuild_one(self) ‑> StatusTurn
Expand source code
async def rebuild_one(self) -> StatusTurn: ...
class StatusTurn (did_work: bool, events: int = 0)
Expand source code
@dataclass(frozen=True)
class StatusTurn:
    did_work: bool
    events: int = 0

StatusTurn(did_work: 'bool', events: 'int' = 0)

Instance variables

var did_work : bool
var events : int