Module adcp.reporting.outbox

Optional transactional reporting notifications, activity and status lifecycle.

Sub-modules

adcp.reporting.outbox.activity

Bounded, principal-scoped projection of reporting HTTP reservations …

adcp.reporting.outbox.identity

One identity rule for authenticated activity reads and trusted registrations.

adcp.reporting.outbox.memory

Process-local outbox sharing the ledger's lock and rollback boundary.

adcp.reporting.outbox.models

Persistence contract for independently leased fanout and HTTP work …

adcp.reporting.outbox.pg

Account-scoped PostgreSQL outbox; HTTP never holds a database transaction.

adcp.reporting.outbox.routing

Trusted account registrations and authenticated, secret-free routing columns.

adcp.reporting.outbox.status

Optional durable status lifecycle API and deterministic worker turns …

adcp.reporting.outbox.status_activity_pg

Closed B+C history union; each queue retains its own HTTP attempt participant.

adcp.reporting.outbox.status_memory

Memory semantic reference for the optional status store; never advertises durability.

adcp.reporting.outbox.status_pg

C-only queues and connection-bound PostgreSQL status transactions …

adcp.reporting.outbox.status_schema

C manifest is independent of the byte-identical A/B required-object manifest.

adcp.reporting.outbox.status_service

Optional status lifecycle composition, without importing PostgreSQL extras.

adcp.reporting.outbox.status_support

Explicit optional C scheduling and per-account advertisement proof.

adcp.reporting.outbox.support

Independent reporting and list-accounts activity capability gates.

adcp.reporting.outbox.worker

Stateless turns over separately durable expansion and HTTP delivery phases.

Functions

def resolve_reporting_consumer(*, auth_info: AuthInfo | None, agent: BuyerAgent | None = None) ‑> str
Expand source code
def resolve_reporting_consumer(
    *, auth_info: AuthInfo | None, agent: BuyerAgent | None = None
) -> str:
    """Resolve the authenticated consumer without choosing between disagreements.

    Inputs must come from verified middleware/the trusted registry. Every
    present flat/registry/signed identity must agree. API/OAuth credentials may
    resolve solely through the registry without a flat principal.
    Normalize aliases in the trusted authentication adapter before this boundary.
    Persist this result as ``ReportingNotificationSubscription.principal_id``.
    """
    from adcp.decisioning.registry import HttpSigCredential

    values = [agent.agent_url] if agent is not None else []
    if auth_info is not None:
        values.extend(
            value for value in (auth_info.principal, auth_info.agent_url) if value is not None
        )
        if isinstance(auth_info.credential, HttpSigCredential):
            values.append(auth_info.credential.agent_url)
    principals = {canonical_consumer(value) for value in values}
    if len(principals) > 1:
        raise ReportingNotificationError("activity_identity_conflict")
    if not principals:
        raise ReportingNotificationError("activity_identity_required")
    return principals.pop()

Resolve the authenticated consumer without choosing between disagreements.

Inputs must come from verified middleware/the trusted registry. Every present flat/registry/signed identity must agree. API/OAuth credentials may resolve solely through the registry without a flat principal. Normalize aliases in the trusted authentication adapter before this boundary. Persist this result as ReportingNotificationSubscription.principal_id.

def sanitize_activity_url(value: str) ‑> str
Expand source code
def sanitize_activity_url(value: str) -> str:
    """Keep the origin and public route words; redact every other path segment.

    This deterministic policy covers UUIDs, JWTs, long hex/base64, mixed tokens,
    percent encoding, and short human-chosen secrets without exposing hashes.
    Userinfo, query, and fragment are always removed. Failures contain no input.
    """
    try:
        if type(value) is not str or len(value) > 8192:
            raise ValueError
        parsed = urlsplit(value)
        if parsed.scheme not in {"https", "http"} or not parsed.hostname:
            raise ValueError
        host = parsed.hostname
        if ":" in host:
            host = f"[{host}]"
        if parsed.port is not None:
            host += f":{parsed.port}"
        segments = []
        for part in parsed.path.split("/"):
            decoded = unquote(part)
            segments.append(
                decoded
                if decoded in _PUBLIC_PATH_SEGMENTS or re.fullmatch(r"v[0-9]{1,3}", decoded)
                else "redacted"
            )
        return str(AnyUrl(urlunsplit((parsed.scheme, host, "/".join(segments), "", ""))))
    except (ValueError, TypeError):
        pass
    # Raise outside the parser's except block so even __context__ cannot retain
    # a raw port/URL from urllib or a validation exception.
    raise ReportingNotificationError("invalid_activity_url")

Keep the origin and public route words; redact every other path segment.

This deterministic policy covers UUIDs, JWTs, long hex/base64, mixed tokens, percent encoding, and short human-chosen secrets without exposing hashes. Userinfo, query, and fragment are always removed. Failures contain no input.

def validate_notification_payload(value: object) ‑> None
Expand source code
def validate_notification_payload(value: object) -> None:
    """Offline full-schema gate, additionally excluding extension metadata."""
    if not isinstance(value, dict) or "ext" in value:
        raise ReportingNotificationError("invalid_payload")
    notification_type = value.get("notification_type")
    path = _SCHEMAS.get(notification_type) if isinstance(notification_type, str) else None
    if path is None:
        raise ReportingNotificationError("invalid_payload")

    def scan(item: object) -> None:
        if isinstance(item, str):
            principal_reference(item)
        elif isinstance(item, list):
            for child in item:
                scan(child)
        elif isinstance(item, dict):
            for child in item.values():
                scan(child)

    try:
        scan(value)
    except ValueError:
        raise ReportingNotificationError("invalid_payload") from None
    validator = get_named_validator(path)
    if validator is None or not validator.is_valid(value):
        raise ReportingNotificationError("invalid_payload")

Offline full-schema gate, additionally excluding extension metadata.

Classes

class ActivityOutcome (status: TerminalActivityStatus,
http_status_code: int | None = None,
response_time_ms: int | None = None)
Expand source code
@dataclass(frozen=True)
class ActivityOutcome:
    status: TerminalActivityStatus
    http_status_code: int | None = None
    response_time_ms: int | None = None

    def __post_init__(self) -> None:
        valid = self.status in {"success", "failed", "timeout", "connection_error"}
        if self.status in {"success", "failed"}:
            valid = valid and (
                type(self.http_status_code) is int
                and 100 <= self.http_status_code <= 599
                and (self.status == "success") == (200 <= self.http_status_code < 300)
                and type(self.response_time_ms) is int
                and self.response_time_ms >= 0
            )
        else:
            valid = valid and self.http_status_code is None and self.response_time_ms is None
        if not valid:
            raise ReportingNotificationError("invalid_activity_outcome")

ActivityOutcome(status: 'TerminalActivityStatus', http_status_code: 'int | None' = None, response_time_ms: 'int | None' = None)

Instance variables

var http_status_code : int | None
var response_time_ms : int | None
var status : Literal['success', 'failed', 'timeout', 'connection_error']
class ActivityRequest (url: str, payload_size_bytes: int)
Expand source code
@dataclass(frozen=True)
class ActivityRequest:
    """Only safe, immutable HTTP request diagnostics cross the SQL boundary."""

    url: str
    payload_size_bytes: int

    def __post_init__(self) -> None:
        object.__setattr__(self, "url", sanitize_activity_url(self.url))
        if type(self.payload_size_bytes) is not int or self.payload_size_bytes < 0:
            raise ReportingNotificationError("invalid_activity_request")

Only safe, immutable HTTP request diagnostics cross the SQL boundary.

Instance variables

var payload_size_bytes : int
var url : str
class AdjustmentPublished (reporting_adjustment_id: str,
consumer_id: str,
adjusts_reporting_revision_id: str,
*,
kind: "Literal['adjustment_published']" = 'adjustment_published')
Expand source code
@dataclass(frozen=True, slots=True)
class AdjustmentPublished(_ClosedValue):
    reporting_adjustment_id: str
    consumer_id: str
    adjusts_reporting_revision_id: str
    kind: Literal["adjustment_published"] = field(default="adjustment_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_adjustment_id, maximum=255)
        consumer_reference(self.consumer_id)
        reporting_identifier(self.adjusts_reporting_revision_id, maximum=255)

AdjustmentPublished(reporting_adjustment_id: 'str', consumer_id: 'str', adjusts_reporting_revision_id: 'str', *, kind: "Literal['adjustment_published']" = 'adjustment_published')

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Instance variables

var adjusts_reporting_revision_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class AdjustmentPublished(_ClosedValue):
    reporting_adjustment_id: str
    consumer_id: str
    adjusts_reporting_revision_id: str
    kind: Literal["adjustment_published"] = field(default="adjustment_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_adjustment_id, maximum=255)
        consumer_reference(self.consumer_id)
        reporting_identifier(self.adjusts_reporting_revision_id, maximum=255)
var consumer_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class AdjustmentPublished(_ClosedValue):
    reporting_adjustment_id: str
    consumer_id: str
    adjusts_reporting_revision_id: str
    kind: Literal["adjustment_published"] = field(default="adjustment_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_adjustment_id, maximum=255)
        consumer_reference(self.consumer_id)
        reporting_identifier(self.adjusts_reporting_revision_id, maximum=255)
var kind : Literal['adjustment_published']
Expand source code
@dataclass(frozen=True, slots=True)
class AdjustmentPublished(_ClosedValue):
    reporting_adjustment_id: str
    consumer_id: str
    adjusts_reporting_revision_id: str
    kind: Literal["adjustment_published"] = field(default="adjustment_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_adjustment_id, maximum=255)
        consumer_reference(self.consumer_id)
        reporting_identifier(self.adjusts_reporting_revision_id, maximum=255)
var reporting_adjustment_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class AdjustmentPublished(_ClosedValue):
    reporting_adjustment_id: str
    consumer_id: str
    adjusts_reporting_revision_id: str
    kind: Literal["adjustment_published"] = field(default="adjustment_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_adjustment_id, maximum=255)
        consumer_reference(self.consumer_id)
        reporting_identifier(self.adjusts_reporting_revision_id, maximum=255)
class DeliveryBinding (account_id: str,
delivery_id: str,
subscriber_id: str,
principal_id: str,
notification_id: str,
notification_type: str,
emission_generation: int,
idempotency_key: str,
destination_sha256: str,
subscription_fingerprint: str,
signing_scope_id: str | None,
cause_kind: str,
cause_id: str,
cause_generation: int,
consumer_namespace: str,
auth_mode: str,
body_sha256: str,
envelope_version: int,
key_version: str)
Expand source code
@dataclass(frozen=True)
class DeliveryBinding:
    """Only non-secret routing identifiers/hashes may be stored in plaintext.

    Every field participates in authenticated encryption. destination_sha256
    binds the full URL (including query); the URL itself is encrypted.
    """

    account_id: str
    delivery_id: str
    subscriber_id: str
    principal_id: str
    notification_id: str
    notification_type: str
    emission_generation: int
    idempotency_key: str
    destination_sha256: str
    subscription_fingerprint: str
    signing_scope_id: str | None
    cause_kind: str
    cause_id: str
    cause_generation: int
    consumer_namespace: str
    auth_mode: str
    body_sha256: str
    envelope_version: int
    key_version: str

Only non-secret routing identifiers/hashes may be stored in plaintext.

Every field participates in authenticated encryption. destination_sha256 binds the full URL (including query); the URL itself is encrypted.

Instance variables

var account_id : str
var auth_mode : str
var body_sha256 : str
var cause_generation : int
var cause_id : str
var cause_kind : str
var consumer_namespace : str
var delivery_id : str
var destination_sha256 : str
var emission_generation : int
var envelope_version : int
var idempotency_key : str
var key_version : str
var notification_id : str
var notification_type : str
var principal_id : str
var signing_scope_id : str | None
var subscriber_id : str
var subscription_fingerprint : str
class DeliveryLease (delivery: StoredDelivery,
token: str,
expires_at: datetime,
attempt_count: int)
Expand source code
@dataclass(frozen=True)
class DeliveryLease:
    delivery: StoredDelivery
    token: str = field(repr=False)
    expires_at: datetime
    attempt_count: int

DeliveryLease(delivery: 'StoredDelivery', token: 'str', expires_at: 'datetime', attempt_count: 'int')

Instance variables

var attempt_count : int
var delivery : StoredDelivery
var expires_at : datetime.datetime
var token : str
class DeliveryStatus (delivery: StoredDelivery,
state: WorkState,
attempt_count: int,
due_at: datetime,
error_code: ErrorCode | None)
Expand source code
@dataclass(frozen=True)
class DeliveryStatus:
    delivery: StoredDelivery
    state: WorkState
    attempt_count: int
    due_at: datetime
    error_code: ErrorCode | None

DeliveryStatus(delivery: 'StoredDelivery', state: 'WorkState', attempt_count: 'int', due_at: 'datetime', error_code: 'ErrorCode | None')

Instance variables

var attempt_count : int
var delivery : StoredDelivery
var due_at : datetime.datetime
var error_code : Literal['network', 'retryable_http', 'permanent_http', 'signing_unavailable', 'permanent_scope', 'subscription_unavailable', 'subscription_changed', 'invalid_configuration', 'invalid_payload', 'integrity_failure', 'lease_expired'] | None
var state : Literal['pending', 'leased', 'complete', 'suppressed', 'quarantined']
class ExpansionLease (account_id: str,
notification_id: str,
emission_generation: int,
token: str,
expires_at: datetime,
attempt_count: int,
event: dict[str, Any],
consumer_namespace: str = '')
Expand source code
@dataclass(frozen=True)
class ExpansionLease:
    account_id: str
    notification_id: str
    emission_generation: int
    token: str = field(repr=False)
    expires_at: datetime
    attempt_count: int
    event: dict[str, Any] = field(repr=False)
    consumer_namespace: str = ""

ExpansionLease(account_id: 'str', notification_id: 'str', emission_generation: 'int', token: 'str', expires_at: 'datetime', attempt_count: 'int', event: 'dict[str, Any]', consumer_namespace: 'str' = '')

Instance variables

var account_id : str
var attempt_count : int
var consumer_namespace : str
var emission_generation : int
var event : dict[str, typing.Any]
var expires_at : datetime.datetime
var notification_id : str
var token : str
class InMemoryReportingOutbox (store: InMemoryReportingLedgerStore)
Expand source code
class InMemoryReportingOutbox:
    """View over an opted-in ledger. Recreating this view preserves its state."""

    def __init__(self, store: InMemoryReportingLedgerStore) -> None:
        if store._notification_state is None:
            raise ValueError("construct the ledger with notifications=True")
        self._store = store

    @property
    def _state(self) -> NotificationState:
        state = self._store._notification_state
        assert state is not None
        return state

    async def list_events(self, *, account_id: str) -> tuple[ReportingDomainEvent, ...]:
        async with self._store._lock:
            return tuple(
                event for (owner, _, _), event in self._state.events.items() if owner == account_id
            )

    async def claim_expansion(
        self, *, account_id: str, now: datetime, lease_seconds: float
    ) -> ExpansionLease | None:
        now = aware_utc(now)
        async with self._store._lock:
            for (owner, consumer, notification, generation), work in self._state.expansions.items():
                if owner == account_id and work.available(now):
                    token, expires = work.claim(now, lease_seconds)
                    event = self._state.events[(owner, consumer, notification)]
                    return ExpansionLease(
                        owner,
                        notification,
                        generation,
                        token,
                        expires,
                        work.attempts,
                        event_storage(event),
                        event.consumer_namespace,
                    )
        return None

    async def complete_expansion(
        self, lease: ExpansionLease, deliveries: tuple[StoredDelivery, ...], *, now: datetime
    ) -> bool:
        async with self._store._mutation():
            work = self._state.expansions.get(
                (
                    lease.account_id,
                    lease.consumer_namespace,
                    lease.notification_id,
                    lease.emission_generation,
                )
            )
            if work is None or not work.held(lease.token, now):
                return False
            for delivery in deliveries:
                binding = delivery.binding
                if (
                    binding.account_id != lease.account_id
                    or binding.notification_id != lease.notification_id
                    or binding.emission_generation != lease.emission_generation
                    or binding.consumer_namespace != lease.consumer_namespace
                ):
                    raise ReportingNotificationError("invalid_configuration")
                key = (
                    binding.account_id,
                    binding.consumer_namespace,
                    binding.notification_id,
                    binding.emission_generation,
                    binding.subscriber_id,
                )
                if key not in self._state.emissions:
                    self._state.deliveries[
                        (binding.account_id, binding.consumer_namespace, binding.delivery_id)
                    ] = (
                        delivery,
                        _Work(now),
                    )
                    self._state.emissions.add(key)
            _finish(work, now=now, state="complete", error_code=None, retry_at=None)
            return True

    async def finish_expansion(
        self,
        lease: ExpansionLease,
        *,
        now: datetime,
        state: WorkState,
        error_code: ErrorCode | None = None,
        retry_at: datetime | None = None,
    ) -> bool:
        validate_finish(state, error_code)
        async with self._store._lock:
            work = self._state.expansions.get(
                (
                    lease.account_id,
                    lease.consumer_namespace,
                    lease.notification_id,
                    lease.emission_generation,
                )
            )
            if work is None or not work.held(lease.token, now):
                return False
            _finish(work, now=now, state=state, error_code=error_code, retry_at=retry_at)
            return True

    async def reemit(
        self,
        *,
        account_id: str,
        notification_id: str,
        now: datetime,
        consumer_namespace: str | None = None,
    ) -> int:
        async with self._store._lock:
            consumers = [
                c
                for a, c, n in self._state.events
                if a == account_id
                and n == notification_id
                and (consumer_namespace is None or c == consumer_namespace)
            ]
            if len(consumers) != 1:
                raise ReportingNotificationError("event_unavailable")
            consumer = consumers[0]
            generation = 1 + max(
                g
                for a, c, n, g in self._state.expansions
                if (a, c, n) == (account_id, consumer, notification_id)
            )
            self._state.expansions[(account_id, consumer, notification_id, generation)] = _Work(now)
            return generation

    async def claim_delivery(
        self, *, account_id: str, now: datetime, lease_seconds: float
    ) -> DeliveryLease | None:
        now = aware_utc(now)
        async with self._store._lock:
            for (owner, _, _), (delivery, work) in self._state.deliveries.items():
                if owner == account_id and work.available(now):
                    token, expires = work.claim(now, lease_seconds)
                    return DeliveryLease(delivery, token, expires, work.attempts)
        return None

    async def delivery_lease_current(self, lease: DeliveryLease, *, now: datetime) -> bool:
        async with self._store._lock:
            binding = lease.delivery.binding
            item = self._state.deliveries.get(
                (binding.account_id, binding.consumer_namespace, binding.delivery_id)
            )
            return item is not None and item[0] == lease.delivery and item[1].held(lease.token, now)

    async def finish_delivery(
        self,
        lease: DeliveryLease,
        *,
        now: datetime,
        state: WorkState,
        error_code: ErrorCode | None = None,
        retry_at: datetime | None = None,
    ) -> bool:
        validate_finish(state, error_code)
        async with self._store._lock:
            binding = lease.delivery.binding
            item = self._state.deliveries.get(
                (binding.account_id, binding.consumer_namespace, binding.delivery_id)
            )
            if item is None or item[0] != lease.delivery or not item[1].held(lease.token, now):
                return False
            _finish(item[1], now=now, state=state, error_code=error_code, retry_at=retry_at)
            return True

    async def reserve_attempt(
        self, lease: DeliveryLease, *, request: ActivityRequest, now: datetime
    ) -> WebhookAttempt | None:
        async with self._store._lock:
            return self._reserve_attempt_locked(lease, request=request, now=now)

    def _reserve_attempt_locked(
        self, lease: DeliveryLease, *, request: ActivityRequest, now: datetime
    ) -> WebhookAttempt | None:
        binding = lease.delivery.binding
        consumer = canonical_consumer(binding.principal_id)
        request = ActivityRequest(request.url, request.payload_size_bytes)
        key = (
            binding.account_id,
            consumer,
            binding.consumer_namespace,
            binding.subscriber_id,
            binding.idempotency_key,
        )
        item = self._state.deliveries.get(
            (binding.account_id, binding.consumer_namespace, binding.delivery_id)
        )
        if (
            item is None
            or item[0] != lease.delivery
            or item[1].expires_at != lease.expires_at
            or not item[1].held(lease.token, now)
        ):
            return None
        if any(
            row.lease_token == lease.token and row.binding == binding
            for row in self._state.activity.values()
        ):
            return None
        number = self._state.activity_heads.get(key, 0) + 1
        attempt = WebhookAttempt(
            binding, number, lease.token, token_hex(32), aware_utc(now), request
        )
        self._state.activity_heads[key] = number
        self._state.activity[(*key, number)] = attempt
        return attempt

    async def complete_attempt(
        self, attempt: WebhookAttempt, *, outcome: ActivityOutcome, now: datetime
    ) -> bool:
        ActivityOutcome.__post_init__(outcome)
        binding = attempt.binding
        consumer = canonical_consumer(binding.principal_id)
        key = (
            binding.account_id,
            consumer,
            binding.consumer_namespace,
            binding.subscriber_id,
            binding.idempotency_key,
            attempt.attempt,
        )
        async with self._store._lock:
            original = self._state.activity.get(key)
            if (
                original != attempt
                or attempt.outcome is not None
                or aware_utc(now) < attempt.fired_at
            ):
                return False
            self._state.activity[key] = replace(
                attempt, outcome=outcome, completed_at=aware_utc(now)
            )
            return True

    async def list_activity(
        self, *, account_id: str, consumer_id: str, limit: int = 50
    ) -> tuple[WebhookAttempt, ...]:
        consumer_id, limit = canonical_consumer(consumer_id), activity_limit(limit)
        async with self._store._lock:
            return tuple(
                sorted(
                    (
                        row
                        for key, row in self._state.activity.items()
                        if key[0] == account_id and key[1] == consumer_id
                    ),
                    key=lambda row: (
                        row.fired_at,
                        row.binding.notification_id,
                        row.binding.idempotency_key,
                        row.binding.subscriber_id,
                        row.attempt,
                        row.binding.delivery_id,
                    ),
                    reverse=True,
                )[:limit]
            )

    async def purge_activity(
        self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
    ) -> int:
        consumer_id = canonical_consumer(consumer_id)
        cutoff = retention_cutoff(now, retention_days)
        async with self._store._lock:
            keys = [
                key
                for key, row in self._state.activity.items()
                if key[0] == account_id
                and key[1] == consumer_id
                and row.completed_at is not None
                and row.completed_at < cutoff
            ]
            for key in keys:
                del self._state.activity[key]
            return len(keys)

    async def list_deliveries(self, *, account_id: str) -> tuple[DeliveryStatus, ...]:
        async with self._store._lock:
            return tuple(
                DeliveryStatus(delivery, work.state, work.attempts, work.due_at, work.error_code)
                for (owner, _, _), (delivery, work) in self._state.deliveries.items()
                if owner == account_id
            )

    async def mark_status_dirty(
        self, scope: ReportingStatusScope, *, reason: DirtyReason, now: datetime
    ) -> None:
        """Additive sweep seam; this slice schedules no clock-driven work."""
        async with self._store._lock:
            self._state.mark_dirty(scope, reason, now)

    async def read_status_dirty(
        self, *, account_id: str, after: int = 0, limit: int = 100
    ) -> tuple[ReportingStatusDirty, ...]:
        if not 1 <= limit <= 1000 or after < 0:
            raise ValueError("invalid dirty checkpoint window")
        async with self._store._lock:
            return tuple(
                deepcopy(
                    [
                        item
                        for item in self._state.dirty
                        if item.scope.account_id == account_id and item.sequence > after
                    ][:limit]
                )
            )

    async def status_checkpoint(self, *, account_id: str, projector_id: str) -> int:
        async with self._store._lock:
            return self._state.checkpoints.get((account_id, projector_id), 0)

    async def advance_status_checkpoint(
        self, *, account_id: str, projector_id: str, expected: int, through: int
    ) -> bool:
        async with self._store._lock:
            maximum = sum(item.scope.account_id == account_id for item in self._state.dirty)
            key = (account_id, projector_id)
            if not 0 <= expected <= through <= maximum:
                raise ValueError("invalid status checkpoint")
            if self._state.checkpoints.get(key, 0) != expected:
                return False
            self._state.checkpoints[key] = through
            return True

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

Subclasses

Methods

async def advance_status_checkpoint(self, *, account_id: str, projector_id: str, expected: int, through: int) ‑> bool
Expand source code
async def advance_status_checkpoint(
    self, *, account_id: str, projector_id: str, expected: int, through: int
) -> bool:
    async with self._store._lock:
        maximum = sum(item.scope.account_id == account_id for item in self._state.dirty)
        key = (account_id, projector_id)
        if not 0 <= expected <= through <= maximum:
            raise ValueError("invalid status checkpoint")
        if self._state.checkpoints.get(key, 0) != expected:
            return False
        self._state.checkpoints[key] = through
        return True
async def claim_delivery(self, *, account_id: str, now: datetime, lease_seconds: float) ‑> DeliveryLease | None
Expand source code
async def claim_delivery(
    self, *, account_id: str, now: datetime, lease_seconds: float
) -> DeliveryLease | None:
    now = aware_utc(now)
    async with self._store._lock:
        for (owner, _, _), (delivery, work) in self._state.deliveries.items():
            if owner == account_id and work.available(now):
                token, expires = work.claim(now, lease_seconds)
                return DeliveryLease(delivery, token, expires, work.attempts)
    return None
async def claim_expansion(self, *, account_id: str, now: datetime, lease_seconds: float) ‑> ExpansionLease | None
Expand source code
async def claim_expansion(
    self, *, account_id: str, now: datetime, lease_seconds: float
) -> ExpansionLease | None:
    now = aware_utc(now)
    async with self._store._lock:
        for (owner, consumer, notification, generation), work in self._state.expansions.items():
            if owner == account_id and work.available(now):
                token, expires = work.claim(now, lease_seconds)
                event = self._state.events[(owner, consumer, notification)]
                return ExpansionLease(
                    owner,
                    notification,
                    generation,
                    token,
                    expires,
                    work.attempts,
                    event_storage(event),
                    event.consumer_namespace,
                )
    return None
async def complete_attempt(self,
attempt: WebhookAttempt,
*,
outcome: ActivityOutcome,
now: datetime) ‑> bool
Expand source code
async def complete_attempt(
    self, attempt: WebhookAttempt, *, outcome: ActivityOutcome, now: datetime
) -> bool:
    ActivityOutcome.__post_init__(outcome)
    binding = attempt.binding
    consumer = canonical_consumer(binding.principal_id)
    key = (
        binding.account_id,
        consumer,
        binding.consumer_namespace,
        binding.subscriber_id,
        binding.idempotency_key,
        attempt.attempt,
    )
    async with self._store._lock:
        original = self._state.activity.get(key)
        if (
            original != attempt
            or attempt.outcome is not None
            or aware_utc(now) < attempt.fired_at
        ):
            return False
        self._state.activity[key] = replace(
            attempt, outcome=outcome, completed_at=aware_utc(now)
        )
        return True
async def complete_expansion(self,
lease: ExpansionLease,
deliveries: tuple[StoredDelivery, ...],
*,
now: datetime) ‑> bool
Expand source code
async def complete_expansion(
    self, lease: ExpansionLease, deliveries: tuple[StoredDelivery, ...], *, now: datetime
) -> bool:
    async with self._store._mutation():
        work = self._state.expansions.get(
            (
                lease.account_id,
                lease.consumer_namespace,
                lease.notification_id,
                lease.emission_generation,
            )
        )
        if work is None or not work.held(lease.token, now):
            return False
        for delivery in deliveries:
            binding = delivery.binding
            if (
                binding.account_id != lease.account_id
                or binding.notification_id != lease.notification_id
                or binding.emission_generation != lease.emission_generation
                or binding.consumer_namespace != lease.consumer_namespace
            ):
                raise ReportingNotificationError("invalid_configuration")
            key = (
                binding.account_id,
                binding.consumer_namespace,
                binding.notification_id,
                binding.emission_generation,
                binding.subscriber_id,
            )
            if key not in self._state.emissions:
                self._state.deliveries[
                    (binding.account_id, binding.consumer_namespace, binding.delivery_id)
                ] = (
                    delivery,
                    _Work(now),
                )
                self._state.emissions.add(key)
        _finish(work, now=now, state="complete", error_code=None, retry_at=None)
        return True
async def delivery_lease_current(self,
lease: DeliveryLease,
*,
now: datetime) ‑> bool
Expand source code
async def delivery_lease_current(self, lease: DeliveryLease, *, now: datetime) -> bool:
    async with self._store._lock:
        binding = lease.delivery.binding
        item = self._state.deliveries.get(
            (binding.account_id, binding.consumer_namespace, binding.delivery_id)
        )
        return item is not None and item[0] == lease.delivery and item[1].held(lease.token, now)
async def finish_delivery(self,
lease: DeliveryLease,
*,
now: datetime,
state: WorkState,
error_code: ErrorCode | None = None,
retry_at: datetime | None = None) ‑> bool
Expand source code
async def finish_delivery(
    self,
    lease: DeliveryLease,
    *,
    now: datetime,
    state: WorkState,
    error_code: ErrorCode | None = None,
    retry_at: datetime | None = None,
) -> bool:
    validate_finish(state, error_code)
    async with self._store._lock:
        binding = lease.delivery.binding
        item = self._state.deliveries.get(
            (binding.account_id, binding.consumer_namespace, binding.delivery_id)
        )
        if item is None or item[0] != lease.delivery or not item[1].held(lease.token, now):
            return False
        _finish(item[1], now=now, state=state, error_code=error_code, retry_at=retry_at)
        return True
async def finish_expansion(self,
lease: ExpansionLease,
*,
now: datetime,
state: WorkState,
error_code: ErrorCode | None = None,
retry_at: datetime | None = None) ‑> bool
Expand source code
async def finish_expansion(
    self,
    lease: ExpansionLease,
    *,
    now: datetime,
    state: WorkState,
    error_code: ErrorCode | None = None,
    retry_at: datetime | None = None,
) -> bool:
    validate_finish(state, error_code)
    async with self._store._lock:
        work = self._state.expansions.get(
            (
                lease.account_id,
                lease.consumer_namespace,
                lease.notification_id,
                lease.emission_generation,
            )
        )
        if work is None or not work.held(lease.token, now):
            return False
        _finish(work, now=now, state=state, error_code=error_code, retry_at=retry_at)
        return True
async def list_activity(self, *, account_id: str, consumer_id: str, limit: int = 50) ‑> tuple[WebhookAttempt, ...]
Expand source code
async def list_activity(
    self, *, account_id: str, consumer_id: str, limit: int = 50
) -> tuple[WebhookAttempt, ...]:
    consumer_id, limit = canonical_consumer(consumer_id), activity_limit(limit)
    async with self._store._lock:
        return tuple(
            sorted(
                (
                    row
                    for key, row in self._state.activity.items()
                    if key[0] == account_id and key[1] == consumer_id
                ),
                key=lambda row: (
                    row.fired_at,
                    row.binding.notification_id,
                    row.binding.idempotency_key,
                    row.binding.subscriber_id,
                    row.attempt,
                    row.binding.delivery_id,
                ),
                reverse=True,
            )[:limit]
        )
async def list_deliveries(self, *, account_id: str) ‑> tuple[DeliveryStatus, ...]
Expand source code
async def list_deliveries(self, *, account_id: str) -> tuple[DeliveryStatus, ...]:
    async with self._store._lock:
        return tuple(
            DeliveryStatus(delivery, work.state, work.attempts, work.due_at, work.error_code)
            for (owner, _, _), (delivery, work) in self._state.deliveries.items()
            if owner == account_id
        )
async def list_events(self, *, account_id: str) ‑> tuple[ReportingDomainEvent, ...]
Expand source code
async def list_events(self, *, account_id: str) -> tuple[ReportingDomainEvent, ...]:
    async with self._store._lock:
        return tuple(
            event for (owner, _, _), event in self._state.events.items() if owner == account_id
        )
async def mark_status_dirty(self,
scope: ReportingStatusScope,
*,
reason: DirtyReason,
now: datetime) ‑> None
Expand source code
async def mark_status_dirty(
    self, scope: ReportingStatusScope, *, reason: DirtyReason, now: datetime
) -> None:
    """Additive sweep seam; this slice schedules no clock-driven work."""
    async with self._store._lock:
        self._state.mark_dirty(scope, reason, now)

Additive sweep seam; this slice schedules no clock-driven work.

async def purge_activity(self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30) ‑> int
Expand source code
async def purge_activity(
    self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
) -> int:
    consumer_id = canonical_consumer(consumer_id)
    cutoff = retention_cutoff(now, retention_days)
    async with self._store._lock:
        keys = [
            key
            for key, row in self._state.activity.items()
            if key[0] == account_id
            and key[1] == consumer_id
            and row.completed_at is not None
            and row.completed_at < cutoff
        ]
        for key in keys:
            del self._state.activity[key]
        return len(keys)
async def read_status_dirty(self, *, account_id: str, after: int = 0, limit: int = 100) ‑> tuple[ReportingStatusDirty, ...]
Expand source code
async def read_status_dirty(
    self, *, account_id: str, after: int = 0, limit: int = 100
) -> tuple[ReportingStatusDirty, ...]:
    if not 1 <= limit <= 1000 or after < 0:
        raise ValueError("invalid dirty checkpoint window")
    async with self._store._lock:
        return tuple(
            deepcopy(
                [
                    item
                    for item in self._state.dirty
                    if item.scope.account_id == account_id and item.sequence > after
                ][:limit]
            )
        )
async def reemit(self,
*,
account_id: str,
notification_id: str,
now: datetime,
consumer_namespace: str | None = None) ‑> int
Expand source code
async def reemit(
    self,
    *,
    account_id: str,
    notification_id: str,
    now: datetime,
    consumer_namespace: str | None = None,
) -> int:
    async with self._store._lock:
        consumers = [
            c
            for a, c, n in self._state.events
            if a == account_id
            and n == notification_id
            and (consumer_namespace is None or c == consumer_namespace)
        ]
        if len(consumers) != 1:
            raise ReportingNotificationError("event_unavailable")
        consumer = consumers[0]
        generation = 1 + max(
            g
            for a, c, n, g in self._state.expansions
            if (a, c, n) == (account_id, consumer, notification_id)
        )
        self._state.expansions[(account_id, consumer, notification_id, generation)] = _Work(now)
        return generation
async def reserve_attempt(self,
lease: DeliveryLease,
*,
request: ActivityRequest,
now: datetime) ‑> WebhookAttempt | None
Expand source code
async def reserve_attempt(
    self, lease: DeliveryLease, *, request: ActivityRequest, now: datetime
) -> WebhookAttempt | None:
    async with self._store._lock:
        return self._reserve_attempt_locked(lease, request=request, now=now)
async def status_checkpoint(self, *, account_id: str, projector_id: str) ‑> int
Expand source code
async def status_checkpoint(self, *, account_id: str, projector_id: str) -> int:
    async with self._store._lock:
        return self._state.checkpoints.get((account_id, projector_id), 0)
class InMemoryReportingStatusOutbox (store: InMemoryReportingLedgerStore)
Expand source code
class InMemoryReportingStatusOutbox(InMemoryReportingOutbox):
    def __init__(self, store: InMemoryReportingLedgerStore) -> None:
        if store._status_notification_state is None:
            raise ValueError("construct a status projection participant first")
        self._store = store

    @property
    def _state(self) -> NotificationState:
        state: _StatusMemoryState = self._store._status_notification_state
        return state.outbox

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

Ancestors

Subclasses

Inherited members

class InMemoryStatusNotificationStore (ledger: InMemoryReportingLedgerStore,
*,
escalation: ReportingDeliveryEscalation | None = None)
Expand source code
class InMemoryStatusNotificationStore:
    def __init__(
        self,
        ledger: InMemoryReportingLedgerStore,
        *,
        escalation: ReportingDeliveryEscalation | None = None,
    ) -> None:
        if ledger._notification_state is None:
            raise ValueError("construct the ledger with notifications=True")
        self.ledger, self.escalation = ledger, escalation
        if ledger._status_notification_state is None:
            ledger._status_notification_state = _StatusMemoryState()
        state = ledger._status_notification_state
        if not hasattr(state, "selector_accounts"):
            state.selector_accounts = {}
        for key, checkpoint in tuple(state.checkpoints.items()):
            # An old shared-state image has no epoch fields. Decode by keyword
            # rather than letting new dataclass class defaults label it v2.
            if "selector_semantics_version" not in vars(checkpoint):
                state.checkpoints[key] = StatusCheckpoint(
                    scope=checkpoint.scope,
                    fingerprint=checkpoint.fingerprint,
                    generation=checkpoint.generation,
                    snapshot=checkpoint.snapshot,
                    next_due_at=checkpoint.next_due_at,
                    source_sequence=checkpoint.source_sequence,
                    baseline=checkpoint.baseline,
                    publishable=checkpoint.publishable,
                    lease_token=checkpoint.lease_token,
                    lease_expires_at=checkpoint.lease_expires_at,
                    selector_semantics_version=1,
                    selector_writer_floor=1,
                )
        self.outbox = InMemoryReportingStatusOutbox(ledger)

    @property
    def _state(self) -> _StatusMemoryState:
        state: _StatusMemoryState = self.ledger._status_notification_state
        return state

    async def create_schema(self) -> None:
        pass

    def _cursor(self, account_id: str) -> int:
        self._projection_fence(account_id)
        account = self._state.accounts.get(account_id)
        if account is None:
            raise ReportingNotificationError("status_baseline_required")
        if account[1] != escalation_identity(self.escalation):
            raise ReportingNotificationError("status_policy_conflict")
        return account[0]

    def _projection_fence(self, account_id: str) -> None:
        if (
            account_id in getattr(self.ledger, "_projection_accounts", {})
            and getattr(self, "_projection_version", 1) != 2
        ):
            raise ReportingNotificationError("status_projection_writer_fenced")

    async def baseline(self, *, account_id: str) -> bool:
        async with self.ledger._mutation():
            self._projection_fence(account_id)
            if account_id in self._state.accounts:
                if not self._needs_rebuild(account_id):
                    self._cursor(account_id)
                return False
            snapshot = settle_memory_snapshot(self.ledger, account_id)
            assert self.ledger._notification_state is not None
            through = max(
                (
                    d.sequence
                    for d in self.ledger._notification_state.dirty
                    if d.scope.account_id == account_id
                ),
                default=0,
            )
            self._apply(snapshot, through=through, baseline=True)
            self._state.replay[account_id] = snapshot
            self._state.accounts[account_id] = (through, escalation_identity(self.escalation))
            self._state.selector_accounts[account_id] = "complete"
            return True

    async def baseline_ready(self, *, account_id: str) -> bool:
        async with self.ledger._mutation():
            if account_id not in self._state.accounts:
                return False
            if self._needs_rebuild(account_id):
                return False
            self._cursor(account_id)
            return True

    def _needs_rebuild(self, account_id: str) -> bool:
        return account_id in self._state.accounts and (
            self._state.selector_accounts.get(account_id) != "complete"
            or any(
                c.scope.account_id == account_id
                and (
                    c.selector_semantics_version != REPORTING_SELECTOR_VERSION
                    or c.selector_writer_floor != REPORTING_SELECTOR_VERSION
                )
                for c in self._state.checkpoints.values()
            )
        )

    def _rebuild(self, account_id: str) -> StatusTurn:
        self._projection_fence(account_id)
        if not self._needs_rebuild(account_id):
            return StatusTurn(False)
        if self._state.selector_accounts.get(account_id) != "transitioning":
            through, policy = self._state.accounts[account_id]
            expected = escalation_identity(self.escalation)
            if policy not in (
                expected,
                {k: v for k, v in expected.items() if k != "selector_semantics_version"},
            ):
                raise ReportingNotificationError("status_policy_conflict")
            self._state.accounts[account_id] = (through, expected)
            self._state.selector_accounts[account_id] = "transitioning"
            for key, checkpoint in tuple(self._state.checkpoints.items()):
                if checkpoint.scope.account_id == account_id:
                    self._state.checkpoints[key] = replace(
                        checkpoint, selector_writer_floor=REPORTING_SELECTOR_VERSION
                    )
            return StatusTurn(True)
        turn = self._project(account_id)
        if turn.did_work:
            return turn
        snapshot = settle_memory_snapshot(self.ledger, account_id)
        deadlines = [
            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
            and c.next_due_at <= snapshot.as_of
        ]
        if deadlines:
            return StatusTurn(
                True,
                self._apply(
                    replace(snapshot, as_of=min(deadlines)), through=self._cursor(account_id)
                ),
            )
        count = self._apply(snapshot, through=self._cursor(account_id))
        self._state.replay[account_id] = snapshot
        self._state.selector_accounts[account_id] = "complete"
        return StatusTurn(True, count)

    async def rebuild_one(self) -> StatusTurn:
        async with self.ledger._mutation():
            expected = escalation_identity(self.escalation)
            legacy = {k: v for k, v in expected.items() if k != "selector_semantics_version"}
            account_id = next(
                (
                    a
                    for a in sorted(self._state.accounts)
                    if self._needs_rebuild(a) and self._state.accounts[a][1] in (expected, legacy)
                ),
                None,
            )
            return self._rebuild(account_id) if account_id is not None else StatusTurn(False)

    def _apply(
        self,
        snapshot: ReportingStatusSnapshot,
        *,
        through: int,
        baseline: bool = False,
        reconciliation: tuple[ReportingDeliveryRecord, ...] | None = None,
        consumer_status_enabled: bool = True,
        enqueue: bool = True,
    ) -> int:
        self._projection_fence(snapshot.account_id)
        snapshot = settled_replay(snapshot)
        events = 0
        scopes = {s.checkpoint_key: s for s in projection_scopes(snapshot)}
        scopes.update(
            (key, c.scope)
            for key, c in self._state.checkpoints.items()
            if c.scope.account_id == snapshot.account_id
        )
        for _, scope in sorted(scopes.items()):
            result = project_status_scope(
                StatusProjectionInput(
                    snapshot,
                    scope,
                    self.escalation,
                    reconciliation=reconciliation,
                    consumer_status_enabled=consumer_status_enabled,
                )
            )
            checkpoint, event = advance_checkpoint(
                self._state.checkpoints.get(scope.checkpoint_key),
                result,
                fired_at=self.ledger._clock(),
                source_sequence=through,
                baseline=baseline,
            )
            self._state.checkpoints[scope.checkpoint_key] = checkpoint
            if event is not None and enqueue:
                self._enqueue_status(event)
                events += 1
        return events

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

    def _project(self, account_id: str) -> StatusTurn:
        cursor = self._cursor(account_id)
        assert self.ledger._notification_state is not None
        boundary = next(
            (
                b
                for b in self.ledger._notification_state.boundaries
                if b.snapshot.account_id == account_id and b.through > cursor
            ),
            None,
        )
        if boundary is None:
            return StatusTurn(False)
        snapshot = settled_replay(
            with_replay_lifecycles(boundary.snapshot, self._state.replay.get(account_id))
        )
        count = self._apply(snapshot, through=boundary.through)
        self._state.replay[account_id] = snapshot
        self._state.accounts[account_id] = (boundary.through, escalation_identity(self.escalation))
        return StatusTurn(True, count)

    async def project_one(self, *, account_id: str) -> StatusTurn:
        async with self.ledger._mutation():
            rebuilt = self._rebuild(account_id)
            if rebuilt.did_work:
                return rebuilt
            return self._project(account_id)

    async def claim_due(
        self, *, account_id: str, lease_seconds: float = 30
    ) -> StatusDueLease | None:
        if lease_seconds <= 0:
            raise ValueError("lease_seconds must be positive")
        async with self.ledger._lock:
            if self._needs_rebuild(account_id):
                return None
            self._cursor(account_id)
            at = self.ledger._clock()
            for key, checkpoint in sorted(self._state.checkpoints.items()):
                if checkpoint.scope.account_id != account_id or checkpoint.next_due_at is None:
                    continue
                if checkpoint.next_due_at > at or (
                    checkpoint.lease_expires_at is not None and checkpoint.lease_expires_at > at
                ):
                    continue
                lease = StatusDueLease(
                    checkpoint.scope,
                    token_hex(32),
                    at + timedelta(seconds=lease_seconds),
                    checkpoint.next_due_at,
                )
                self._state.checkpoints[key] = replace(
                    checkpoint, lease_token=lease.token, lease_expires_at=lease.expires_at
                )
                return lease
        return None

    async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
        try:
            return await self._complete_due(lease)
        except _ExpiredStatusLeaseError:
            return StatusTurn(False)

    async def _complete_due(self, lease: StatusDueLease) -> StatusTurn:
        async with self.ledger._mutation():
            if self._needs_rebuild(lease.scope.account_id):
                return StatusTurn(False)
            checkpoint = self._state.checkpoints.get(lease.scope.checkpoint_key)
            at = self.ledger._clock()
            if (
                checkpoint is None
                or checkpoint.lease_token != lease.token
                or checkpoint.lease_expires_at != lease.expires_at
                or lease.expires_at <= at
            ):
                return StatusTurn(False)
            # Pending committed source transactions win over a previously claimed deadline.
            count = 0
            repaired = False
            while (turn := self._project(lease.scope.account_id)).did_work:
                count += turn.events
                repaired = True
            snapshot = settle_memory_snapshot(self.ledger, lease.scope.account_id)
            if repaired:
                count += self._apply(snapshot, through=self._cursor(lease.scope.account_id))
            else:
                # Walk distinct semantic deadlines in order, including overdue sweeps.
                for _ in range(1000):
                    due = [
                        c.next_due_at
                        for c in self._state.checkpoints.values()
                        if c.scope.account_id == lease.scope.account_id
                        and c.next_due_at is not None
                        and c.next_due_at <= at
                    ]
                    if not due:
                        break
                    count += self._apply(
                        replace(snapshot, as_of=min(due)),
                        through=self._cursor(lease.scope.account_id),
                    )
                else:
                    raise ReportingNotificationError("status_deadline_limit")
            current = self._state.checkpoints[lease.scope.checkpoint_key]
            if current.lease_token != lease.token or lease.expires_at <= self.ledger._clock():
                raise _ExpiredStatusLeaseError
            self._state.checkpoints[lease.scope.checkpoint_key] = replace(
                current, lease_token=None, lease_expires_at=None
            )
            return StatusTurn(True, count)

    async def release_due(self, lease: StatusDueLease) -> bool:
        async with self.ledger._lock:
            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 False
            self._state.checkpoints[lease.scope.checkpoint_key] = replace(
                checkpoint, lease_token=None, lease_expires_at=None
            )
            return True

    async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]:
        async with self.ledger._lock:
            return tuple(
                c
                for _, c in sorted(self._state.checkpoints.items())
                if c.scope.account_id == account_id
            )

Subclasses

Methods

async def baseline(self, *, account_id: str) ‑> bool
Expand source code
async def baseline(self, *, account_id: str) -> bool:
    async with self.ledger._mutation():
        self._projection_fence(account_id)
        if account_id in self._state.accounts:
            if not self._needs_rebuild(account_id):
                self._cursor(account_id)
            return False
        snapshot = settle_memory_snapshot(self.ledger, account_id)
        assert self.ledger._notification_state is not None
        through = max(
            (
                d.sequence
                for d in self.ledger._notification_state.dirty
                if d.scope.account_id == account_id
            ),
            default=0,
        )
        self._apply(snapshot, through=through, baseline=True)
        self._state.replay[account_id] = snapshot
        self._state.accounts[account_id] = (through, escalation_identity(self.escalation))
        self._state.selector_accounts[account_id] = "complete"
        return True
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._mutation():
        if account_id not in self._state.accounts:
            return False
        if self._needs_rebuild(account_id):
            return False
        self._cursor(account_id)
        return True
async def checkpoints(self, *, account_id: str) ‑> tuple[StatusCheckpoint, ...]
Expand source code
async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]:
    async with self.ledger._lock:
        return tuple(
            c
            for _, c in sorted(self._state.checkpoints.items())
            if c.scope.account_id == account_id
        )
async def claim_due(self, *, account_id: str, lease_seconds: float = 30) ‑> StatusDueLease | None
Expand source code
async def claim_due(
    self, *, account_id: str, lease_seconds: float = 30
) -> StatusDueLease | None:
    if lease_seconds <= 0:
        raise ValueError("lease_seconds must be positive")
    async with self.ledger._lock:
        if self._needs_rebuild(account_id):
            return None
        self._cursor(account_id)
        at = self.ledger._clock()
        for key, checkpoint in sorted(self._state.checkpoints.items()):
            if checkpoint.scope.account_id != account_id or checkpoint.next_due_at is None:
                continue
            if checkpoint.next_due_at > at or (
                checkpoint.lease_expires_at is not None and checkpoint.lease_expires_at > at
            ):
                continue
            lease = StatusDueLease(
                checkpoint.scope,
                token_hex(32),
                at + timedelta(seconds=lease_seconds),
                checkpoint.next_due_at,
            )
            self._state.checkpoints[key] = replace(
                checkpoint, lease_token=lease.token, lease_expires_at=lease.expires_at
            )
            return lease
    return None
async def complete_due(self,
lease: StatusDueLease) ‑> StatusTurn
Expand source code
async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
    try:
        return await self._complete_due(lease)
    except _ExpiredStatusLeaseError:
        return StatusTurn(False)
async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    pass
async def project_one(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def project_one(self, *, account_id: str) -> StatusTurn:
    async with self.ledger._mutation():
        rebuilt = self._rebuild(account_id)
        if rebuilt.did_work:
            return rebuilt
        return self._project(account_id)
async def rebuild_one(self) ‑> StatusTurn
Expand source code
async def rebuild_one(self) -> StatusTurn:
    async with self.ledger._mutation():
        expected = escalation_identity(self.escalation)
        legacy = {k: v for k, v in expected.items() if k != "selector_semantics_version"}
        account_id = next(
            (
                a
                for a in sorted(self._state.accounts)
                if self._needs_rebuild(a) and self._state.accounts[a][1] in (expected, legacy)
            ),
            None,
        )
        return self._rebuild(account_id) if account_id is not None else StatusTurn(False)
async def release_due(self,
lease: StatusDueLease) ‑> bool
Expand source code
async def release_due(self, lease: StatusDueLease) -> bool:
    async with self.ledger._lock:
        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 False
        self._state.checkpoints[lease.scope.checkpoint_key] = replace(
            checkpoint, lease_token=None, lease_expires_at=None
        )
        return True
class MaterializationReady (generation_key: ReportingConfigurationGenerationKey,
consumer_id: str,
reporting_obligation_id: str,
destination_ref: str,
method: "Literal['file_transfer', 'dataset_share', 'warehouse_materialization']",
reporting_revision_id: str,
reporting_materialization_id: str,
readiness: "Literal['available', 'delivered']",
finality: ReportingFinality,
data_through: datetime | None,
feed_purpose: FeedPurpose,
*,
kind: "Literal['materialization_ready']" = 'materialization_ready')
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))

Snapshot of a verified, frozen Managed binding; never a capability string.

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Instance variables

var consumer_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var data_through : datetime.datetime | None
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var destination_ref : str
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var feed_purpose : Literal['pacing', 'analytics', 'billing']
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var finality : Literal['snapshot', 'official']
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var generation_key : ReportingConfigurationGenerationKey
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var kind : Literal['materialization_ready']
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var method : Literal['file_transfer', 'dataset_share', 'warehouse_materialization']
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var readiness : Literal['available', 'delivered']
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var reporting_materialization_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var reporting_obligation_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
var reporting_revision_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class MaterializationReady(_ClosedValue):
    """Snapshot of a verified, frozen Managed binding; never a capability string."""

    generation_key: ReportingConfigurationGenerationKey
    consumer_id: str
    reporting_obligation_id: str
    destination_ref: str
    method: Literal["file_transfer", "dataset_share", "warehouse_materialization"]
    reporting_revision_id: str
    reporting_materialization_id: str
    readiness: Literal["available", "delivered"]
    finality: ReportingFinality
    data_through: datetime | None
    feed_purpose: FeedPurpose
    kind: Literal["materialization_ready"] = field(default="materialization_ready", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        consumer_reference(self.consumer_id)
        if self.consumer_id != self.generation_key.consumer_id:
            raise ReportingNotificationError("invalid_status_scope")
        for value in (
            self.reporting_obligation_id,
            self.reporting_revision_id,
            self.reporting_materialization_id,
        ):
            reporting_identifier(value, maximum=255)
        reporting_identifier(self.generation_key.delivery_config_id, maximum=64)
        if self.data_through is not None:
            object.__setattr__(self, "data_through", aware_utc(self.data_through))
class PgReportingActivityUnionStore (b_outbox: PgReportingOutbox,
status_outbox: PgReportingStatusOutbox)
Expand source code
class PgReportingActivityUnionStore:
    """Globally order both durable histories before applying the requested limit."""

    def __init__(self, b_outbox: PgReportingOutbox, status_outbox: PgReportingStatusOutbox) -> None:
        if not PG_AVAILABLE:
            raise ImportError("PgReportingActivityUnionStore requires PostgreSQL; install adcp[pg]")
        if (
            type(b_outbox) is not PgReportingOutbox
            or type(status_outbox) is not PgReportingStatusOutbox
            or b_outbox._pool is not status_outbox._pool
            or b_outbox._clock is not status_outbox._clock
        ):
            raise ReportingNotificationError("activity_union_requires_shared_pool")
        self.b_outbox, self.status_outbox = b_outbox, status_outbox

    async def list_activity(
        self, *, account_id: str, consumer_id: str, limit: int = 50
    ) -> tuple[WebhookAttempt, ...]:
        consumer_id, limit = canonical_consumer(consumer_id), activity_limit(limit)
        async with self.b_outbox._activity_transaction() as connection:
            rows = await (
                await connection.execute(
                    f"SELECT {_ACTIVITY_COLUMNS} FROM ("  # nosec B608
                    " SELECT *, 0 AS queue_rank FROM reporting_webhook_attempts"
                    " WHERE account_id=%s AND principal_id=%s UNION ALL"
                    " SELECT *, 1 AS queue_rank FROM reporting_status_webhook_attempts"
                    " WHERE account_id=%s AND principal_id=%s) AS history"
                    " ORDER BY fired_at DESC, notification_id DESC, idempotency_key DESC,"
                    " subscriber_id DESC, attempt DESC, delivery_id DESC, queue_rank DESC LIMIT %s",
                    (account_id, consumer_id, account_id, consumer_id, limit),
                )
            ).fetchall()
        return tuple(_attempt(row) for row in rows)

    async def purge_activity(
        self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
    ) -> int:
        consumer_id = canonical_consumer(consumer_id)
        retention_cutoff(now, retention_days)
        async with self.b_outbox._activity_transaction() as connection:
            cutoff = retention_cutoff(
                await database_now(connection, self.b_outbox._clock), retention_days
            )
            count = 0
            for table in ("reporting_webhook_attempts", "reporting_status_webhook_attempts"):
                cursor = await connection.execute(
                    f"DELETE FROM {table} WHERE account_id=%s AND principal_id=%s"  # nosec B608
                    " AND completed_at IS NOT NULL AND completed_at < %s",
                    (account_id, consumer_id, cutoff),
                )
                count += cursor.rowcount
        return count

Globally order both durable histories before applying the requested limit.

Methods

async def list_activity(self, *, account_id: str, consumer_id: str, limit: int = 50) ‑> tuple[WebhookAttempt, ...]
Expand source code
async def list_activity(
    self, *, account_id: str, consumer_id: str, limit: int = 50
) -> tuple[WebhookAttempt, ...]:
    consumer_id, limit = canonical_consumer(consumer_id), activity_limit(limit)
    async with self.b_outbox._activity_transaction() as connection:
        rows = await (
            await connection.execute(
                f"SELECT {_ACTIVITY_COLUMNS} FROM ("  # nosec B608
                " SELECT *, 0 AS queue_rank FROM reporting_webhook_attempts"
                " WHERE account_id=%s AND principal_id=%s UNION ALL"
                " SELECT *, 1 AS queue_rank FROM reporting_status_webhook_attempts"
                " WHERE account_id=%s AND principal_id=%s) AS history"
                " ORDER BY fired_at DESC, notification_id DESC, idempotency_key DESC,"
                " subscriber_id DESC, attempt DESC, delivery_id DESC, queue_rank DESC LIMIT %s",
                (account_id, consumer_id, account_id, consumer_id, limit),
            )
        ).fetchall()
    return tuple(_attempt(row) for row in rows)
async def purge_activity(self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30) ‑> int
Expand source code
async def purge_activity(
    self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
) -> int:
    consumer_id = canonical_consumer(consumer_id)
    retention_cutoff(now, retention_days)
    async with self.b_outbox._activity_transaction() as connection:
        cutoff = retention_cutoff(
            await database_now(connection, self.b_outbox._clock), retention_days
        )
        count = 0
        for table in ("reporting_webhook_attempts", "reporting_status_webhook_attempts"):
            cursor = await connection.execute(
                f"DELETE FROM {table} WHERE account_id=%s AND principal_id=%s"  # nosec B608
                " AND completed_at IS NOT NULL AND completed_at < %s",
                (account_id, consumer_id, cutoff),
            )
            count += cursor.rowcount
    return count
class PgReportingOutbox (*, pool: AsyncConnectionPool, clock: Callable[[], datetime] | None = None)
Expand source code
class PgReportingOutbox(_PgReportingActivity):
    """Caller-owned pool, database time, expiring random tokens, fenced writes.

    ``now`` arguments implement the shared memory protocol. In PostgreSQL they
    never override the database clock; only the explicit constructor ``clock``
    seam does, for deterministic tests. Retry *durations* are applied to DB time.
    Claims are counted for diagnostics, never as an HTTP retry limit. Evidence
    and prepared bindings are retained indefinitely. The additive activity
    participant purges only completed HTTP history beyond its retention floor.
    """

    def __init__(
        self, *, pool: AsyncConnectionPool, clock: Callable[[], datetime] | None = None
    ) -> None:
        # Same actionable hint the sibling PG stores raise, so a base-install
        # adopter is told to add the extra instead of meeting a driver error
        # from inside the first query.
        from adcp.reporting.ledger.pg import _INSTALL_HINT, PG_AVAILABLE

        if not PG_AVAILABLE:
            raise ImportError(_INSTALL_HINT)
        self._pool, self._clock = pool, clock

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

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

    async def list_events(self, *, account_id: str) -> tuple[ReportingDomainEvent, ...]:
        async with self._connection() as conn:
            rows = await (
                await conn.execute(
                    "SELECT snapshot FROM reporting_notification_events WHERE account_id = %s"
                    " ORDER BY fired_at, notification_id",
                    (account_id,),
                )
            ).fetchall()
        events = tuple(decode_event(row[0]) for row in rows)
        if any(event.account_id != account_id for event in events):
            raise ReportingNotificationError("invalid_event")
        return events

    async def claim_expansion(
        self, *, account_id: str, now: datetime, lease_seconds: float
    ) -> ExpansionLease | None:
        if lease_seconds <= 0:
            raise ValueError("lease_seconds must be positive")
        async with self._connection() as conn, conn.transaction():
            at = await database_now(conn, self._clock)
            row = await (
                await conn.execute(
                    "SELECT x.notification_id, x.emission_generation, x.claim_count, e.snapshot,"
                    " x.consumer_namespace"
                    " FROM reporting_notification_expansions x JOIN reporting_notification_events e"
                    " ON e.account_id = x.account_id AND e.notification_id = x.notification_id"
                    " AND e.consumer_namespace = x.consumer_namespace"
                    " WHERE x.account_id = %s AND x.due_at <= %s AND (x.state = 'pending' OR"
                    " (x.state = 'leased' AND x.lease_expires_at <= %s))"
                    " ORDER BY x.due_at, x.notification_id, x.emission_generation"
                    " FOR UPDATE OF x SKIP LOCKED LIMIT 1",
                    (account_id, at, at),
                )
            ).fetchone()
            if row is None:
                return None
            token, expires = token_hex(32), at + timedelta(seconds=lease_seconds)
            await conn.execute(
                "UPDATE reporting_notification_expansions SET state = 'leased', lease_token = %s,"
                " lease_expires_at = %s, claim_count = claim_count + 1"
                " WHERE account_id = %s AND consumer_namespace = %s"
                " AND notification_id = %s AND emission_generation = %s",
                (token, expires, account_id, row[4], row[0], row[1]),
            )
            return ExpansionLease(
                account_id, row[0], row[1], token, expires, row[2] + 1, row[3], row[4]
            )

    async def complete_expansion(
        self, lease: ExpansionLease, deliveries: tuple[StoredDelivery, ...], *, now: datetime
    ) -> bool:
        try:
            async with self._connection() as conn, conn.transaction():
                at = await database_now(conn, self._clock)
                row = await (
                    await conn.execute(
                        "SELECT 1 FROM reporting_notification_expansions WHERE account_id = %s"
                        " AND consumer_namespace = %s"
                        " AND notification_id = %s AND emission_generation = %s"
                        " AND state = 'leased'"
                        " AND lease_token = %s AND lease_expires_at > %s FOR UPDATE",
                        (
                            lease.account_id,
                            lease.consumer_namespace,
                            lease.notification_id,
                            lease.emission_generation,
                            lease.token,
                            at,
                        ),
                    )
                ).fetchone()
                if row is None:
                    return False
                for delivery in deliveries:
                    binding = delivery.binding
                    if (
                        binding.account_id != lease.account_id
                        or binding.notification_id != lease.notification_id
                        or binding.emission_generation != lease.emission_generation
                        or binding.consumer_namespace != lease.consumer_namespace
                    ):
                        raise ReportingNotificationError("invalid_configuration")
                    await self._insert_delivery(conn, delivery, at)
                if not await self._finish_expansion(
                    conn, lease, state="complete", error_code=None, delay=0
                ):
                    # Do not commit N subscriber rows if the final fence failed.
                    raise _LostLeaseError
            return True
        except _LostLeaseError:
            return False

    async def _insert_delivery(self, conn: Any, delivery: StoredDelivery, at: datetime) -> None:
        placeholders = ",".join(["%s"] * (_BINDING_SIZE + 2))
        await conn.execute(
            f"INSERT INTO reporting_notification_deliveries ({_DELIVERY_COLUMNS}, due_at)"  # nosec B608
            f" VALUES ({placeholders}) ON CONFLICT"  # nosec B608
            " (account_id, consumer_namespace, notification_id, emission_generation, subscriber_id)"
            " DO NOTHING",
            (*astuple(delivery.binding), delivery.envelope, at),
        )

    async def _finish_expansion(
        self,
        conn: Any,
        lease: ExpansionLease,
        *,
        state: WorkState,
        error_code: ErrorCode | None,
        delay: float,
    ) -> bool:
        at = await database_now(conn, self._clock)
        cursor = await conn.execute(
            "UPDATE reporting_notification_expansions SET state = %s, error_code = %s,"
            " due_at = %s, lease_token = NULL, lease_expires_at = NULL"
            " WHERE account_id = %s AND notification_id = %s AND emission_generation = %s"
            " AND consumer_namespace = %s"
            " AND state = 'leased' AND lease_token = %s AND lease_expires_at > %s",
            (
                state,
                error_code,
                at + timedelta(seconds=delay),
                lease.account_id,
                lease.notification_id,
                lease.emission_generation,
                lease.consumer_namespace,
                lease.token,
                at,
            ),
        )
        return bool(cursor.rowcount)

    async def finish_expansion(
        self,
        lease: ExpansionLease,
        *,
        now: datetime,
        state: WorkState,
        error_code: ErrorCode | None = None,
        retry_at: datetime | None = None,
    ) -> bool:
        validate_finish(state, error_code)
        delay = max(0.0, (retry_at - now).total_seconds()) if retry_at is not None else 0.0
        async with self._connection() as conn, conn.transaction():
            return await self._finish_expansion(
                conn, lease, state=state, error_code=error_code, delay=delay
            )

    async def reemit(
        self,
        *,
        account_id: str,
        notification_id: str,
        now: datetime,
        consumer_namespace: str | None = None,
    ) -> int:
        async with self._connection() as conn, conn.transaction():
            rows = await (
                await conn.execute(
                    "SELECT consumer_namespace FROM reporting_notification_events"
                    " WHERE account_id = %s"
                    " AND notification_id = %s AND (%s::text IS NULL OR consumer_namespace = %s)"
                    " ORDER BY consumer_namespace FOR UPDATE",
                    (account_id, notification_id, consumer_namespace, consumer_namespace),
                )
            ).fetchall()
            if len(rows) != 1:
                raise ReportingNotificationError("event_unavailable")
            consumer = rows[0][0]
            row = await (
                await conn.execute(
                    "SELECT max(emission_generation) + 1 FROM reporting_notification_expansions"
                    " WHERE account_id = %s AND consumer_namespace = %s AND notification_id = %s",
                    (account_id, consumer, notification_id),
                )
            ).fetchone()
            assert row is not None
            generation = int(row[0])
            await conn.execute(
                "INSERT INTO reporting_notification_expansions"
                " (account_id, consumer_namespace, notification_id, emission_generation, due_at)"
                " VALUES (%s,%s,%s,%s,%s)",
                (
                    account_id,
                    consumer,
                    notification_id,
                    generation,
                    await database_now(conn, self._clock),
                ),
            )
            return generation

    async def claim_delivery(
        self, *, account_id: str, now: datetime, lease_seconds: float
    ) -> DeliveryLease | None:
        if lease_seconds <= 0:
            raise ValueError("lease_seconds must be positive")
        async with self._connection() as conn, conn.transaction():
            at = await database_now(conn, self._clock)
            row = await (
                await conn.execute(
                    f"SELECT {_DELIVERY_COLUMNS}, claim_count"  # nosec B608
                    " FROM reporting_notification_deliveries"
                    " WHERE account_id = %s AND due_at <= %s AND (state = 'pending' OR"
                    " (state = 'leased' AND lease_expires_at <= %s)) ORDER BY due_at, delivery_id"
                    " FOR UPDATE SKIP LOCKED LIMIT 1",
                    (account_id, at, at),
                )
            ).fetchone()
            if row is None:
                return None
            token, expires = token_hex(32), at + timedelta(seconds=lease_seconds)
            await conn.execute(
                "UPDATE reporting_notification_deliveries SET state = 'leased', lease_token = %s,"
                " lease_expires_at = %s, claim_count = claim_count + 1"
                " WHERE account_id = %s AND consumer_namespace = %s AND delivery_id = %s",
                (
                    token,
                    expires,
                    account_id,
                    row[_BINDING_NAMES.index("consumer_namespace")],
                    row[1],
                ),
            )
            return DeliveryLease(_delivery(row), token, expires, row[-1] + 1)

    async def delivery_lease_current(self, lease: DeliveryLease, *, now: datetime) -> bool:
        binding = lease.delivery.binding
        async with self._connection() as conn:
            at = await database_now(conn, self._clock)
            row = await (
                await conn.execute(
                    "SELECT 1 FROM reporting_notification_deliveries WHERE account_id = %s"
                    " AND consumer_namespace = %s AND principal_id = %s"
                    " AND delivery_id = %s AND state = 'leased' AND lease_token = %s"
                    " AND lease_expires_at > %s",
                    (
                        binding.account_id,
                        binding.consumer_namespace,
                        binding.principal_id,
                        binding.delivery_id,
                        lease.token,
                        at,
                    ),
                )
            ).fetchone()
        return row is not None

    async def finish_delivery(
        self,
        lease: DeliveryLease,
        *,
        now: datetime,
        state: WorkState,
        error_code: ErrorCode | None = None,
        retry_at: datetime | None = None,
    ) -> bool:
        validate_finish(state, error_code)
        delay = max(0.0, (retry_at - now).total_seconds()) if retry_at is not None else 0.0
        binding = lease.delivery.binding
        async with self._connection() as conn, conn.transaction():
            at = await database_now(conn, self._clock)
            cursor = await conn.execute(
                "UPDATE reporting_notification_deliveries SET state = %s, error_code = %s,"
                " due_at = %s, lease_token = NULL, lease_expires_at = NULL"
                " WHERE account_id = %s AND delivery_id = %s AND state = 'leased'"
                " AND consumer_namespace = %s AND principal_id = %s"
                " AND lease_token = %s AND lease_expires_at > %s",
                (
                    state,
                    error_code,
                    at + timedelta(seconds=delay),
                    binding.account_id,
                    binding.delivery_id,
                    binding.consumer_namespace,
                    binding.principal_id,
                    lease.token,
                    at,
                ),
            )
            return bool(cursor.rowcount)

    async def list_deliveries(self, *, account_id: str) -> tuple[DeliveryStatus, ...]:
        async with self._connection() as conn:
            rows = await (
                await conn.execute(
                    f"SELECT {_DELIVERY_COLUMNS}, state, claim_count, due_at, error_code"  # nosec B608
                    " FROM reporting_notification_deliveries WHERE account_id = %s"
                    " ORDER BY due_at, delivery_id",
                    (account_id,),
                )
            ).fetchall()
        return tuple(DeliveryStatus(_delivery(row), *row[-4:]) for row in rows)

    async def mark_status_dirty(
        self, scope: ReportingStatusScope, *, reason: DirtyReason, now: datetime
    ) -> None:
        """Additive clock-sweep handoff, not a scheduler or a status projector."""
        async with self._connection() as conn, conn.transaction():
            await mark_dirty(conn, scope, reason, await database_now(conn, self._clock))

    async def read_status_dirty(
        self, *, account_id: str, after: int = 0, limit: int = 100
    ) -> tuple[ReportingStatusDirty, ...]:
        if not 1 <= limit <= 1000 or after < 0:
            raise ValueError("invalid dirty checkpoint window")
        async with self._connection() as conn:
            rows = await (
                await conn.execute(
                    "SELECT snapshot FROM reporting_status_dirty WHERE account_id = %s"
                    " AND sequence > %s ORDER BY sequence LIMIT %s",
                    (account_id, after, limit),
                )
            ).fetchall()
        records = tuple(decode_dirty(row[0]) for row in rows)
        if any(record.scope.account_id != account_id for record in records):
            raise ReportingNotificationError("invalid_status_evidence")
        return records

    async def status_checkpoint(self, *, account_id: str, projector_id: str) -> int:
        async with self._connection() as conn:
            row = await (
                await conn.execute(
                    "SELECT sequence FROM reporting_status_checkpoints WHERE account_id = %s"
                    " AND projector_id = %s",
                    (account_id, projector_id),
                )
            ).fetchone()
        return int(row[0]) if row is not None else 0

    async def advance_status_checkpoint(
        self, *, account_id: str, projector_id: str, expected: int, through: int
    ) -> bool:
        async with self._connection() as conn, conn.transaction():
            row = await (
                await conn.execute(
                    "SELECT max_sequence FROM reporting_status_dirty_heads WHERE account_id = %s",
                    (account_id,),
                )
            ).fetchone()
            maximum = int(row[0]) if row is not None else 0
            if not 0 <= expected <= through <= maximum:
                raise ValueError("invalid status checkpoint")
            await conn.execute(
                "INSERT INTO reporting_status_checkpoints (account_id, projector_id, sequence)"
                " VALUES (%s,%s,0) ON CONFLICT (account_id, projector_id) DO NOTHING",
                (account_id, projector_id),
            )
            cursor = await conn.execute(
                "UPDATE reporting_status_checkpoints SET sequence = %s WHERE account_id = %s"
                " AND projector_id = %s AND sequence = %s",
                (through, account_id, projector_id, expected),
            )
            return bool(cursor.rowcount)

Caller-owned pool, database time, expiring random tokens, fenced writes.

now arguments implement the shared memory protocol. In PostgreSQL they never override the database clock; only the explicit constructor clock seam does, for deterministic tests. Retry durations are applied to DB time. Claims are counted for diagnostics, never as an HTTP retry limit. Evidence and prepared bindings are retained indefinitely. The additive activity participant purges only completed HTTP history beyond its retention floor.

Ancestors

  • adcp.reporting.outbox._activity_pg._PgReportingActivity

Subclasses

Methods

async def advance_status_checkpoint(self, *, account_id: str, projector_id: str, expected: int, through: int) ‑> bool
Expand source code
async def advance_status_checkpoint(
    self, *, account_id: str, projector_id: str, expected: int, through: int
) -> bool:
    async with self._connection() as conn, conn.transaction():
        row = await (
            await conn.execute(
                "SELECT max_sequence FROM reporting_status_dirty_heads WHERE account_id = %s",
                (account_id,),
            )
        ).fetchone()
        maximum = int(row[0]) if row is not None else 0
        if not 0 <= expected <= through <= maximum:
            raise ValueError("invalid status checkpoint")
        await conn.execute(
            "INSERT INTO reporting_status_checkpoints (account_id, projector_id, sequence)"
            " VALUES (%s,%s,0) ON CONFLICT (account_id, projector_id) DO NOTHING",
            (account_id, projector_id),
        )
        cursor = await conn.execute(
            "UPDATE reporting_status_checkpoints SET sequence = %s WHERE account_id = %s"
            " AND projector_id = %s AND sequence = %s",
            (through, account_id, projector_id, expected),
        )
        return bool(cursor.rowcount)
async def claim_delivery(self, *, account_id: str, now: datetime, lease_seconds: float) ‑> DeliveryLease | None
Expand source code
async def claim_delivery(
    self, *, account_id: str, now: datetime, lease_seconds: float
) -> DeliveryLease | None:
    if lease_seconds <= 0:
        raise ValueError("lease_seconds must be positive")
    async with self._connection() as conn, conn.transaction():
        at = await database_now(conn, self._clock)
        row = await (
            await conn.execute(
                f"SELECT {_DELIVERY_COLUMNS}, claim_count"  # nosec B608
                " FROM reporting_notification_deliveries"
                " WHERE account_id = %s AND due_at <= %s AND (state = 'pending' OR"
                " (state = 'leased' AND lease_expires_at <= %s)) ORDER BY due_at, delivery_id"
                " FOR UPDATE SKIP LOCKED LIMIT 1",
                (account_id, at, at),
            )
        ).fetchone()
        if row is None:
            return None
        token, expires = token_hex(32), at + timedelta(seconds=lease_seconds)
        await conn.execute(
            "UPDATE reporting_notification_deliveries SET state = 'leased', lease_token = %s,"
            " lease_expires_at = %s, claim_count = claim_count + 1"
            " WHERE account_id = %s AND consumer_namespace = %s AND delivery_id = %s",
            (
                token,
                expires,
                account_id,
                row[_BINDING_NAMES.index("consumer_namespace")],
                row[1],
            ),
        )
        return DeliveryLease(_delivery(row), token, expires, row[-1] + 1)
async def claim_expansion(self, *, account_id: str, now: datetime, lease_seconds: float) ‑> ExpansionLease | None
Expand source code
async def claim_expansion(
    self, *, account_id: str, now: datetime, lease_seconds: float
) -> ExpansionLease | None:
    if lease_seconds <= 0:
        raise ValueError("lease_seconds must be positive")
    async with self._connection() as conn, conn.transaction():
        at = await database_now(conn, self._clock)
        row = await (
            await conn.execute(
                "SELECT x.notification_id, x.emission_generation, x.claim_count, e.snapshot,"
                " x.consumer_namespace"
                " FROM reporting_notification_expansions x JOIN reporting_notification_events e"
                " ON e.account_id = x.account_id AND e.notification_id = x.notification_id"
                " AND e.consumer_namespace = x.consumer_namespace"
                " WHERE x.account_id = %s AND x.due_at <= %s AND (x.state = 'pending' OR"
                " (x.state = 'leased' AND x.lease_expires_at <= %s))"
                " ORDER BY x.due_at, x.notification_id, x.emission_generation"
                " FOR UPDATE OF x SKIP LOCKED LIMIT 1",
                (account_id, at, at),
            )
        ).fetchone()
        if row is None:
            return None
        token, expires = token_hex(32), at + timedelta(seconds=lease_seconds)
        await conn.execute(
            "UPDATE reporting_notification_expansions SET state = 'leased', lease_token = %s,"
            " lease_expires_at = %s, claim_count = claim_count + 1"
            " WHERE account_id = %s AND consumer_namespace = %s"
            " AND notification_id = %s AND emission_generation = %s",
            (token, expires, account_id, row[4], row[0], row[1]),
        )
        return ExpansionLease(
            account_id, row[0], row[1], token, expires, row[2] + 1, row[3], row[4]
        )
async def complete_expansion(self,
lease: ExpansionLease,
deliveries: tuple[StoredDelivery, ...],
*,
now: datetime) ‑> bool
Expand source code
async def complete_expansion(
    self, lease: ExpansionLease, deliveries: tuple[StoredDelivery, ...], *, now: datetime
) -> bool:
    try:
        async with self._connection() as conn, conn.transaction():
            at = await database_now(conn, self._clock)
            row = await (
                await conn.execute(
                    "SELECT 1 FROM reporting_notification_expansions WHERE account_id = %s"
                    " AND consumer_namespace = %s"
                    " AND notification_id = %s AND emission_generation = %s"
                    " AND state = 'leased'"
                    " AND lease_token = %s AND lease_expires_at > %s FOR UPDATE",
                    (
                        lease.account_id,
                        lease.consumer_namespace,
                        lease.notification_id,
                        lease.emission_generation,
                        lease.token,
                        at,
                    ),
                )
            ).fetchone()
            if row is None:
                return False
            for delivery in deliveries:
                binding = delivery.binding
                if (
                    binding.account_id != lease.account_id
                    or binding.notification_id != lease.notification_id
                    or binding.emission_generation != lease.emission_generation
                    or binding.consumer_namespace != lease.consumer_namespace
                ):
                    raise ReportingNotificationError("invalid_configuration")
                await self._insert_delivery(conn, delivery, at)
            if not await self._finish_expansion(
                conn, lease, state="complete", error_code=None, delay=0
            ):
                # Do not commit N subscriber rows if the final fence failed.
                raise _LostLeaseError
        return True
    except _LostLeaseError:
        return False
async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    from adcp.reporting.ledger.pg import PgReportingLedgerStore

    await PgReportingLedgerStore(pool=self._pool, notifications=True).create_schema()
async def delivery_lease_current(self,
lease: DeliveryLease,
*,
now: datetime) ‑> bool
Expand source code
async def delivery_lease_current(self, lease: DeliveryLease, *, now: datetime) -> bool:
    binding = lease.delivery.binding
    async with self._connection() as conn:
        at = await database_now(conn, self._clock)
        row = await (
            await conn.execute(
                "SELECT 1 FROM reporting_notification_deliveries WHERE account_id = %s"
                " AND consumer_namespace = %s AND principal_id = %s"
                " AND delivery_id = %s AND state = 'leased' AND lease_token = %s"
                " AND lease_expires_at > %s",
                (
                    binding.account_id,
                    binding.consumer_namespace,
                    binding.principal_id,
                    binding.delivery_id,
                    lease.token,
                    at,
                ),
            )
        ).fetchone()
    return row is not None
async def finish_delivery(self,
lease: DeliveryLease,
*,
now: datetime,
state: WorkState,
error_code: ErrorCode | None = None,
retry_at: datetime | None = None) ‑> bool
Expand source code
async def finish_delivery(
    self,
    lease: DeliveryLease,
    *,
    now: datetime,
    state: WorkState,
    error_code: ErrorCode | None = None,
    retry_at: datetime | None = None,
) -> bool:
    validate_finish(state, error_code)
    delay = max(0.0, (retry_at - now).total_seconds()) if retry_at is not None else 0.0
    binding = lease.delivery.binding
    async with self._connection() as conn, conn.transaction():
        at = await database_now(conn, self._clock)
        cursor = await conn.execute(
            "UPDATE reporting_notification_deliveries SET state = %s, error_code = %s,"
            " due_at = %s, lease_token = NULL, lease_expires_at = NULL"
            " WHERE account_id = %s AND delivery_id = %s AND state = 'leased'"
            " AND consumer_namespace = %s AND principal_id = %s"
            " AND lease_token = %s AND lease_expires_at > %s",
            (
                state,
                error_code,
                at + timedelta(seconds=delay),
                binding.account_id,
                binding.delivery_id,
                binding.consumer_namespace,
                binding.principal_id,
                lease.token,
                at,
            ),
        )
        return bool(cursor.rowcount)
async def finish_expansion(self,
lease: ExpansionLease,
*,
now: datetime,
state: WorkState,
error_code: ErrorCode | None = None,
retry_at: datetime | None = None) ‑> bool
Expand source code
async def finish_expansion(
    self,
    lease: ExpansionLease,
    *,
    now: datetime,
    state: WorkState,
    error_code: ErrorCode | None = None,
    retry_at: datetime | None = None,
) -> bool:
    validate_finish(state, error_code)
    delay = max(0.0, (retry_at - now).total_seconds()) if retry_at is not None else 0.0
    async with self._connection() as conn, conn.transaction():
        return await self._finish_expansion(
            conn, lease, state=state, error_code=error_code, delay=delay
        )
async def list_deliveries(self, *, account_id: str) ‑> tuple[DeliveryStatus, ...]
Expand source code
async def list_deliveries(self, *, account_id: str) -> tuple[DeliveryStatus, ...]:
    async with self._connection() as conn:
        rows = await (
            await conn.execute(
                f"SELECT {_DELIVERY_COLUMNS}, state, claim_count, due_at, error_code"  # nosec B608
                " FROM reporting_notification_deliveries WHERE account_id = %s"
                " ORDER BY due_at, delivery_id",
                (account_id,),
            )
        ).fetchall()
    return tuple(DeliveryStatus(_delivery(row), *row[-4:]) for row in rows)
async def list_events(self, *, account_id: str) ‑> tuple[ReportingDomainEvent, ...]
Expand source code
async def list_events(self, *, account_id: str) -> tuple[ReportingDomainEvent, ...]:
    async with self._connection() as conn:
        rows = await (
            await conn.execute(
                "SELECT snapshot FROM reporting_notification_events WHERE account_id = %s"
                " ORDER BY fired_at, notification_id",
                (account_id,),
            )
        ).fetchall()
    events = tuple(decode_event(row[0]) for row in rows)
    if any(event.account_id != account_id for event in events):
        raise ReportingNotificationError("invalid_event")
    return events
async def mark_status_dirty(self,
scope: ReportingStatusScope,
*,
reason: DirtyReason,
now: datetime) ‑> None
Expand source code
async def mark_status_dirty(
    self, scope: ReportingStatusScope, *, reason: DirtyReason, now: datetime
) -> None:
    """Additive clock-sweep handoff, not a scheduler or a status projector."""
    async with self._connection() as conn, conn.transaction():
        await mark_dirty(conn, scope, reason, await database_now(conn, self._clock))

Additive clock-sweep handoff, not a scheduler or a status projector.

async def read_status_dirty(self, *, account_id: str, after: int = 0, limit: int = 100) ‑> tuple[ReportingStatusDirty, ...]
Expand source code
async def read_status_dirty(
    self, *, account_id: str, after: int = 0, limit: int = 100
) -> tuple[ReportingStatusDirty, ...]:
    if not 1 <= limit <= 1000 or after < 0:
        raise ValueError("invalid dirty checkpoint window")
    async with self._connection() as conn:
        rows = await (
            await conn.execute(
                "SELECT snapshot FROM reporting_status_dirty WHERE account_id = %s"
                " AND sequence > %s ORDER BY sequence LIMIT %s",
                (account_id, after, limit),
            )
        ).fetchall()
    records = tuple(decode_dirty(row[0]) for row in rows)
    if any(record.scope.account_id != account_id for record in records):
        raise ReportingNotificationError("invalid_status_evidence")
    return records
async def reemit(self,
*,
account_id: str,
notification_id: str,
now: datetime,
consumer_namespace: str | None = None) ‑> int
Expand source code
async def reemit(
    self,
    *,
    account_id: str,
    notification_id: str,
    now: datetime,
    consumer_namespace: str | None = None,
) -> int:
    async with self._connection() as conn, conn.transaction():
        rows = await (
            await conn.execute(
                "SELECT consumer_namespace FROM reporting_notification_events"
                " WHERE account_id = %s"
                " AND notification_id = %s AND (%s::text IS NULL OR consumer_namespace = %s)"
                " ORDER BY consumer_namespace FOR UPDATE",
                (account_id, notification_id, consumer_namespace, consumer_namespace),
            )
        ).fetchall()
        if len(rows) != 1:
            raise ReportingNotificationError("event_unavailable")
        consumer = rows[0][0]
        row = await (
            await conn.execute(
                "SELECT max(emission_generation) + 1 FROM reporting_notification_expansions"
                " WHERE account_id = %s AND consumer_namespace = %s AND notification_id = %s",
                (account_id, consumer, notification_id),
            )
        ).fetchone()
        assert row is not None
        generation = int(row[0])
        await conn.execute(
            "INSERT INTO reporting_notification_expansions"
            " (account_id, consumer_namespace, notification_id, emission_generation, due_at)"
            " VALUES (%s,%s,%s,%s,%s)",
            (
                account_id,
                consumer,
                notification_id,
                generation,
                await database_now(conn, self._clock),
            ),
        )
        return generation
async def status_checkpoint(self, *, account_id: str, projector_id: str) ‑> int
Expand source code
async def status_checkpoint(self, *, account_id: str, projector_id: str) -> int:
    async with self._connection() as conn:
        row = await (
            await conn.execute(
                "SELECT sequence FROM reporting_status_checkpoints WHERE account_id = %s"
                " AND projector_id = %s",
                (account_id, projector_id),
            )
        ).fetchone()
    return int(row[0]) if row is not None else 0
class PgReportingStatusOutbox (*, pool: AsyncConnectionPool, clock: Callable[[], datetime] | None = None)
Expand source code
class PgReportingStatusOutbox(PgReportingOutbox):
    """Separate C events, expansion, delivery and activity, with generic HTTP logic."""

    @asynccontextmanager
    async def _connection(self) -> AsyncIterator[Any]:
        async with self._pool.connection() as connection:
            yield _StatusQueueConnection(connection)

    async def create_schema(self) -> None:
        await PgStatusNotificationStore(
            PgReportingLedgerStore(pool=self._pool, clock=self._clock, notifications=True)
        ).create_schema()

    async def _next_attempt_on(self, conn: Any, binding: DeliveryBinding) -> int:
        b = binding
        row = await (
            await conn.execute(
                "INSERT INTO reporting_webhook_attempt_heads (account_id, consumer_namespace,"
                " principal_id, subscriber_id, idempotency_key, last_attempt)"
                " VALUES (%s,%s,%s,%s,%s,1) ON CONFLICT"
                " (account_id, consumer_namespace, principal_id, subscriber_id, idempotency_key)"
                " DO UPDATE SET last_attempt=reporting_webhook_attempt_heads.last_attempt+1"
                " WHERE reporting_webhook_attempt_heads.account_id=%s"
                " AND reporting_webhook_attempt_heads.consumer_namespace=%s"
                " AND reporting_webhook_attempt_heads.principal_id=%s RETURNING last_attempt",
                (
                    b.account_id,
                    b.consumer_namespace,
                    b.principal_id,
                    b.subscriber_id,
                    b.idempotency_key,
                    b.account_id,
                    b.consumer_namespace,
                    b.principal_id,
                ),
            )
        ).fetchone()
        assert row is not None
        return int(row[0])

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

Ancestors

Subclasses

Methods

async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    await PgStatusNotificationStore(
        PgReportingLedgerStore(pool=self._pool, clock=self._clock, notifications=True)
    ).create_schema()

Inherited members

class PgStatusNotificationStore (ledger: PgReportingLedgerStore,
*,
escalation: ReportingDeliveryEscalation | None = None)
Expand source code
class PgStatusNotificationStore:
    """Optional C transaction participant over an existing notification-enabled ledger."""

    def __init__(
        self,
        ledger: PgReportingLedgerStore,
        *,
        escalation: ReportingDeliveryEscalation | None = None,
    ) -> None:
        if not PG_AVAILABLE:
            raise ImportError("PgStatusNotificationStore requires PostgreSQL; install adcp[pg]")
        if not isinstance(ledger, PgReportingLedgerStore) or not ledger._notifications_enabled:
            raise ReportingNotificationError("status_chain_unready")
        self.ledger, self.escalation = ledger, escalation
        self.outbox = PgReportingStatusOutbox(pool=ledger._pool, clock=ledger._clock)

    async def create_schema(self) -> None:
        from adcp.reporting.outbox.status_schema import validate_status_schema

        async with self.ledger._pool.connection() as connection, connection.transaction():
            await self.ledger._create_schema_on(connection)
            await connection.execute(
                files("adcp.reporting.ledger")
                .joinpath("reporting_status_notifications.sql")
                .read_text()
            )
            await connection.execute(
                files("adcp.reporting.ledger")
                .joinpath("reporting_status_selector_version.sql")
                .read_text()
            )
            await validate_status_schema(connection)

    @asynccontextmanager
    async def _transaction(self, account_id: str) -> AsyncIterator[Any]:
        async with self.ledger._pool.connection() as connection, connection.transaction():
            await connection.execute(
                "SELECT set_config('adcp.reporting.selector_semantics_version', '2', true)"
            )
            await self.ledger._lock_account(connection, account_id)
            yield connection

    async def _account_on(self, connection: Any, account_id: str) -> int:
        row = await (
            await connection.execute(
                "SELECT dirty_sequence, policy, baseline_complete FROM reporting_status_accounts"
                " WHERE account_id=%s FOR UPDATE",
                (account_id,),
            )
        ).fetchone()
        if row is None or not row[2]:
            raise ReportingNotificationError("status_baseline_required")
        if row[1] != escalation_identity(self.escalation):
            raise ReportingNotificationError("status_policy_conflict")
        return int(row[0])

    async def _lock_scopes_on(self, connection: Any, snapshot: ReportingStatusSnapshot) -> None:
        # Insert placeholders before allocating random IDs. New and existing
        # typed scopes have exactly the same row-lock/notification boundary.
        for scope in projection_scopes(snapshot):
            await connection.execute(
                f"INSERT INTO reporting_status_scope_checkpoints ({_KEY}, scope, fingerprint,"  # nosec B608
                " snapshot, source_sequence, baseline, publishable, selector_writer_floor)"
                ' VALUES (%s,%s,%s,%s,%s,%s,%s::jsonb,%s,\'{"health":"waiting"}\',0,FALSE,FALSE,2)'
                f" ON CONFLICT ({_KEY}) DO NOTHING",  # nosec B608
                (*scope.checkpoint_key, json.dumps(asdict(scope)), "0" * 64),
            )
        await (
            await connection.execute(
                "SELECT 1 FROM reporting_status_scope_checkpoints WHERE account_id=%s"
                f" ORDER BY {_KEY}"  # nosec B608
                " FOR UPDATE",
                (snapshot.account_id,),
            )
        ).fetchall()

    async def _write_on(self, connection: Any, checkpoint: StatusCheckpoint) -> None:
        await connection.execute(
            "UPDATE reporting_status_scope_checkpoints SET scope=%s::jsonb, fingerprint=%s,"
            " generation=%s, snapshot=%s::jsonb, next_due_at=%s, source_sequence=%s,"
            " baseline=%s, publishable=%s, initialized=TRUE, selector_semantics_version=%s,"
            " selector_writer_floor=%s"
            f" WHERE {_WHERE}",  # nosec B608
            (
                json.dumps(asdict(checkpoint.scope)),
                checkpoint.fingerprint,
                checkpoint.generation,
                json.dumps(checkpoint.snapshot),
                checkpoint.next_due_at,
                checkpoint.source_sequence,
                checkpoint.baseline,
                checkpoint.publishable,
                checkpoint.selector_semantics_version,
                checkpoint.selector_writer_floor,
                *checkpoint.scope.checkpoint_key,
            ),
        )

    async def _apply_on(
        self,
        connection: Any,
        snapshot: ReportingStatusSnapshot,
        *,
        through: int,
        baseline: bool = False,
        reconciliation: tuple[ReportingDeliveryRecord, ...] | None = None,
        consumer_status_enabled: bool = True,
        enqueue: bool = True,
    ) -> int:
        await self._lock_scopes_on(connection, snapshot)
        snapshot = settled_replay(snapshot)
        count = 0
        rows = await (
            await connection.execute(
                "SELECT scope FROM reporting_status_scope_checkpoints WHERE account_id=%s"
                " ORDER BY account_id, consumer_namespace, delivery_config_id, version,"
                " scope_kind, obligation_namespace",
                (snapshot.account_id,),
            )
        ).fetchall()
        for scope in (decode_status_scope(row[0]) for row in rows):
            row = await (
                await connection.execute(
                    f"SELECT {_CHECKPOINT} FROM reporting_status_scope_checkpoints WHERE {_WHERE}"  # nosec B608
                    " FOR UPDATE",
                    scope.checkpoint_key,
                )
            ).fetchone()
            result = project_status_scope(
                StatusProjectionInput(
                    snapshot,
                    scope,
                    self.escalation,
                    reconciliation=reconciliation,
                    consumer_status_enabled=consumer_status_enabled,
                )
            )
            checkpoint, event = advance_checkpoint(
                _checkpoint(row),
                result,
                fired_at=await database_now(connection, self.ledger._clock),
                source_sequence=through,
                baseline=baseline,
            )
            await self._write_on(connection, checkpoint)
            if event is not None and enqueue:
                await self._enqueue_status_on(connection, event)
                count += 1
        return count

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

    async def baseline(self, *, account_id: str) -> bool:
        from adcp.reporting.outbox.status_schema import validate_status_schema

        async with self._transaction(account_id) as connection:
            await validate_status_schema(connection)
            row = await (
                await connection.execute(
                    "SELECT baseline_complete FROM reporting_status_accounts WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            if row is not None and row[0]:
                if not await self._needs_rebuild_on(connection, account_id):
                    await self._account_on(connection, account_id)
                return False
            await connection.execute(
                "INSERT INTO reporting_status_accounts (account_id, policy) VALUES (%s,%s::jsonb)"
                " ON CONFLICT (account_id) DO NOTHING",
                (account_id, json.dumps(escalation_identity(self.escalation))),
            )
            snapshot = await settle_snapshot_on(self.ledger, connection, account_id=account_id)
            row = await (
                await connection.execute(
                    "SELECT max_sequence FROM reporting_status_dirty_heads WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            through = int(row[0]) if row is not None else 0
            await self._apply_on(connection, snapshot, through=through, baseline=True)
            await connection.execute(
                "UPDATE reporting_status_accounts SET baseline_complete=TRUE,"
                " baseline_highwater=%s,"
                " dirty_sequence=%s, baseline_at=%s, replay_lifecycles=%s::jsonb"
                ", selector_target_version=2, selector_transition='complete'"
                " WHERE account_id=%s",
                (through, through, snapshot.as_of, _replay_storage(snapshot), account_id),
            )
            return True

    async def baseline_ready(self, *, account_id: str) -> bool:
        from adcp.reporting.outbox.status_schema import validate_status_schema

        async with self._transaction(account_id) as connection:
            await validate_status_schema(connection)
            row = await (
                await connection.execute(
                    "SELECT baseline_complete, policy FROM reporting_status_accounts"
                    " WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            if row is None or not row[0] or await self._needs_rebuild_on(connection, account_id):
                return False
            if row[1] != escalation_identity(self.escalation):
                raise ReportingNotificationError("status_policy_conflict")
            # Readiness describes the durable lifecycle, independently of the
            # account's current business health or waived issue occurrences.
            return True

    async def _needs_rebuild_on(self, connection: Any, account_id: str) -> bool:
        row = await (
            await connection.execute(
                "SELECT selector_target_version <> 2 OR selector_transition <> 'complete'"
                " OR EXISTS(SELECT 1 FROM reporting_status_scope_checkpoints c"
                " WHERE c.account_id=a.account_id AND"
                " (c.selector_semantics_version <> 2 OR c.selector_writer_floor <> 2))"
                " FROM reporting_status_accounts a WHERE account_id=%s AND baseline_complete",
                (account_id,),
            )
        ).fetchone()
        return bool(row and row[0])

    async def _rebuild_on(self, connection: Any, account_id: str) -> StatusTurn:
        if not await self._needs_rebuild_on(connection, account_id):
            return StatusTurn(False)
        row = await (
            await connection.execute(
                "SELECT policy, selector_transition FROM reporting_status_accounts"
                " WHERE account_id=%s FOR UPDATE",
                (account_id,),
            )
        ).fetchone()
        expected = escalation_identity(self.escalation)
        legacy = {k: v for k, v in expected.items() if k != "selector_semantics_version"}
        if row[0] not in (expected, legacy):
            raise ReportingNotificationError("status_policy_conflict")
        if row[1] != "transitioning":
            # Phase one commits a checkpoint-local writer floor. The guard
            # never reads/locks an account, including for old due claimers.
            await (
                await connection.execute(
                    "SELECT 1 FROM reporting_status_scope_checkpoints WHERE account_id=%s"
                    " ORDER BY account_id, consumer_namespace, delivery_config_id, version,"
                    " scope_kind, obligation_namespace FOR UPDATE",
                    (account_id,),
                )
            ).fetchall()
            await connection.execute(
                "UPDATE reporting_status_scope_checkpoints SET selector_writer_floor=2"
                " WHERE account_id=%s AND selector_writer_floor <> 2",
                (account_id,),
            )
            await connection.execute(
                "UPDATE reporting_status_accounts SET selector_target_version=2,"
                " selector_transition='transitioning', policy=%s::jsonb WHERE account_id=%s",
                (json.dumps(expected), account_id),
            )
            return StatusTurn(True)
        # Phase two advances exactly one immutable boundary or deadline per
        # transaction. A crash leaves the durable cursor at the last commit.
        turn = await self._project_on(connection, account_id)
        if turn.did_work:
            return turn
        snapshot = await settle_snapshot_on(self.ledger, connection, account_id=account_id)
        through = await self._account_on(connection, account_id)
        due = await (
            await connection.execute(
                "SELECT min(next_due_at) FROM reporting_status_scope_checkpoints"
                " WHERE account_id=%s AND next_due_at <= %s",
                (account_id, snapshot.as_of),
            )
        ).fetchone()
        if due[0] is not None:
            count = await self._apply_on(
                connection, replace(snapshot, as_of=due[0]), through=through
            )
            return StatusTurn(True, count)
        count = await self._apply_on(connection, snapshot, through=through)
        await connection.execute(
            "UPDATE reporting_status_accounts SET selector_transition='complete',"
            " replay_lifecycles=%s::jsonb WHERE account_id=%s",
            (_replay_storage(snapshot), account_id),
        )
        return StatusTurn(True, count)

    async def rebuild_one(self) -> StatusTurn:
        """Discover incomplete C cutovers by index, without enumerating accounts."""
        expected = escalation_identity(self.escalation)
        policies = (
            json.dumps(expected),
            json.dumps({k: v for k, v in expected.items() if k != "selector_semantics_version"}),
        )
        async with self.ledger._pool.connection() as connection:
            row = await (
                await connection.execute(
                    "SELECT account_id FROM reporting_status_accounts WHERE baseline_complete"
                    " AND (selector_target_version <> 2 OR selector_transition <> 'complete')"
                    " AND policy IN (%s::jsonb,%s::jsonb)"
                    " UNION SELECT c.account_id FROM reporting_status_scope_checkpoints c"
                    " JOIN reporting_status_accounts a ON a.account_id=c.account_id"
                    " WHERE (c.selector_semantics_version <> 2 OR c.selector_writer_floor <> 2)"
                    " AND a.baseline_complete AND a.policy IN (%s::jsonb,%s::jsonb)"
                    " ORDER BY account_id LIMIT 1",
                    (*policies, *policies),
                )
            ).fetchone()
        if row is None:
            return StatusTurn(False)
        async with self._transaction(row[0]) as connection:
            return await self._rebuild_on(connection, row[0])

    async def _project_on(self, connection: Any, account_id: str) -> StatusTurn:
        through = await self._account_on(connection, account_id)
        row = await (
            await connection.execute(
                "SELECT through, input FROM reporting_status_boundaries WHERE account_id=%s"
                " AND through > %s ORDER BY through LIMIT 1",
                (account_id, through),
            )
        ).fetchone()
        if row is None:
            return StatusTurn(False)
        snapshot = snapshot_from_storage(row[1])
        await self._lock_scopes_on(connection, snapshot)
        replay = await (
            await connection.execute(
                "SELECT replay_lifecycles FROM reporting_status_accounts WHERE account_id=%s",
                (account_id,),
            )
        ).fetchone()
        snapshot = _with_replay(snapshot, replay[0])
        await persist_replay_lifecycles_on(self.ledger, connection, snapshot)
        count = await self._apply_on(connection, snapshot, through=row[0])
        await connection.execute(
            "UPDATE reporting_status_accounts SET dirty_sequence=%s, replay_lifecycles=%s::jsonb"
            " WHERE account_id=%s",
            (row[0], _replay_storage(snapshot), account_id),
        )
        return StatusTurn(True, count)

    async def project_one(self, *, account_id: str) -> StatusTurn:
        async with self._transaction(account_id) as connection:
            rebuilt = await self._rebuild_on(connection, account_id)
            if rebuilt.did_work:
                return rebuilt
            return await self._project_on(connection, account_id)

    async def claim_due(
        self, *, account_id: str, lease_seconds: float = 30
    ) -> StatusDueLease | None:
        if lease_seconds <= 0:
            raise ValueError("lease_seconds must be positive")
        async with self._transaction(account_id) as connection:
            if await self._needs_rebuild_on(connection, account_id):
                return None
            await self._account_on(connection, account_id)
            at = await database_now(connection, self.ledger._clock)
            row = await (
                await connection.execute(
                    "SELECT scope, next_due_at FROM reporting_status_scope_checkpoints"
                    " WHERE account_id=%s AND next_due_at <= %s"
                    " AND (lease_expires_at IS NULL OR lease_expires_at <= %s)"
                    f" ORDER BY next_due_at, {_KEY} FOR UPDATE SKIP LOCKED LIMIT 1",  # nosec B608
                    (account_id, at, at),
                )
            ).fetchone()
            if row is None:
                return None
            # The claim clock is read after the nonblocking row lock as well.
            at = await database_now(connection, self.ledger._clock)
            lease = StatusDueLease(
                decode_status_scope(row[0]),
                token_hex(32),
                at + timedelta(seconds=lease_seconds),
                row[1],
            )
            await connection.execute(
                "UPDATE reporting_status_scope_checkpoints SET lease_token=%s, lease_expires_at=%s"
                f" WHERE {_WHERE}",  # nosec B608
                (lease.token, lease.expires_at, *lease.scope.checkpoint_key),
            )
            return lease

    async def _held_on(self, connection: Any, lease: StatusDueLease) -> bool:
        at = await database_now(connection, self.ledger._clock)
        row = await (
            await connection.execute(
                f"SELECT 1 FROM reporting_status_scope_checkpoints WHERE {_WHERE}"  # nosec B608
                " AND lease_token=%s AND lease_expires_at=%s AND lease_expires_at > %s FOR UPDATE",
                (*lease.scope.checkpoint_key, lease.token, lease.expires_at, at),
            )
        ).fetchone()
        return row is not None

    async def _ack_on(self, connection: Any, lease: StatusDueLease) -> bool:
        at = await database_now(connection, self.ledger._clock)
        cursor = await connection.execute(
            "UPDATE reporting_status_scope_checkpoints SET lease_token=NULL, lease_expires_at=NULL"
            f" WHERE {_WHERE} AND lease_token=%s AND lease_expires_at > %s",  # nosec B608
            (*lease.scope.checkpoint_key, lease.token, at),
        )
        return bool(cursor.rowcount)

    async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
        try:
            async with self._transaction(lease.scope.account_id) as connection:
                if await self._needs_rebuild_on(connection, lease.scope.account_id):
                    return StatusTurn(False)
                await self._account_on(connection, lease.scope.account_id)
                if not await self._held_on(connection, lease):
                    return StatusTurn(False)
                count = 0
                repaired = False
                while (turn := await self._project_on(connection, lease.scope.account_id)).did_work:
                    count += turn.events
                    repaired = True
                snapshot = await settle_snapshot_on(
                    self.ledger, connection, account_id=lease.scope.account_id
                )
                through = await self._account_on(connection, lease.scope.account_id)
                if repaired:
                    count += await self._apply_on(connection, snapshot, through=through)
                else:
                    for _ in range(1000):
                        row = await (
                            await connection.execute(
                                "SELECT min(next_due_at) FROM reporting_status_scope_checkpoints"
                                " WHERE account_id=%s AND next_due_at <= %s",
                                (lease.scope.account_id, snapshot.as_of),
                            )
                        ).fetchone()
                        if row is None or row[0] is None:
                            break
                        count += await self._apply_on(
                            connection, replace(snapshot, as_of=row[0]), through=through
                        )
                    else:
                        raise ReportingNotificationError("status_deadline_limit")
                if not await self._ack_on(connection, lease):
                    raise _ExpiredStatusLeaseError
                return StatusTurn(True, count)
        except _ExpiredStatusLeaseError:
            return StatusTurn(False)

    async def release_due(self, lease: StatusDueLease) -> bool:
        async with self._transaction(lease.scope.account_id) as connection:
            return await self._ack_on(connection, lease)

    async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]:
        async with self.ledger._pool.connection() as connection:
            rows = await (
                await connection.execute(
                    f"SELECT {_CHECKPOINT} FROM reporting_status_scope_checkpoints"  # nosec B608
                    f" WHERE account_id=%s ORDER BY {_KEY}",
                    (account_id,),  # nosec B608
                )
            ).fetchall()
        return tuple(c for row in rows if (c := _checkpoint(row)) is not None)

Optional C transaction participant over an existing notification-enabled ledger.

Subclasses

Methods

async def baseline(self, *, account_id: str) ‑> bool
Expand source code
async def baseline(self, *, account_id: str) -> bool:
    from adcp.reporting.outbox.status_schema import validate_status_schema

    async with self._transaction(account_id) as connection:
        await validate_status_schema(connection)
        row = await (
            await connection.execute(
                "SELECT baseline_complete FROM reporting_status_accounts WHERE account_id=%s",
                (account_id,),
            )
        ).fetchone()
        if row is not None and row[0]:
            if not await self._needs_rebuild_on(connection, account_id):
                await self._account_on(connection, account_id)
            return False
        await connection.execute(
            "INSERT INTO reporting_status_accounts (account_id, policy) VALUES (%s,%s::jsonb)"
            " ON CONFLICT (account_id) DO NOTHING",
            (account_id, json.dumps(escalation_identity(self.escalation))),
        )
        snapshot = await settle_snapshot_on(self.ledger, connection, account_id=account_id)
        row = await (
            await connection.execute(
                "SELECT max_sequence FROM reporting_status_dirty_heads WHERE account_id=%s",
                (account_id,),
            )
        ).fetchone()
        through = int(row[0]) if row is not None else 0
        await self._apply_on(connection, snapshot, through=through, baseline=True)
        await connection.execute(
            "UPDATE reporting_status_accounts SET baseline_complete=TRUE,"
            " baseline_highwater=%s,"
            " dirty_sequence=%s, baseline_at=%s, replay_lifecycles=%s::jsonb"
            ", selector_target_version=2, selector_transition='complete'"
            " WHERE account_id=%s",
            (through, through, snapshot.as_of, _replay_storage(snapshot), account_id),
        )
        return True
async def baseline_ready(self, *, account_id: str) ‑> bool
Expand source code
async def baseline_ready(self, *, account_id: str) -> bool:
    from adcp.reporting.outbox.status_schema import validate_status_schema

    async with self._transaction(account_id) as connection:
        await validate_status_schema(connection)
        row = await (
            await connection.execute(
                "SELECT baseline_complete, policy FROM reporting_status_accounts"
                " WHERE account_id=%s",
                (account_id,),
            )
        ).fetchone()
        if row is None or not row[0] or await self._needs_rebuild_on(connection, account_id):
            return False
        if row[1] != escalation_identity(self.escalation):
            raise ReportingNotificationError("status_policy_conflict")
        # Readiness describes the durable lifecycle, independently of the
        # account's current business health or waived issue occurrences.
        return True
async def checkpoints(self, *, account_id: str) ‑> tuple[StatusCheckpoint, ...]
Expand source code
async def checkpoints(self, *, account_id: str) -> tuple[StatusCheckpoint, ...]:
    async with self.ledger._pool.connection() as connection:
        rows = await (
            await connection.execute(
                f"SELECT {_CHECKPOINT} FROM reporting_status_scope_checkpoints"  # nosec B608
                f" WHERE account_id=%s ORDER BY {_KEY}",
                (account_id,),  # nosec B608
            )
        ).fetchall()
    return tuple(c for row in rows if (c := _checkpoint(row)) is not None)
async def claim_due(self, *, account_id: str, lease_seconds: float = 30) ‑> StatusDueLease | None
Expand source code
async def claim_due(
    self, *, account_id: str, lease_seconds: float = 30
) -> StatusDueLease | None:
    if lease_seconds <= 0:
        raise ValueError("lease_seconds must be positive")
    async with self._transaction(account_id) as connection:
        if await self._needs_rebuild_on(connection, account_id):
            return None
        await self._account_on(connection, account_id)
        at = await database_now(connection, self.ledger._clock)
        row = await (
            await connection.execute(
                "SELECT scope, next_due_at FROM reporting_status_scope_checkpoints"
                " WHERE account_id=%s AND next_due_at <= %s"
                " AND (lease_expires_at IS NULL OR lease_expires_at <= %s)"
                f" ORDER BY next_due_at, {_KEY} FOR UPDATE SKIP LOCKED LIMIT 1",  # nosec B608
                (account_id, at, at),
            )
        ).fetchone()
        if row is None:
            return None
        # The claim clock is read after the nonblocking row lock as well.
        at = await database_now(connection, self.ledger._clock)
        lease = StatusDueLease(
            decode_status_scope(row[0]),
            token_hex(32),
            at + timedelta(seconds=lease_seconds),
            row[1],
        )
        await connection.execute(
            "UPDATE reporting_status_scope_checkpoints SET lease_token=%s, lease_expires_at=%s"
            f" WHERE {_WHERE}",  # nosec B608
            (lease.token, lease.expires_at, *lease.scope.checkpoint_key),
        )
        return lease
async def complete_due(self,
lease: StatusDueLease) ‑> StatusTurn
Expand source code
async def complete_due(self, lease: StatusDueLease) -> StatusTurn:
    try:
        async with self._transaction(lease.scope.account_id) as connection:
            if await self._needs_rebuild_on(connection, lease.scope.account_id):
                return StatusTurn(False)
            await self._account_on(connection, lease.scope.account_id)
            if not await self._held_on(connection, lease):
                return StatusTurn(False)
            count = 0
            repaired = False
            while (turn := await self._project_on(connection, lease.scope.account_id)).did_work:
                count += turn.events
                repaired = True
            snapshot = await settle_snapshot_on(
                self.ledger, connection, account_id=lease.scope.account_id
            )
            through = await self._account_on(connection, lease.scope.account_id)
            if repaired:
                count += await self._apply_on(connection, snapshot, through=through)
            else:
                for _ in range(1000):
                    row = await (
                        await connection.execute(
                            "SELECT min(next_due_at) FROM reporting_status_scope_checkpoints"
                            " WHERE account_id=%s AND next_due_at <= %s",
                            (lease.scope.account_id, snapshot.as_of),
                        )
                    ).fetchone()
                    if row is None or row[0] is None:
                        break
                    count += await self._apply_on(
                        connection, replace(snapshot, as_of=row[0]), through=through
                    )
                else:
                    raise ReportingNotificationError("status_deadline_limit")
            if not await self._ack_on(connection, lease):
                raise _ExpiredStatusLeaseError
            return StatusTurn(True, count)
    except _ExpiredStatusLeaseError:
        return StatusTurn(False)
async def create_schema(self) ‑> None
Expand source code
async def create_schema(self) -> None:
    from adcp.reporting.outbox.status_schema import validate_status_schema

    async with self.ledger._pool.connection() as connection, connection.transaction():
        await self.ledger._create_schema_on(connection)
        await connection.execute(
            files("adcp.reporting.ledger")
            .joinpath("reporting_status_notifications.sql")
            .read_text()
        )
        await connection.execute(
            files("adcp.reporting.ledger")
            .joinpath("reporting_status_selector_version.sql")
            .read_text()
        )
        await validate_status_schema(connection)
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:
        rebuilt = await self._rebuild_on(connection, account_id)
        if rebuilt.did_work:
            return rebuilt
        return await self._project_on(connection, account_id)
async def rebuild_one(self) ‑> StatusTurn
Expand source code
async def rebuild_one(self) -> StatusTurn:
    """Discover incomplete C cutovers by index, without enumerating accounts."""
    expected = escalation_identity(self.escalation)
    policies = (
        json.dumps(expected),
        json.dumps({k: v for k, v in expected.items() if k != "selector_semantics_version"}),
    )
    async with self.ledger._pool.connection() as connection:
        row = await (
            await connection.execute(
                "SELECT account_id FROM reporting_status_accounts WHERE baseline_complete"
                " AND (selector_target_version <> 2 OR selector_transition <> 'complete')"
                " AND policy IN (%s::jsonb,%s::jsonb)"
                " UNION SELECT c.account_id FROM reporting_status_scope_checkpoints c"
                " JOIN reporting_status_accounts a ON a.account_id=c.account_id"
                " WHERE (c.selector_semantics_version <> 2 OR c.selector_writer_floor <> 2)"
                " AND a.baseline_complete AND a.policy IN (%s::jsonb,%s::jsonb)"
                " ORDER BY account_id LIMIT 1",
                (*policies, *policies),
            )
        ).fetchone()
    if row is None:
        return StatusTurn(False)
    async with self._transaction(row[0]) as connection:
        return await self._rebuild_on(connection, row[0])

Discover incomplete C cutovers by index, without enumerating accounts.

async def release_due(self,
lease: StatusDueLease) ‑> bool
Expand source code
async def release_due(self, lease: StatusDueLease) -> bool:
    async with self._transaction(lease.scope.account_id) as connection:
        return await self._ack_on(connection, lease)
class ReportingActivityProjector (store: ReportingActivityReader)
Expand source code
class ReportingActivityProjector:
    """Optional decorator mounted *after* AccountStore's visibility/filtering.

    Pass this instance as ``account_activity=`` to the platform server factory.
    It never resolves additional accounts or accepts a consumer from the body.
    Memory stores support conformance but cannot justify durable capabilities.
    """

    def __init__(self, store: ReportingActivityReader) -> None:
        self.store = store

    async def for_account(
        self, *, account_id: str, context: ResolveContext, limit: int = 50
    ) -> list[dict[str, Any]]:
        consumer = resolve_reporting_consumer(auth_info=context.auth_info, agent=context.agent)
        attempts = await self.store.list_activity(
            account_id=account_id, consumer_id=consumer, limit=activity_limit(limit)
        )
        return [attempt.to_wire() for attempt in attempts]

    async def enrich(
        self, envelope: Any, *, context: ResolveContext, include: bool, limit: int = 50
    ) -> Any:
        # Preserve legacy envelopes, pagination, and account ordering. Only
        # already-returned accounts can ever reach the activity read boundary.
        if not isinstance(envelope, dict) or not isinstance(envelope.get("accounts"), list):
            return envelope
        if include:
            activity_limit(limit)
            resolve_reporting_consumer(auth_info=context.auth_info, agent=context.agent)
        accounts = []
        for original in envelope["accounts"]:
            if not isinstance(original, dict):
                accounts.append(original)
                continue
            account = dict(original)
            account.pop("webhook_activity", None)
            if include and isinstance(account.get("account_id"), str):
                account["webhook_activity"] = await self.for_account(
                    account_id=account["account_id"], context=context, limit=limit
                )
            accounts.append(account)
        return {**envelope, "accounts": accounts}

Optional decorator mounted after AccountStore's visibility/filtering.

Pass this instance as account_activity= to the platform server factory. It never resolves additional accounts or accepts a consumer from the body. Memory stores support conformance but cannot justify durable capabilities.

Methods

async def enrich(self, envelope: Any, *, context: ResolveContext, include: bool, limit: int = 50) ‑> Any
Expand source code
async def enrich(
    self, envelope: Any, *, context: ResolveContext, include: bool, limit: int = 50
) -> Any:
    # Preserve legacy envelopes, pagination, and account ordering. Only
    # already-returned accounts can ever reach the activity read boundary.
    if not isinstance(envelope, dict) or not isinstance(envelope.get("accounts"), list):
        return envelope
    if include:
        activity_limit(limit)
        resolve_reporting_consumer(auth_info=context.auth_info, agent=context.agent)
    accounts = []
    for original in envelope["accounts"]:
        if not isinstance(original, dict):
            accounts.append(original)
            continue
        account = dict(original)
        account.pop("webhook_activity", None)
        if include and isinstance(account.get("account_id"), str):
            account["webhook_activity"] = await self.for_account(
                account_id=account["account_id"], context=context, limit=limit
            )
        accounts.append(account)
    return {**envelope, "accounts": accounts}
async def for_account(self, *, account_id: str, context: ResolveContext, limit: int = 50) ‑> list[dict[str, Any]]
Expand source code
async def for_account(
    self, *, account_id: str, context: ResolveContext, limit: int = 50
) -> list[dict[str, Any]]:
    consumer = resolve_reporting_consumer(auth_info=context.auth_info, agent=context.agent)
    attempts = await self.store.list_activity(
        account_id=account_id, consumer_id=consumer, limit=activity_limit(limit)
    )
    return [attempt.to_wire() for attempt in attempts]
class ReportingActivityReader (*args, **kwargs)
Expand source code
class ReportingActivityReader(Protocol):
    """Optional read/purge projection, including the closed B+C activity union."""

    async def list_activity(
        self, *, account_id: str, consumer_id: str, limit: int = 50
    ) -> tuple[WebhookAttempt, ...]: ...

    async def purge_activity(
        self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
    ) -> int: ...

Optional read/purge projection, including the closed B+C activity union.

Ancestors

  • typing.Protocol
  • typing.Generic

Subclasses

Methods

async def list_activity(self, *, account_id: str, consumer_id: str, limit: int = 50) ‑> tuple[WebhookAttempt, ...]
Expand source code
async def list_activity(
    self, *, account_id: str, consumer_id: str, limit: int = 50
) -> tuple[WebhookAttempt, ...]: ...
async def purge_activity(self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30) ‑> int
Expand source code
async def purge_activity(
    self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
) -> int: ...
class ReportingActivityStore (*args, **kwargs)
Expand source code
class ReportingActivityStore(ReportingActivityReader, Protocol):
    async def reserve_attempt(
        self, lease: DeliveryLease, *, request: ActivityRequest, now: datetime
    ) -> WebhookAttempt | None: ...

    async def complete_attempt(
        self, attempt: WebhookAttempt, *, outcome: ActivityOutcome, now: datetime
    ) -> bool: ...

    async def list_activity(
        self, *, account_id: str, consumer_id: str, limit: int = 50
    ) -> tuple[WebhookAttempt, ...]: ...

    async def purge_activity(
        self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
    ) -> int: ...

Optional read/purge projection, including the closed B+C activity union.

Ancestors

Methods

async def complete_attempt(self,
attempt: WebhookAttempt,
*,
outcome: ActivityOutcome,
now: datetime) ‑> bool
Expand source code
async def complete_attempt(
    self, attempt: WebhookAttempt, *, outcome: ActivityOutcome, now: datetime
) -> bool: ...
async def list_activity(self, *, account_id: str, consumer_id: str, limit: int = 50) ‑> tuple[WebhookAttempt, ...]
Expand source code
async def list_activity(
    self, *, account_id: str, consumer_id: str, limit: int = 50
) -> tuple[WebhookAttempt, ...]: ...
async def purge_activity(self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30) ‑> int
Expand source code
async def purge_activity(
    self, *, account_id: str, consumer_id: str, now: datetime, retention_days: int = 30
) -> int: ...
async def reserve_attempt(self,
lease: DeliveryLease,
*,
request: ActivityRequest,
now: datetime) ‑> WebhookAttempt | None
Expand source code
async def reserve_attempt(
    self, lease: DeliveryLease, *, request: ActivityRequest, now: datetime
) -> WebhookAttempt | None: ...
class ReportingActivitySupport (worker: ReportingNotificationWorker,
ledger: ReportingLedgerStore,
projector: ReportingActivityProjector | None = None,
status: ReportingStatusSupport | None = None)
Expand source code
@dataclass(frozen=True)
class ReportingActivitySupport:
    """The concrete writer/store/projector chain scheduled by the adopter.

    Mount this as ``reporting_activity=`` on the platform server factory. Mount
    ``projector`` separately as ``account_activity=`` to enable list-accounts
    enrichment. Neither relationship notification claims nor memory components
    supply evidence of durable reporting activity.
    """

    worker: ReportingNotificationWorker
    ledger: ReportingLedgerStore
    projector: ReportingActivityProjector | None = None
    status: ReportingStatusSupport | None = None
    _schema_validation: _SchemaValidation = field(
        default_factory=_SchemaValidation, init=False, repr=False, compare=False
    )

    def invalidate_schema_validation(self) -> None:
        """Discard schema evidence before controlled DDL, after draining callers.

        An older in-flight scan cannot publish into the new epoch. A fresh
        support/server instance after migration is the normal deployment path;
        serving-time DDL is not automatically detected by discovery.
        """
        state = self._schema_validation
        state.epoch += 1
        state.positive = None
        state.flight = None

    async def _schema_proven(self, pool: Any, *, composite: bool) -> bool:
        from adcp.reporting.outbox import _schema

        assert self.projector is not None
        wiring: tuple[object, ...] = (
            self.worker,
            self.ledger,
            self.projector,
            self.projector.store,
            self.worker.outbox,
            self.worker.activity,
            self.worker.cipher,
            self.worker.subscriptions,
            self.worker.signing,
            pool,
        )
        status_contract = ""
        if composite:
            from adcp.reporting.outbox import status_schema

            assert self.status is not None
            wiring += (
                self.status,
                self.status.store,
                self.status.worker,
                self.status.worker.outbox,
            )
            status_contract = _schema._digest(status_schema.REQUIRED_STATUS_OBJECTS)
        key: _ProofKey = (
            "B+C" if composite else "B",
            tuple(id(component) for component in wiring),
            _schema._digest(_schema.REQUIRED_OBJECTS),
            status_contract,
        )
        state = self._schema_validation
        if state.positive == key:
            return True
        flight = state.flight
        if flight is None:
            flight = _SchemaFlight(key, state.epoch)
            state.flight = flight
            flight.task = asyncio.create_task(self._scan_schema(pool, flight, composite=composite))
            flight.task.add_done_callback(_consume_schema_failure)
        task = flight.task
        if task is None or flight.key != key or task.get_loop() is not asyncio.get_running_loop():
            # Overlapping, differently wired/looped startup is not evidence.
            # Completed startup leaves no task or loop in this instance.
            return False
        if not await asyncio.shield(task):
            return False
        # Wiring may have changed while the catalog checkout was suspended.
        # Rerun the same cheap checks before accepting the completed proof.
        return await self.durable()

    async def _scan_schema(self, pool: Any, flight: _SchemaFlight, *, composite: bool) -> bool:
        from adcp.reporting.outbox import _schema

        state = self._schema_validation
        try:
            async with pool.connection() as connection:
                try:
                    installed = await _schema.schema_objects(connection)
                except Exception:
                    raise ReportingNotificationError(
                        "notification_schema_unready:catalog_unavailable"
                    ) from None
                # One physical catalog capture, each required contract once.
                _schema._validate_schema_objects(installed, activity=True)
                if composite:
                    from adcp.reporting.outbox import status_schema

                    status_schema._validate_status_objects(installed, activity=True, status=False)
            if state.epoch != flight.epoch or state.flight is not flight:
                return False
            state.positive = flight.key
            return True
        finally:
            if state.flight is flight:
                state.flight = None

    async def durable(self) -> bool:
        from adcp.reporting.ledger.delivery import InMemoryReportingReconciliationStore
        from adcp.reporting.ledger.store import InMemoryReportingLedgerStore
        from adcp.reporting.outbox.memory import InMemoryReportingOutbox
        from adcp.reporting.outbox.routing import ReportingEnvelopeCipher
        from adcp.reporting.outbox.worker import ReportingNotificationWorker

        if self.projector is None or self.worker.activity is None:
            return False
        if self.status is not None:
            return await self._composite_durable()
        if (
            type(self.worker) is not ReportingNotificationWorker
            or type(self.worker.cipher) is not ReportingEnvelopeCipher
            or type(self.projector) is not ReportingActivityProjector
            or id(self.worker.activity) != id(self.worker.outbox)
            or id(self.projector.store) != id(self.worker.outbox)
        ):
            raise ReportingNotificationError("activity_chain_unready")
        if type(self.worker.outbox) is InMemoryReportingOutbox:
            if (
                type(self.ledger)
                not in {InMemoryReportingLedgerStore, InMemoryReportingReconciliationStore}
                or self.worker.outbox._store is not self.ledger
            ):
                raise ReportingNotificationError("activity_chain_unready")
            return False
        # Lazy imports retain base-install operation without the [pg] extra.
        from adcp.reporting.ledger.delivery_pg import PgReportingReconciliationStore
        from adcp.reporting.ledger.pg import PgReportingLedgerStore
        from adcp.reporting.outbox.pg import PgReportingOutbox

        if (
            type(self.worker.outbox) is not PgReportingOutbox
            or not isinstance(self.worker.outbox, PgReportingOutbox)
            or type(self.ledger) not in {PgReportingLedgerStore, PgReportingReconciliationStore}
            or not isinstance(self.ledger, PgReportingLedgerStore)
            or self.ledger._pool is not self.worker.outbox._pool
            or not self.ledger._notifications_enabled
        ):
            raise ReportingNotificationError("activity_chain_unready")
        return await self._schema_proven(self.ledger._pool, composite=False)

    async def _composite_durable(self) -> bool:
        """Only the closed B+C read/purge union may replace the original B reader."""
        assert self.status is not None
        if not self.status.scheduled:
            return False
        from adcp.reporting.ledger.delivery_pg import PgReportingReconciliationStore
        from adcp.reporting.ledger.pg import PgReportingLedgerStore
        from adcp.reporting.outbox.pg import PgReportingOutbox
        from adcp.reporting.outbox.routing import ReportingEnvelopeCipher
        from adcp.reporting.outbox.status_activity_pg import PgReportingActivityUnionStore
        from adcp.reporting.outbox.status_pg import (
            PgReportingStatusOutbox,
            PgStatusNotificationStore,
        )
        from adcp.reporting.outbox.worker import ReportingNotificationWorker

        store = self.status.store
        if not isinstance(store, PgStatusNotificationStore):
            return False
        reader = self.projector.store if self.projector is not None else None
        if (
            type(self.worker) is not ReportingNotificationWorker
            or type(store) is not PgStatusNotificationStore
            or type(self.ledger) not in {PgReportingLedgerStore, PgReportingReconciliationStore}
            or not isinstance(self.ledger, PgReportingLedgerStore)
            or type(self.worker.outbox) is not PgReportingOutbox
            or not isinstance(self.worker.outbox, PgReportingOutbox)
            or type(self.worker.cipher) is not ReportingEnvelopeCipher
            or type(self.status.worker) is not ReportingNotificationWorker
            or type(self.projector) is not ReportingActivityProjector
            or type(reader) is not PgReportingActivityUnionStore
            or not isinstance(reader, PgReportingActivityUnionStore)
            or reader.b_outbox is not self.worker.outbox
            or reader.status_outbox is not store.outbox
            or self.worker.activity is not self.worker.outbox
            or self.status.worker.activity is not store.outbox
            or self.status.worker.outbox is not store.outbox
            or self.ledger is not store.ledger
            or self.worker.outbox._pool is not store.ledger._pool
            or type(store.outbox) is not PgReportingStatusOutbox
            or store.outbox._pool is not store.ledger._pool
            or self.worker.outbox._clock is not store.outbox._clock
            or not store.ledger._notifications_enabled
            or self.worker.cipher is not self.status.worker.cipher
            or self.worker.subscriptions is not self.status.worker.subscriptions
            or self.worker.signing is not self.status.worker.signing
        ):
            raise ReportingNotificationError("activity_chain_unready")
        return await self._schema_proven(store.ledger._pool, composite=True)

    async def capability_flags(
        self, *, account_activity: ReportingActivityProjector | None = None
    ) -> dict[str, bool]:
        durable = await self.durable()
        return {
            "reporting": durable,
            "account_notifications": durable and account_activity is self.projector,
        }

The concrete writer/store/projector chain scheduled by the adopter.

Mount this as reporting_activity= on the platform server factory. Mount projector separately as account_activity= to enable list-accounts enrichment. Neither relationship notification claims nor memory components supply evidence of durable reporting activity.

Instance variables

var ledger : ReportingLedgerStore
var projector : ReportingActivityProjector | None
var status : ReportingStatusSupport | None
var worker : ReportingNotificationWorker

Methods

async def capability_flags(self,
*,
account_activity: ReportingActivityProjector | None = None) ‑> dict[str, bool]
Expand source code
async def capability_flags(
    self, *, account_activity: ReportingActivityProjector | None = None
) -> dict[str, bool]:
    durable = await self.durable()
    return {
        "reporting": durable,
        "account_notifications": durable and account_activity is self.projector,
    }
async def durable(self) ‑> bool
Expand source code
async def durable(self) -> bool:
    from adcp.reporting.ledger.delivery import InMemoryReportingReconciliationStore
    from adcp.reporting.ledger.store import InMemoryReportingLedgerStore
    from adcp.reporting.outbox.memory import InMemoryReportingOutbox
    from adcp.reporting.outbox.routing import ReportingEnvelopeCipher
    from adcp.reporting.outbox.worker import ReportingNotificationWorker

    if self.projector is None or self.worker.activity is None:
        return False
    if self.status is not None:
        return await self._composite_durable()
    if (
        type(self.worker) is not ReportingNotificationWorker
        or type(self.worker.cipher) is not ReportingEnvelopeCipher
        or type(self.projector) is not ReportingActivityProjector
        or id(self.worker.activity) != id(self.worker.outbox)
        or id(self.projector.store) != id(self.worker.outbox)
    ):
        raise ReportingNotificationError("activity_chain_unready")
    if type(self.worker.outbox) is InMemoryReportingOutbox:
        if (
            type(self.ledger)
            not in {InMemoryReportingLedgerStore, InMemoryReportingReconciliationStore}
            or self.worker.outbox._store is not self.ledger
        ):
            raise ReportingNotificationError("activity_chain_unready")
        return False
    # Lazy imports retain base-install operation without the [pg] extra.
    from adcp.reporting.ledger.delivery_pg import PgReportingReconciliationStore
    from adcp.reporting.ledger.pg import PgReportingLedgerStore
    from adcp.reporting.outbox.pg import PgReportingOutbox

    if (
        type(self.worker.outbox) is not PgReportingOutbox
        or not isinstance(self.worker.outbox, PgReportingOutbox)
        or type(self.ledger) not in {PgReportingLedgerStore, PgReportingReconciliationStore}
        or not isinstance(self.ledger, PgReportingLedgerStore)
        or self.ledger._pool is not self.worker.outbox._pool
        or not self.ledger._notifications_enabled
    ):
        raise ReportingNotificationError("activity_chain_unready")
    return await self._schema_proven(self.ledger._pool, composite=False)
def invalidate_schema_validation(self) ‑> None
Expand source code
def invalidate_schema_validation(self) -> None:
    """Discard schema evidence before controlled DDL, after draining callers.

    An older in-flight scan cannot publish into the new epoch. A fresh
    support/server instance after migration is the normal deployment path;
    serving-time DDL is not automatically detected by discovery.
    """
    state = self._schema_validation
    state.epoch += 1
    state.positive = None
    state.flight = None

Discard schema evidence before controlled DDL, after draining callers.

An older in-flight scan cannot publish into the new epoch. A fresh support/server instance after migration is the normal deployment path; serving-time DDL is not automatically detected by discovery.

class ReportingDomainEvent (account_id: str,
notification_id: str,
fired_at: datetime,
cause: NotificationCause,
cause_generation: int = 1)
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingDomainEvent(_ClosedValue):
    account_id: str
    notification_id: str
    fired_at: datetime
    cause: NotificationCause
    cause_generation: int = 1

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        reporting_identifier(self.notification_id, maximum=255)
        object.__setattr__(self, "fired_at", aware_utc(self.fired_at))
        if self.cause_generation < 1:
            raise ReportingNotificationError()
        if isinstance(self.cause, MaterializationReady):
            if self.cause.generation_key.account_id != self.account_id:
                raise ReportingNotificationError()
        if isinstance(self.cause, StatusChanged) and (
            self.cause.scope.account_id != self.account_id
            or self.cause.checkpoint_generation != self.cause_generation
        ):
            raise ReportingNotificationError()

    @property
    def notification_type(self) -> NotificationType:
        if isinstance(self.cause, StatusChanged):
            return "reporting.status_changed"
        return (
            "reporting.delivery_ready"
            if isinstance(self.cause, MaterializationReady)
            else "reporting.ledger_changed"
        )

    @property
    def cause_id(self) -> str:
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            return cause.reporting_revision_id
        if isinstance(cause, AdjustmentPublished):
            return cause.reporting_adjustment_id
        if isinstance(cause, StatusChanged):
            from adcp.reporting.canonical_json import canonical_json_utf8_v1

            return (
                "rpsc_"
                + hashlib.sha256(canonical_json_utf8_v1(cause.scope.checkpoint_key)).hexdigest()
            )
        # Structured encoding: consumers and materializations may reuse IDs.
        return json.dumps(
            [cause.consumer_id, cause.reporting_materialization_id], separators=(",", ":")
        )

    @property
    def consumer_namespace(self) -> str:
        if isinstance(self.cause, StatusChanged):
            return self.cause.scope.consumer_id or ""
        return self.cause.consumer_id

    @property
    def causal_key(self) -> tuple[str, str, str, str, str, int]:
        return (
            self.account_id,
            self.consumer_namespace,
            self.notification_type,
            self.cause.kind,
            self.cause_id,
            self.cause_generation,
        )

    def body(self, *, subscriber_id: str, idempotency_key: str) -> bytes:
        """Explicit wire allowlist, validated against all rc.3 conditionals."""
        from adcp.reporting.canonical_json import canonical_json_utf8_v1

        value: dict[str, Any] = {
            "account_id": self.account_id,
            "notification_id": self.notification_id,
            "notification_type": self.notification_type,
            "fired_at": iso(self.fired_at),
            "subscriber_id": subscriber_id,
            "idempotency_key": idempotency_key,
        }
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            value.update(
                change_kind=cause.kind,
                reporting_revision_id=cause.reporting_revision_id,
                finality=cause.finality,
            )
            if cause.supersedes_reporting_revision_id is not None:
                value["supersedes_reporting_revision_id"] = cause.supersedes_reporting_revision_id
        elif isinstance(cause, AdjustmentPublished):
            value.update(
                change_kind=cause.kind,
                reporting_adjustment_id=cause.reporting_adjustment_id,
                adjusts_reporting_revision_id=cause.adjusts_reporting_revision_id,
            )
        elif isinstance(cause, StatusChanged):
            key = cause.scope.generation_key
            assert key is not None
            value.update(
                delivery_config_id=key.delivery_config_id,
                delivery_config_version=key.delivery_config_version,
                feed_purpose=cause.scope.feed_purpose,
                health=cause.health,
            )
            if cause.scope.reporting_obligation_id is not None:
                value["reporting_obligation_id"] = cause.scope.reporting_obligation_id
            if cause.previous_health is not None:
                value["previous_health"] = cause.previous_health
            if cause.issue_ids:
                value["issue_ids"] = list(cause.issue_ids)
        else:
            value.update(
                delivery_config_id=cause.generation_key.delivery_config_id,
                delivery_config_version=cause.generation_key.delivery_config_version,
                reporting_revision_id=cause.reporting_revision_id,
                reporting_materialization_id=cause.reporting_materialization_id,
                readiness=cause.readiness,
                finality=cause.finality,
                data_through=iso(cause.data_through) if cause.data_through is not None else None,
                feed_purpose=cause.feed_purpose,
            )
        validate_notification_payload(value)
        return canonical_json_utf8_v1(value)

ReportingDomainEvent(account_id: 'str', notification_id: 'str', fired_at: 'datetime', cause: 'NotificationCause', cause_generation: 'int' = 1)

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Instance variables

var account_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingDomainEvent(_ClosedValue):
    account_id: str
    notification_id: str
    fired_at: datetime
    cause: NotificationCause
    cause_generation: int = 1

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        reporting_identifier(self.notification_id, maximum=255)
        object.__setattr__(self, "fired_at", aware_utc(self.fired_at))
        if self.cause_generation < 1:
            raise ReportingNotificationError()
        if isinstance(self.cause, MaterializationReady):
            if self.cause.generation_key.account_id != self.account_id:
                raise ReportingNotificationError()
        if isinstance(self.cause, StatusChanged) and (
            self.cause.scope.account_id != self.account_id
            or self.cause.checkpoint_generation != self.cause_generation
        ):
            raise ReportingNotificationError()

    @property
    def notification_type(self) -> NotificationType:
        if isinstance(self.cause, StatusChanged):
            return "reporting.status_changed"
        return (
            "reporting.delivery_ready"
            if isinstance(self.cause, MaterializationReady)
            else "reporting.ledger_changed"
        )

    @property
    def cause_id(self) -> str:
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            return cause.reporting_revision_id
        if isinstance(cause, AdjustmentPublished):
            return cause.reporting_adjustment_id
        if isinstance(cause, StatusChanged):
            from adcp.reporting.canonical_json import canonical_json_utf8_v1

            return (
                "rpsc_"
                + hashlib.sha256(canonical_json_utf8_v1(cause.scope.checkpoint_key)).hexdigest()
            )
        # Structured encoding: consumers and materializations may reuse IDs.
        return json.dumps(
            [cause.consumer_id, cause.reporting_materialization_id], separators=(",", ":")
        )

    @property
    def consumer_namespace(self) -> str:
        if isinstance(self.cause, StatusChanged):
            return self.cause.scope.consumer_id or ""
        return self.cause.consumer_id

    @property
    def causal_key(self) -> tuple[str, str, str, str, str, int]:
        return (
            self.account_id,
            self.consumer_namespace,
            self.notification_type,
            self.cause.kind,
            self.cause_id,
            self.cause_generation,
        )

    def body(self, *, subscriber_id: str, idempotency_key: str) -> bytes:
        """Explicit wire allowlist, validated against all rc.3 conditionals."""
        from adcp.reporting.canonical_json import canonical_json_utf8_v1

        value: dict[str, Any] = {
            "account_id": self.account_id,
            "notification_id": self.notification_id,
            "notification_type": self.notification_type,
            "fired_at": iso(self.fired_at),
            "subscriber_id": subscriber_id,
            "idempotency_key": idempotency_key,
        }
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            value.update(
                change_kind=cause.kind,
                reporting_revision_id=cause.reporting_revision_id,
                finality=cause.finality,
            )
            if cause.supersedes_reporting_revision_id is not None:
                value["supersedes_reporting_revision_id"] = cause.supersedes_reporting_revision_id
        elif isinstance(cause, AdjustmentPublished):
            value.update(
                change_kind=cause.kind,
                reporting_adjustment_id=cause.reporting_adjustment_id,
                adjusts_reporting_revision_id=cause.adjusts_reporting_revision_id,
            )
        elif isinstance(cause, StatusChanged):
            key = cause.scope.generation_key
            assert key is not None
            value.update(
                delivery_config_id=key.delivery_config_id,
                delivery_config_version=key.delivery_config_version,
                feed_purpose=cause.scope.feed_purpose,
                health=cause.health,
            )
            if cause.scope.reporting_obligation_id is not None:
                value["reporting_obligation_id"] = cause.scope.reporting_obligation_id
            if cause.previous_health is not None:
                value["previous_health"] = cause.previous_health
            if cause.issue_ids:
                value["issue_ids"] = list(cause.issue_ids)
        else:
            value.update(
                delivery_config_id=cause.generation_key.delivery_config_id,
                delivery_config_version=cause.generation_key.delivery_config_version,
                reporting_revision_id=cause.reporting_revision_id,
                reporting_materialization_id=cause.reporting_materialization_id,
                readiness=cause.readiness,
                finality=cause.finality,
                data_through=iso(cause.data_through) if cause.data_through is not None else None,
                feed_purpose=cause.feed_purpose,
            )
        validate_notification_payload(value)
        return canonical_json_utf8_v1(value)
prop causal_key : tuple[str, str, str, str, str, int]
Expand source code
@property
def causal_key(self) -> tuple[str, str, str, str, str, int]:
    return (
        self.account_id,
        self.consumer_namespace,
        self.notification_type,
        self.cause.kind,
        self.cause_id,
        self.cause_generation,
    )
var cause : RevisionPublished | AdjustmentPublished | MaterializationReady | StatusChanged
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingDomainEvent(_ClosedValue):
    account_id: str
    notification_id: str
    fired_at: datetime
    cause: NotificationCause
    cause_generation: int = 1

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        reporting_identifier(self.notification_id, maximum=255)
        object.__setattr__(self, "fired_at", aware_utc(self.fired_at))
        if self.cause_generation < 1:
            raise ReportingNotificationError()
        if isinstance(self.cause, MaterializationReady):
            if self.cause.generation_key.account_id != self.account_id:
                raise ReportingNotificationError()
        if isinstance(self.cause, StatusChanged) and (
            self.cause.scope.account_id != self.account_id
            or self.cause.checkpoint_generation != self.cause_generation
        ):
            raise ReportingNotificationError()

    @property
    def notification_type(self) -> NotificationType:
        if isinstance(self.cause, StatusChanged):
            return "reporting.status_changed"
        return (
            "reporting.delivery_ready"
            if isinstance(self.cause, MaterializationReady)
            else "reporting.ledger_changed"
        )

    @property
    def cause_id(self) -> str:
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            return cause.reporting_revision_id
        if isinstance(cause, AdjustmentPublished):
            return cause.reporting_adjustment_id
        if isinstance(cause, StatusChanged):
            from adcp.reporting.canonical_json import canonical_json_utf8_v1

            return (
                "rpsc_"
                + hashlib.sha256(canonical_json_utf8_v1(cause.scope.checkpoint_key)).hexdigest()
            )
        # Structured encoding: consumers and materializations may reuse IDs.
        return json.dumps(
            [cause.consumer_id, cause.reporting_materialization_id], separators=(",", ":")
        )

    @property
    def consumer_namespace(self) -> str:
        if isinstance(self.cause, StatusChanged):
            return self.cause.scope.consumer_id or ""
        return self.cause.consumer_id

    @property
    def causal_key(self) -> tuple[str, str, str, str, str, int]:
        return (
            self.account_id,
            self.consumer_namespace,
            self.notification_type,
            self.cause.kind,
            self.cause_id,
            self.cause_generation,
        )

    def body(self, *, subscriber_id: str, idempotency_key: str) -> bytes:
        """Explicit wire allowlist, validated against all rc.3 conditionals."""
        from adcp.reporting.canonical_json import canonical_json_utf8_v1

        value: dict[str, Any] = {
            "account_id": self.account_id,
            "notification_id": self.notification_id,
            "notification_type": self.notification_type,
            "fired_at": iso(self.fired_at),
            "subscriber_id": subscriber_id,
            "idempotency_key": idempotency_key,
        }
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            value.update(
                change_kind=cause.kind,
                reporting_revision_id=cause.reporting_revision_id,
                finality=cause.finality,
            )
            if cause.supersedes_reporting_revision_id is not None:
                value["supersedes_reporting_revision_id"] = cause.supersedes_reporting_revision_id
        elif isinstance(cause, AdjustmentPublished):
            value.update(
                change_kind=cause.kind,
                reporting_adjustment_id=cause.reporting_adjustment_id,
                adjusts_reporting_revision_id=cause.adjusts_reporting_revision_id,
            )
        elif isinstance(cause, StatusChanged):
            key = cause.scope.generation_key
            assert key is not None
            value.update(
                delivery_config_id=key.delivery_config_id,
                delivery_config_version=key.delivery_config_version,
                feed_purpose=cause.scope.feed_purpose,
                health=cause.health,
            )
            if cause.scope.reporting_obligation_id is not None:
                value["reporting_obligation_id"] = cause.scope.reporting_obligation_id
            if cause.previous_health is not None:
                value["previous_health"] = cause.previous_health
            if cause.issue_ids:
                value["issue_ids"] = list(cause.issue_ids)
        else:
            value.update(
                delivery_config_id=cause.generation_key.delivery_config_id,
                delivery_config_version=cause.generation_key.delivery_config_version,
                reporting_revision_id=cause.reporting_revision_id,
                reporting_materialization_id=cause.reporting_materialization_id,
                readiness=cause.readiness,
                finality=cause.finality,
                data_through=iso(cause.data_through) if cause.data_through is not None else None,
                feed_purpose=cause.feed_purpose,
            )
        validate_notification_payload(value)
        return canonical_json_utf8_v1(value)
var cause_generation : int
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingDomainEvent(_ClosedValue):
    account_id: str
    notification_id: str
    fired_at: datetime
    cause: NotificationCause
    cause_generation: int = 1

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        reporting_identifier(self.notification_id, maximum=255)
        object.__setattr__(self, "fired_at", aware_utc(self.fired_at))
        if self.cause_generation < 1:
            raise ReportingNotificationError()
        if isinstance(self.cause, MaterializationReady):
            if self.cause.generation_key.account_id != self.account_id:
                raise ReportingNotificationError()
        if isinstance(self.cause, StatusChanged) and (
            self.cause.scope.account_id != self.account_id
            or self.cause.checkpoint_generation != self.cause_generation
        ):
            raise ReportingNotificationError()

    @property
    def notification_type(self) -> NotificationType:
        if isinstance(self.cause, StatusChanged):
            return "reporting.status_changed"
        return (
            "reporting.delivery_ready"
            if isinstance(self.cause, MaterializationReady)
            else "reporting.ledger_changed"
        )

    @property
    def cause_id(self) -> str:
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            return cause.reporting_revision_id
        if isinstance(cause, AdjustmentPublished):
            return cause.reporting_adjustment_id
        if isinstance(cause, StatusChanged):
            from adcp.reporting.canonical_json import canonical_json_utf8_v1

            return (
                "rpsc_"
                + hashlib.sha256(canonical_json_utf8_v1(cause.scope.checkpoint_key)).hexdigest()
            )
        # Structured encoding: consumers and materializations may reuse IDs.
        return json.dumps(
            [cause.consumer_id, cause.reporting_materialization_id], separators=(",", ":")
        )

    @property
    def consumer_namespace(self) -> str:
        if isinstance(self.cause, StatusChanged):
            return self.cause.scope.consumer_id or ""
        return self.cause.consumer_id

    @property
    def causal_key(self) -> tuple[str, str, str, str, str, int]:
        return (
            self.account_id,
            self.consumer_namespace,
            self.notification_type,
            self.cause.kind,
            self.cause_id,
            self.cause_generation,
        )

    def body(self, *, subscriber_id: str, idempotency_key: str) -> bytes:
        """Explicit wire allowlist, validated against all rc.3 conditionals."""
        from adcp.reporting.canonical_json import canonical_json_utf8_v1

        value: dict[str, Any] = {
            "account_id": self.account_id,
            "notification_id": self.notification_id,
            "notification_type": self.notification_type,
            "fired_at": iso(self.fired_at),
            "subscriber_id": subscriber_id,
            "idempotency_key": idempotency_key,
        }
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            value.update(
                change_kind=cause.kind,
                reporting_revision_id=cause.reporting_revision_id,
                finality=cause.finality,
            )
            if cause.supersedes_reporting_revision_id is not None:
                value["supersedes_reporting_revision_id"] = cause.supersedes_reporting_revision_id
        elif isinstance(cause, AdjustmentPublished):
            value.update(
                change_kind=cause.kind,
                reporting_adjustment_id=cause.reporting_adjustment_id,
                adjusts_reporting_revision_id=cause.adjusts_reporting_revision_id,
            )
        elif isinstance(cause, StatusChanged):
            key = cause.scope.generation_key
            assert key is not None
            value.update(
                delivery_config_id=key.delivery_config_id,
                delivery_config_version=key.delivery_config_version,
                feed_purpose=cause.scope.feed_purpose,
                health=cause.health,
            )
            if cause.scope.reporting_obligation_id is not None:
                value["reporting_obligation_id"] = cause.scope.reporting_obligation_id
            if cause.previous_health is not None:
                value["previous_health"] = cause.previous_health
            if cause.issue_ids:
                value["issue_ids"] = list(cause.issue_ids)
        else:
            value.update(
                delivery_config_id=cause.generation_key.delivery_config_id,
                delivery_config_version=cause.generation_key.delivery_config_version,
                reporting_revision_id=cause.reporting_revision_id,
                reporting_materialization_id=cause.reporting_materialization_id,
                readiness=cause.readiness,
                finality=cause.finality,
                data_through=iso(cause.data_through) if cause.data_through is not None else None,
                feed_purpose=cause.feed_purpose,
            )
        validate_notification_payload(value)
        return canonical_json_utf8_v1(value)
prop cause_id : str
Expand source code
@property
def cause_id(self) -> str:
    cause = self.cause
    if isinstance(cause, RevisionPublished):
        return cause.reporting_revision_id
    if isinstance(cause, AdjustmentPublished):
        return cause.reporting_adjustment_id
    if isinstance(cause, StatusChanged):
        from adcp.reporting.canonical_json import canonical_json_utf8_v1

        return (
            "rpsc_"
            + hashlib.sha256(canonical_json_utf8_v1(cause.scope.checkpoint_key)).hexdigest()
        )
    # Structured encoding: consumers and materializations may reuse IDs.
    return json.dumps(
        [cause.consumer_id, cause.reporting_materialization_id], separators=(",", ":")
    )
prop consumer_namespace : str
Expand source code
@property
def consumer_namespace(self) -> str:
    if isinstance(self.cause, StatusChanged):
        return self.cause.scope.consumer_id or ""
    return self.cause.consumer_id
var fired_at : datetime.datetime
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingDomainEvent(_ClosedValue):
    account_id: str
    notification_id: str
    fired_at: datetime
    cause: NotificationCause
    cause_generation: int = 1

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        reporting_identifier(self.notification_id, maximum=255)
        object.__setattr__(self, "fired_at", aware_utc(self.fired_at))
        if self.cause_generation < 1:
            raise ReportingNotificationError()
        if isinstance(self.cause, MaterializationReady):
            if self.cause.generation_key.account_id != self.account_id:
                raise ReportingNotificationError()
        if isinstance(self.cause, StatusChanged) and (
            self.cause.scope.account_id != self.account_id
            or self.cause.checkpoint_generation != self.cause_generation
        ):
            raise ReportingNotificationError()

    @property
    def notification_type(self) -> NotificationType:
        if isinstance(self.cause, StatusChanged):
            return "reporting.status_changed"
        return (
            "reporting.delivery_ready"
            if isinstance(self.cause, MaterializationReady)
            else "reporting.ledger_changed"
        )

    @property
    def cause_id(self) -> str:
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            return cause.reporting_revision_id
        if isinstance(cause, AdjustmentPublished):
            return cause.reporting_adjustment_id
        if isinstance(cause, StatusChanged):
            from adcp.reporting.canonical_json import canonical_json_utf8_v1

            return (
                "rpsc_"
                + hashlib.sha256(canonical_json_utf8_v1(cause.scope.checkpoint_key)).hexdigest()
            )
        # Structured encoding: consumers and materializations may reuse IDs.
        return json.dumps(
            [cause.consumer_id, cause.reporting_materialization_id], separators=(",", ":")
        )

    @property
    def consumer_namespace(self) -> str:
        if isinstance(self.cause, StatusChanged):
            return self.cause.scope.consumer_id or ""
        return self.cause.consumer_id

    @property
    def causal_key(self) -> tuple[str, str, str, str, str, int]:
        return (
            self.account_id,
            self.consumer_namespace,
            self.notification_type,
            self.cause.kind,
            self.cause_id,
            self.cause_generation,
        )

    def body(self, *, subscriber_id: str, idempotency_key: str) -> bytes:
        """Explicit wire allowlist, validated against all rc.3 conditionals."""
        from adcp.reporting.canonical_json import canonical_json_utf8_v1

        value: dict[str, Any] = {
            "account_id": self.account_id,
            "notification_id": self.notification_id,
            "notification_type": self.notification_type,
            "fired_at": iso(self.fired_at),
            "subscriber_id": subscriber_id,
            "idempotency_key": idempotency_key,
        }
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            value.update(
                change_kind=cause.kind,
                reporting_revision_id=cause.reporting_revision_id,
                finality=cause.finality,
            )
            if cause.supersedes_reporting_revision_id is not None:
                value["supersedes_reporting_revision_id"] = cause.supersedes_reporting_revision_id
        elif isinstance(cause, AdjustmentPublished):
            value.update(
                change_kind=cause.kind,
                reporting_adjustment_id=cause.reporting_adjustment_id,
                adjusts_reporting_revision_id=cause.adjusts_reporting_revision_id,
            )
        elif isinstance(cause, StatusChanged):
            key = cause.scope.generation_key
            assert key is not None
            value.update(
                delivery_config_id=key.delivery_config_id,
                delivery_config_version=key.delivery_config_version,
                feed_purpose=cause.scope.feed_purpose,
                health=cause.health,
            )
            if cause.scope.reporting_obligation_id is not None:
                value["reporting_obligation_id"] = cause.scope.reporting_obligation_id
            if cause.previous_health is not None:
                value["previous_health"] = cause.previous_health
            if cause.issue_ids:
                value["issue_ids"] = list(cause.issue_ids)
        else:
            value.update(
                delivery_config_id=cause.generation_key.delivery_config_id,
                delivery_config_version=cause.generation_key.delivery_config_version,
                reporting_revision_id=cause.reporting_revision_id,
                reporting_materialization_id=cause.reporting_materialization_id,
                readiness=cause.readiness,
                finality=cause.finality,
                data_through=iso(cause.data_through) if cause.data_through is not None else None,
                feed_purpose=cause.feed_purpose,
            )
        validate_notification_payload(value)
        return canonical_json_utf8_v1(value)
var notification_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingDomainEvent(_ClosedValue):
    account_id: str
    notification_id: str
    fired_at: datetime
    cause: NotificationCause
    cause_generation: int = 1

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        reporting_identifier(self.notification_id, maximum=255)
        object.__setattr__(self, "fired_at", aware_utc(self.fired_at))
        if self.cause_generation < 1:
            raise ReportingNotificationError()
        if isinstance(self.cause, MaterializationReady):
            if self.cause.generation_key.account_id != self.account_id:
                raise ReportingNotificationError()
        if isinstance(self.cause, StatusChanged) and (
            self.cause.scope.account_id != self.account_id
            or self.cause.checkpoint_generation != self.cause_generation
        ):
            raise ReportingNotificationError()

    @property
    def notification_type(self) -> NotificationType:
        if isinstance(self.cause, StatusChanged):
            return "reporting.status_changed"
        return (
            "reporting.delivery_ready"
            if isinstance(self.cause, MaterializationReady)
            else "reporting.ledger_changed"
        )

    @property
    def cause_id(self) -> str:
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            return cause.reporting_revision_id
        if isinstance(cause, AdjustmentPublished):
            return cause.reporting_adjustment_id
        if isinstance(cause, StatusChanged):
            from adcp.reporting.canonical_json import canonical_json_utf8_v1

            return (
                "rpsc_"
                + hashlib.sha256(canonical_json_utf8_v1(cause.scope.checkpoint_key)).hexdigest()
            )
        # Structured encoding: consumers and materializations may reuse IDs.
        return json.dumps(
            [cause.consumer_id, cause.reporting_materialization_id], separators=(",", ":")
        )

    @property
    def consumer_namespace(self) -> str:
        if isinstance(self.cause, StatusChanged):
            return self.cause.scope.consumer_id or ""
        return self.cause.consumer_id

    @property
    def causal_key(self) -> tuple[str, str, str, str, str, int]:
        return (
            self.account_id,
            self.consumer_namespace,
            self.notification_type,
            self.cause.kind,
            self.cause_id,
            self.cause_generation,
        )

    def body(self, *, subscriber_id: str, idempotency_key: str) -> bytes:
        """Explicit wire allowlist, validated against all rc.3 conditionals."""
        from adcp.reporting.canonical_json import canonical_json_utf8_v1

        value: dict[str, Any] = {
            "account_id": self.account_id,
            "notification_id": self.notification_id,
            "notification_type": self.notification_type,
            "fired_at": iso(self.fired_at),
            "subscriber_id": subscriber_id,
            "idempotency_key": idempotency_key,
        }
        cause = self.cause
        if isinstance(cause, RevisionPublished):
            value.update(
                change_kind=cause.kind,
                reporting_revision_id=cause.reporting_revision_id,
                finality=cause.finality,
            )
            if cause.supersedes_reporting_revision_id is not None:
                value["supersedes_reporting_revision_id"] = cause.supersedes_reporting_revision_id
        elif isinstance(cause, AdjustmentPublished):
            value.update(
                change_kind=cause.kind,
                reporting_adjustment_id=cause.reporting_adjustment_id,
                adjusts_reporting_revision_id=cause.adjusts_reporting_revision_id,
            )
        elif isinstance(cause, StatusChanged):
            key = cause.scope.generation_key
            assert key is not None
            value.update(
                delivery_config_id=key.delivery_config_id,
                delivery_config_version=key.delivery_config_version,
                feed_purpose=cause.scope.feed_purpose,
                health=cause.health,
            )
            if cause.scope.reporting_obligation_id is not None:
                value["reporting_obligation_id"] = cause.scope.reporting_obligation_id
            if cause.previous_health is not None:
                value["previous_health"] = cause.previous_health
            if cause.issue_ids:
                value["issue_ids"] = list(cause.issue_ids)
        else:
            value.update(
                delivery_config_id=cause.generation_key.delivery_config_id,
                delivery_config_version=cause.generation_key.delivery_config_version,
                reporting_revision_id=cause.reporting_revision_id,
                reporting_materialization_id=cause.reporting_materialization_id,
                readiness=cause.readiness,
                finality=cause.finality,
                data_through=iso(cause.data_through) if cause.data_through is not None else None,
                feed_purpose=cause.feed_purpose,
            )
        validate_notification_payload(value)
        return canonical_json_utf8_v1(value)
prop notification_type : NotificationType
Expand source code
@property
def notification_type(self) -> NotificationType:
    if isinstance(self.cause, StatusChanged):
        return "reporting.status_changed"
    return (
        "reporting.delivery_ready"
        if isinstance(self.cause, MaterializationReady)
        else "reporting.ledger_changed"
    )

Methods

def body(self, *, subscriber_id: str, idempotency_key: str) ‑> bytes
Expand source code
def body(self, *, subscriber_id: str, idempotency_key: str) -> bytes:
    """Explicit wire allowlist, validated against all rc.3 conditionals."""
    from adcp.reporting.canonical_json import canonical_json_utf8_v1

    value: dict[str, Any] = {
        "account_id": self.account_id,
        "notification_id": self.notification_id,
        "notification_type": self.notification_type,
        "fired_at": iso(self.fired_at),
        "subscriber_id": subscriber_id,
        "idempotency_key": idempotency_key,
    }
    cause = self.cause
    if isinstance(cause, RevisionPublished):
        value.update(
            change_kind=cause.kind,
            reporting_revision_id=cause.reporting_revision_id,
            finality=cause.finality,
        )
        if cause.supersedes_reporting_revision_id is not None:
            value["supersedes_reporting_revision_id"] = cause.supersedes_reporting_revision_id
    elif isinstance(cause, AdjustmentPublished):
        value.update(
            change_kind=cause.kind,
            reporting_adjustment_id=cause.reporting_adjustment_id,
            adjusts_reporting_revision_id=cause.adjusts_reporting_revision_id,
        )
    elif isinstance(cause, StatusChanged):
        key = cause.scope.generation_key
        assert key is not None
        value.update(
            delivery_config_id=key.delivery_config_id,
            delivery_config_version=key.delivery_config_version,
            feed_purpose=cause.scope.feed_purpose,
            health=cause.health,
        )
        if cause.scope.reporting_obligation_id is not None:
            value["reporting_obligation_id"] = cause.scope.reporting_obligation_id
        if cause.previous_health is not None:
            value["previous_health"] = cause.previous_health
        if cause.issue_ids:
            value["issue_ids"] = list(cause.issue_ids)
    else:
        value.update(
            delivery_config_id=cause.generation_key.delivery_config_id,
            delivery_config_version=cause.generation_key.delivery_config_version,
            reporting_revision_id=cause.reporting_revision_id,
            reporting_materialization_id=cause.reporting_materialization_id,
            readiness=cause.readiness,
            finality=cause.finality,
            data_through=iso(cause.data_through) if cause.data_through is not None else None,
            feed_purpose=cause.feed_purpose,
        )
    validate_notification_payload(value)
    return canonical_json_utf8_v1(value)

Explicit wire allowlist, validated against all rc.3 conditionals.

class ReportingEnvelopeCipher (key: bytes,
*,
key_version: str = 'v1',
previous_keys: Mapping[str, bytes] | None = None)
Expand source code
class ReportingEnvelopeCipher:
    """AES-256-GCM with all routing columns and body identity in canonical AAD.

    Key versions enable rolling encryption-key rotation while retaining old
    decryptors. Signing-key rotation is independent and resolved per attempt.
    Neither ciphertext nor plaintext routing is ever logged by this module.
    """

    def __init__(
        self,
        key: bytes,
        *,
        key_version: str = "v1",
        previous_keys: Mapping[str, bytes] | None = None,
    ) -> None:
        keys = dict(previous_keys or {})
        keys[key_version] = key
        if any(type(value) is not bytes or len(value) != 32 for value in keys.values()):
            raise ValueError("reporting outbox encryption keys must contain 32 bytes")
        for version in keys:
            reporting_identifier(version, maximum=64)
        self._keys, self.key_version = keys, key_version

    @staticmethod
    def _aad(binding: DeliveryBinding) -> bytes:
        return canonical_json_utf8_v1({"purpose": "adcp.reporting.delivery", **asdict(binding)})

    def prepare(
        self,
        event: ReportingDomainEvent,
        subscription: ReportingNotificationSubscription,
        emission_generation: int,
    ) -> StoredDelivery:
        if type(subscription) is not ReportingNotificationSubscription or not subscription.matches(
            event
        ):
            raise ReportingNotificationError("invalid_configuration")
        if type(emission_generation) is not int or emission_generation < 1:
            raise ReportingNotificationError("invalid_configuration")
        idempotency_key = str(uuid4())
        body = event.body(subscriber_id=subscription.subscriber_id, idempotency_key=idempotency_key)
        binding = DeliveryBinding(
            account_id=event.account_id,
            delivery_id=str(uuid4()),
            subscriber_id=subscription.subscriber_id,
            principal_id=subscription.principal_id,
            notification_id=event.notification_id,
            notification_type=event.notification_type,
            emission_generation=emission_generation,
            idempotency_key=idempotency_key,
            destination_sha256=hashlib.sha256(subscription.url.encode("utf-8")).hexdigest(),
            subscription_fingerprint=subscription.fingerprint,
            signing_scope_id=subscription.signing_scope_id,
            cause_kind=event.cause.kind,
            cause_id=event.cause_id,
            cause_generation=event.cause_generation,
            consumer_namespace=event.consumer_namespace,
            auth_mode=subscription.auth_mode,
            body_sha256=hashlib.sha256(body).hexdigest(),
            envelope_version=1,
            key_version=self.key_version,
        )
        plaintext = canonical_json_utf8_v1(
            {
                "subscription": asdict(subscription),
                "event": event_storage(event),
                "body": base64.b64encode(body).decode("ascii"),
            }
        )
        nonce = secrets.token_bytes(12)
        encrypted = AESGCM(self._keys[self.key_version]).encrypt(
            nonce, plaintext, self._aad(binding)
        )
        return StoredDelivery(binding, nonce + encrypted)

    def open(self, delivery: StoredDelivery) -> OpenedReportingDelivery:
        """Authenticate all bound columns before any resolver/DNS/signing/HTTP."""
        try:
            binding = delivery.binding
            if (
                type(delivery.envelope) is not bytes
                or not 28 <= len(delivery.envelope) <= 1024 * 1024
            ):
                raise ValueError
            key = self._keys[binding.key_version]
            plaintext = AESGCM(key).decrypt(
                delivery.envelope[:12], delivery.envelope[12:], self._aad(binding)
            )
            if binding.envelope_version != 1:
                raise ValueError
            value = json.loads(plaintext)
            if set(value) != {"subscription", "event", "body"}:
                raise ValueError
            subscription = _decode_subscription(value["subscription"])
            event = decode_event(value["event"])
            body = base64.b64decode(value["body"], validate=True)
            payload = json.loads(body)
            validate_notification_payload(payload)
            if (
                not subscription.matches(event)
                or subscription.account_id != binding.account_id
                or subscription.subscriber_id != binding.subscriber_id
                or subscription.principal_id != binding.principal_id
                or subscription.signing_scope_id != binding.signing_scope_id
                or subscription.auth_mode != binding.auth_mode
                or subscription.fingerprint != binding.subscription_fingerprint
                or hashlib.sha256(subscription.url.encode("utf-8")).hexdigest()
                != binding.destination_sha256
                or event.notification_id != binding.notification_id
                or event.notification_type != binding.notification_type
                or event.cause.kind != binding.cause_kind
                or event.cause_id != binding.cause_id
                or event.cause_generation != binding.cause_generation
                or event.consumer_namespace != binding.consumer_namespace
                or hashlib.sha256(body).hexdigest() != binding.body_sha256
                or body
                != event.body(
                    subscriber_id=binding.subscriber_id, idempotency_key=binding.idempotency_key
                )
            ):
                raise ValueError
            return OpenedReportingDelivery(
                subscription,
                event,
                PreparedWebhook(subscription.url, binding.idempotency_key, body),
            )
        except Exception:
            # This phase is entirely local. Deeply malformed/oversized poison
            # is terminal too; it must not strand a worker before later rows.
            raise ReportingNotificationError("integrity_failure") from None

AES-256-GCM with all routing columns and body identity in canonical AAD.

Key versions enable rolling encryption-key rotation while retaining old decryptors. Signing-key rotation is independent and resolved per attempt. Neither ciphertext nor plaintext routing is ever logged by this module.

Methods

def open(self,
delivery: StoredDelivery) ‑> OpenedReportingDelivery
Expand source code
def open(self, delivery: StoredDelivery) -> OpenedReportingDelivery:
    """Authenticate all bound columns before any resolver/DNS/signing/HTTP."""
    try:
        binding = delivery.binding
        if (
            type(delivery.envelope) is not bytes
            or not 28 <= len(delivery.envelope) <= 1024 * 1024
        ):
            raise ValueError
        key = self._keys[binding.key_version]
        plaintext = AESGCM(key).decrypt(
            delivery.envelope[:12], delivery.envelope[12:], self._aad(binding)
        )
        if binding.envelope_version != 1:
            raise ValueError
        value = json.loads(plaintext)
        if set(value) != {"subscription", "event", "body"}:
            raise ValueError
        subscription = _decode_subscription(value["subscription"])
        event = decode_event(value["event"])
        body = base64.b64decode(value["body"], validate=True)
        payload = json.loads(body)
        validate_notification_payload(payload)
        if (
            not subscription.matches(event)
            or subscription.account_id != binding.account_id
            or subscription.subscriber_id != binding.subscriber_id
            or subscription.principal_id != binding.principal_id
            or subscription.signing_scope_id != binding.signing_scope_id
            or subscription.auth_mode != binding.auth_mode
            or subscription.fingerprint != binding.subscription_fingerprint
            or hashlib.sha256(subscription.url.encode("utf-8")).hexdigest()
            != binding.destination_sha256
            or event.notification_id != binding.notification_id
            or event.notification_type != binding.notification_type
            or event.cause.kind != binding.cause_kind
            or event.cause_id != binding.cause_id
            or event.cause_generation != binding.cause_generation
            or event.consumer_namespace != binding.consumer_namespace
            or hashlib.sha256(body).hexdigest() != binding.body_sha256
            or body
            != event.body(
                subscriber_id=binding.subscriber_id, idempotency_key=binding.idempotency_key
            )
        ):
            raise ValueError
        return OpenedReportingDelivery(
            subscription,
            event,
            PreparedWebhook(subscription.url, binding.idempotency_key, body),
        )
    except Exception:
        # This phase is entirely local. Deeply malformed/oversized poison
        # is terminal too; it must not strand a worker before later rows.
        raise ReportingNotificationError("integrity_failure") from None

Authenticate all bound columns before any resolver/DNS/signing/HTTP.

def prepare(self,
event: ReportingDomainEvent,
subscription: ReportingNotificationSubscription,
emission_generation: int) ‑> StoredDelivery
Expand source code
def prepare(
    self,
    event: ReportingDomainEvent,
    subscription: ReportingNotificationSubscription,
    emission_generation: int,
) -> StoredDelivery:
    if type(subscription) is not ReportingNotificationSubscription or not subscription.matches(
        event
    ):
        raise ReportingNotificationError("invalid_configuration")
    if type(emission_generation) is not int or emission_generation < 1:
        raise ReportingNotificationError("invalid_configuration")
    idempotency_key = str(uuid4())
    body = event.body(subscriber_id=subscription.subscriber_id, idempotency_key=idempotency_key)
    binding = DeliveryBinding(
        account_id=event.account_id,
        delivery_id=str(uuid4()),
        subscriber_id=subscription.subscriber_id,
        principal_id=subscription.principal_id,
        notification_id=event.notification_id,
        notification_type=event.notification_type,
        emission_generation=emission_generation,
        idempotency_key=idempotency_key,
        destination_sha256=hashlib.sha256(subscription.url.encode("utf-8")).hexdigest(),
        subscription_fingerprint=subscription.fingerprint,
        signing_scope_id=subscription.signing_scope_id,
        cause_kind=event.cause.kind,
        cause_id=event.cause_id,
        cause_generation=event.cause_generation,
        consumer_namespace=event.consumer_namespace,
        auth_mode=subscription.auth_mode,
        body_sha256=hashlib.sha256(body).hexdigest(),
        envelope_version=1,
        key_version=self.key_version,
    )
    plaintext = canonical_json_utf8_v1(
        {
            "subscription": asdict(subscription),
            "event": event_storage(event),
            "body": base64.b64encode(body).decode("ascii"),
        }
    )
    nonce = secrets.token_bytes(12)
    encrypted = AESGCM(self._keys[self.key_version]).encrypt(
        nonce, plaintext, self._aad(binding)
    )
    return StoredDelivery(binding, nonce + encrypted)
class ReportingLegacyAuthentication (scheme: "Literal['Bearer', 'HMAC-SHA256']", credentials: str)
Expand source code
@dataclass(frozen=True)
class ReportingLegacyAuthentication:
    scheme: Literal["Bearer", "HMAC-SHA256"]
    credentials: str = field(repr=False)

    def __post_init__(self) -> None:
        if (
            self.scheme not in {"Bearer", "HMAC-SHA256"}
            or type(self.credentials) is not str
            or not self.credentials
            or not self.credentials.isprintable()
        ):
            raise ReportingNotificationError("invalid_configuration")

ReportingLegacyAuthentication(scheme: "Literal['Bearer', 'HMAC-SHA256']", credentials: 'str')

Instance variables

var credentials : str
var scheme : Literal['Bearer', 'HMAC-SHA256']
class ReportingNotificationError (code: str = 'invalid_notification')
Expand source code
class ReportingNotificationError(ValueError):
    """A sanitized local classification, never an external diagnostic."""

    def __init__(self, code: str = "invalid_notification") -> None:
        self.code = code
        super().__init__(code)

A sanitized local classification, never an external diagnostic.

Ancestors

  • builtins.ValueError
  • builtins.Exception
  • builtins.BaseException
class ReportingNotificationOutbox (*args, **kwargs)
Expand source code
class ReportingNotificationOutbox(Protocol):
    async def list_events(self, *, account_id: str) -> tuple[ReportingDomainEvent, ...]: ...

    async def claim_expansion(
        self, *, account_id: str, now: datetime, lease_seconds: float
    ) -> ExpansionLease | None: ...

    async def complete_expansion(
        self, lease: ExpansionLease, deliveries: tuple[StoredDelivery, ...], *, now: datetime
    ) -> bool: ...

    async def finish_expansion(
        self,
        lease: ExpansionLease,
        *,
        now: datetime,
        state: WorkState,
        error_code: ErrorCode | None = None,
        retry_at: datetime | None = None,
    ) -> bool: ...

    async def reemit(self, *, account_id: str, notification_id: str, now: datetime) -> int: ...

    async def claim_delivery(
        self, *, account_id: str, now: datetime, lease_seconds: float
    ) -> DeliveryLease | None: ...

    async def delivery_lease_current(self, lease: DeliveryLease, *, now: datetime) -> bool: ...

    async def finish_delivery(
        self,
        lease: DeliveryLease,
        *,
        now: datetime,
        state: WorkState,
        error_code: ErrorCode | None = None,
        retry_at: datetime | None = None,
    ) -> bool: ...

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

    async def read_status_dirty(
        self, *, account_id: str, after: int = 0, limit: int = 100
    ) -> tuple[ReportingStatusDirty, ...]: ...

    async def status_checkpoint(self, *, account_id: str, projector_id: str) -> int: ...

    async def advance_status_checkpoint(
        self, *, account_id: str, projector_id: str, expected: int, through: int
    ) -> bool: ...

Base class for protocol classes.

Protocol classes are defined as::

class Proto(Protocol):
    def meth(self) -> int:
        ...

Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).

For example::

class C:
    def meth(self) -> int:
        return 0

def func(x: Proto) -> int:
    return x.meth()

func(C())  # Passes static type check

See PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::

class GenProto(Protocol[T]):
    def meth(self) -> T:
        ...

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def advance_status_checkpoint(self, *, account_id: str, projector_id: str, expected: int, through: int) ‑> bool
Expand source code
async def advance_status_checkpoint(
    self, *, account_id: str, projector_id: str, expected: int, through: int
) -> bool: ...
async def claim_delivery(self, *, account_id: str, now: datetime, lease_seconds: float) ‑> DeliveryLease | None
Expand source code
async def claim_delivery(
    self, *, account_id: str, now: datetime, lease_seconds: float
) -> DeliveryLease | None: ...
async def claim_expansion(self, *, account_id: str, now: datetime, lease_seconds: float) ‑> ExpansionLease | None
Expand source code
async def claim_expansion(
    self, *, account_id: str, now: datetime, lease_seconds: float
) -> ExpansionLease | None: ...
async def complete_expansion(self,
lease: ExpansionLease,
deliveries: tuple[StoredDelivery, ...],
*,
now: datetime) ‑> bool
Expand source code
async def complete_expansion(
    self, lease: ExpansionLease, deliveries: tuple[StoredDelivery, ...], *, now: datetime
) -> bool: ...
async def delivery_lease_current(self,
lease: DeliveryLease,
*,
now: datetime) ‑> bool
Expand source code
async def delivery_lease_current(self, lease: DeliveryLease, *, now: datetime) -> bool: ...
async def finish_delivery(self,
lease: DeliveryLease,
*,
now: datetime,
state: WorkState,
error_code: ErrorCode | None = None,
retry_at: datetime | None = None) ‑> bool
Expand source code
async def finish_delivery(
    self,
    lease: DeliveryLease,
    *,
    now: datetime,
    state: WorkState,
    error_code: ErrorCode | None = None,
    retry_at: datetime | None = None,
) -> bool: ...
async def finish_expansion(self,
lease: ExpansionLease,
*,
now: datetime,
state: WorkState,
error_code: ErrorCode | None = None,
retry_at: datetime | None = None) ‑> bool
Expand source code
async def finish_expansion(
    self,
    lease: ExpansionLease,
    *,
    now: datetime,
    state: WorkState,
    error_code: ErrorCode | None = None,
    retry_at: datetime | None = None,
) -> bool: ...
async def list_deliveries(self, *, account_id: str) ‑> tuple[DeliveryStatus, ...]
Expand source code
async def list_deliveries(self, *, account_id: str) -> tuple[DeliveryStatus, ...]: ...
async def list_events(self, *, account_id: str) ‑> tuple[ReportingDomainEvent, ...]
Expand source code
async def list_events(self, *, account_id: str) -> tuple[ReportingDomainEvent, ...]: ...
async def read_status_dirty(self, *, account_id: str, after: int = 0, limit: int = 100) ‑> tuple[ReportingStatusDirty, ...]
Expand source code
async def read_status_dirty(
    self, *, account_id: str, after: int = 0, limit: int = 100
) -> tuple[ReportingStatusDirty, ...]: ...
async def reemit(self, *, account_id: str, notification_id: str, now: datetime) ‑> int
Expand source code
async def reemit(self, *, account_id: str, notification_id: str, now: datetime) -> int: ...
async def status_checkpoint(self, *, account_id: str, projector_id: str) ‑> int
Expand source code
async def status_checkpoint(self, *, account_id: str, projector_id: str) -> int: ...
class ReportingNotificationSubscription (account_id: str,
subscriber_id: str,
principal_id: str,
url: str,
event_types: tuple[str, ...],
configuration_revision: str,
authorization_ref: str,
proof_of_control_ref: str,
signing_scope_id: str | None = None,
authentication: ReportingLegacyAuthentication | None = None,
*,
active: bool,
authorized: bool,
proof_valid: bool)
Expand source code
@dataclass(frozen=True)
class ReportingNotificationSubscription:
    """One *trusted* normalized active account notification configuration.

    The resolver must read these fields together, from authorized principal
    state and a successful account/subscriber/URL proof-of-control registration.
    References are mandatory and part of the fingerprint; they are not a
    substitute for performing those checks in the registration service.
    Never construct this from a webhook body or an unauthenticated sync request.
    """

    account_id: str
    subscriber_id: str
    principal_id: str
    url: str = field(repr=False)
    event_types: tuple[str, ...]
    configuration_revision: str
    authorization_ref: str
    proof_of_control_ref: str
    signing_scope_id: str | None = None
    authentication: ReportingLegacyAuthentication | None = field(default=None, repr=False)
    active: bool = field(kw_only=True)
    authorized: bool = field(kw_only=True)
    proof_valid: bool = field(kw_only=True)

    def __post_init__(self) -> None:
        try:
            principal_reference(self.account_id)
            canonical_consumer(self.principal_id)
            reporting_identifier(self.subscriber_id, maximum=64)
            for value in (
                self.configuration_revision,
                self.authorization_ref,
                self.proof_of_control_ref,
            ):
                reporting_identifier(value, maximum=255)
            if self.signing_scope_id is not None:
                principal_reference(self.signing_scope_id)
            if any(
                type(value) is not bool
                for value in (self.active, self.authorized, self.proof_valid)
            ):
                raise ValueError
            events = tuple(sorted(self.event_types))
            if not events or len(set(events)) != len(events):
                raise ValueError
            for event in events:
                reporting_identifier(event, maximum=255)
            object.__setattr__(self, "event_types", events)
            if (self.authentication is None) == (self.signing_scope_id is None):
                raise ValueError
            if (
                self.authentication is not None
                and type(self.authentication) is not ReportingLegacyAuthentication
            ):
                raise ValueError
            if type(self.url) is not str or not self.url.isprintable() or len(self.url) > 8192:
                raise ValueError
            parsed = httpx.URL(self.url)
            # Percent-encoding expands the canonical form, so the bound has to
            # hold on the value that is actually stored, sanitized for activity
            # and length-checked in SQL. Rejecting it here keeps a deterministic
            # failure at registration instead of quarantining every expansion.
            canonical = str(parsed)
            if (
                parsed.scheme != "https"
                or not parsed.host
                or parsed.userinfo
                or parsed.fragment
                or parsed.port not in (None, 443)
                or len(canonical) > 8192
            ):
                raise ValueError
            object.__setattr__(self, "url", canonical)
        except (ValueError, TypeError, httpx.InvalidURL):
            # Clear parser context before raising the closed configuration error below.
            pass
        else:
            return
        raise ReportingNotificationError("invalid_configuration")

    @property
    def auth_mode(self) -> str:
        return self.authentication.scheme if self.authentication is not None else "rfc9421"

    @property
    def fingerprint(self) -> str:
        return hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest()

    def matches(self, event: ReportingDomainEvent) -> bool:
        return (
            self.active
            and self.authorized
            and self.proof_valid
            and self.account_id == event.account_id
            and event.notification_type in self.event_types
            and (not event.consumer_namespace or self.principal_id == event.consumer_namespace)
        )

One trusted normalized active account notification configuration.

The resolver must read these fields together, from authorized principal state and a successful account/subscriber/URL proof-of-control registration. References are mandatory and part of the fingerprint; they are not a substitute for performing those checks in the registration service. Never construct this from a webhook body or an unauthenticated sync request.

Instance variables

var account_id : str
var active : bool
prop auth_mode : str
Expand source code
@property
def auth_mode(self) -> str:
    return self.authentication.scheme if self.authentication is not None else "rfc9421"
var authentication : ReportingLegacyAuthentication | None
var authorization_ref : str
var authorized : bool
var configuration_revision : str
var event_types : tuple[str, ...]
prop fingerprint : str
Expand source code
@property
def fingerprint(self) -> str:
    return hashlib.sha256(canonical_json_utf8_v1(asdict(self))).hexdigest()
var principal_id : str
var proof_of_control_ref : str
var proof_valid : bool
var signing_scope_id : str | None
var subscriber_id : str
var url : str

Methods

def matches(self,
event: ReportingDomainEvent) ‑> bool
Expand source code
def matches(self, event: ReportingDomainEvent) -> bool:
    return (
        self.active
        and self.authorized
        and self.proof_valid
        and self.account_id == event.account_id
        and event.notification_type in self.event_types
        and (not event.consumer_namespace or self.principal_id == event.consumer_namespace)
    )
class ReportingNotificationWorker (*,
outbox: ReportingNotificationOutbox,
subscriptions: ReportingSubscriptionResolver,
cipher: ReportingEnvelopeCipher,
signing: ReportingSigningResolver | None = None,
clock: Callable[[], datetime] | None = None,
lease_seconds: float = 60,
retry_seconds: float = 5,
activity: ReportingActivityStore | None = None,
delivery_window: ReportingDeliveryWindow | None = None)
Expand source code
class ReportingNotificationWorker:
    """At-least-once delivery with immutable bytes and per-attempt signing.

    A worker has no general-purpose sender injection. It constructs the exact
    SDK sender from trusted current key material, with the SDK-owned IP-pinned
    transport, public HTTPS/443 only, no rewrite hooks, and no redirects.

    Cancellation/process death intentionally leaves the lease to expire. Claims
    never exhaust a retry budget. A receiver must dedupe by authenticated sender
    and idempotency key: HTTP acceptance and our ACK cannot be one transaction.
    """

    def __init__(
        self,
        *,
        outbox: ReportingNotificationOutbox,
        subscriptions: ReportingSubscriptionResolver,
        cipher: ReportingEnvelopeCipher,
        signing: ReportingSigningResolver | None = None,
        clock: Callable[[], datetime] | None = None,
        lease_seconds: float = 60,
        retry_seconds: float = 5,
        activity: ReportingActivityStore | None = None,
        delivery_window: ReportingDeliveryWindow | None = None,
    ) -> None:
        if lease_seconds < 1 or retry_seconds <= 0:
            raise ValueError("positive retry and at least one second of lease are required")
        self.outbox, self.subscriptions, self.cipher, self.signing = (
            outbox,
            subscriptions,
            cipher,
            signing,
        )
        self._clock = clock or (lambda: datetime.now(timezone.utc))
        self.lease_seconds, self.retry_seconds = lease_seconds, retry_seconds
        if activity is not None and id(activity) != id(outbox):
            raise ReportingNotificationError("activity_requires_reporting_outbox")
        self.activity = activity
        self.delivery_window = delivery_window

    async def advertised_notifications(
        self,
        ledger: ReportingLedgerStore,
        *,
        account_id: str,
        ready_scope: ReportingDeliveryScope | None = None,
        activity_projector: ReportingActivityProjector | None = None,
    ) -> dict[str, str | bool]:
        """Check the installed chain at startup before publishing these fields.

        Merge the result into the producer's capability block while this worker
        is scheduled. Ready support additionally requires a retained Managed
        scope. Activity additionally requires the mounted durable projector;
        status notification projection belongs to the subsequent slice.
        """
        from adcp.reporting.outbox._capabilities import advertised_notifications

        return await advertised_notifications(
            self,
            ledger,
            account_id=account_id,
            ready_scope=ready_scope,
            activity_projector=activity_projector,
        )

    async def expand_one(self, *, account_id: str) -> bool:
        lease = await self.outbox.claim_expansion(
            account_id=account_id, now=self._clock(), lease_seconds=self.lease_seconds
        )
        if lease is None:
            return False
        try:
            event = decode_event(lease.event)
            if (
                event.account_id != lease.account_id
                or event.notification_id != lease.notification_id
                or event.consumer_namespace != lease.consumer_namespace
            ):
                raise ReportingNotificationError("invalid_event")
        except Exception:
            await self.outbox.finish_expansion(
                lease, now=self._clock(), state="quarantined", error_code="invalid_payload"
            )
            return True

        try:
            # One external snapshot. Its membership is fixed only when the
            # subsequent all-N insertion + complete checkpoint commits.
            resolved = await asyncio.wait_for(
                self.subscriptions.list_active(
                    account_id=account_id, notification_type=event.notification_type
                ),
                timeout=self.lease_seconds * 0.8,
            )
        except Exception:
            now = self._clock()
            await self.outbox.finish_expansion(
                lease,
                now=now,
                state="pending",
                error_code="subscription_unavailable",
                retry_at=now + timedelta(seconds=self.retry_seconds),
            )
            return True

        try:
            if not isinstance(resolved, (list, tuple)):
                raise ReportingNotificationError()
            snapshot = tuple(resolved)
            subscribers: set[str] = set()
            for subscription in snapshot:
                if (
                    type(subscription) is not ReportingNotificationSubscription
                    or subscription.account_id != account_id
                    or subscription.subscriber_id in subscribers
                ):
                    raise ReportingNotificationError()
                # Reconstruct the closed normalized type rather than trusting
                # an overridden fingerprint/matches property on a subclass.
                ReportingNotificationSubscription.__post_init__(subscription)
                subscribers.add(subscription.subscriber_id)
            deliveries = tuple(
                self.cipher.prepare(event, subscription, lease.emission_generation)
                for subscription in snapshot
                if subscription.matches(event)
            )
        except (ValueError, TypeError, AttributeError):
            await self.outbox.finish_expansion(
                lease, now=self._clock(), state="quarantined", error_code="invalid_configuration"
            )
            return True
        # Empty valid membership completes too. Failure here rolls back the
        # entire snapshot. A restarted worker may then observe new membership;
        # it can never combine committed rows from two snapshots.
        await self.outbox.complete_expansion(lease, deliveries, now=self._clock())
        return True

    async def deliver_one(self, *, account_id: str) -> bool:
        lease = await self.outbox.claim_delivery(
            account_id=account_id, now=self._clock(), lease_seconds=self.lease_seconds
        )
        if lease is None:
            return False
        observation = _HttpObservation()
        if self.delivery_window is not None:
            observation.retry_deadline, observation.retry_window_expired = (
                await self.delivery_window.inspect(lease, now=self._clock())
            )
        try:
            opened = self.cipher.open(lease.delivery)
        except (ReportingNotificationError, ValueError, TypeError):
            outcome = _Outcome("quarantined", "integrity_failure")
        else:
            try:
                if observation.retry_window_expired:
                    outcome = _Outcome("suppressed", "lease_expired")
                else:
                    outcome = await asyncio.wait_for(
                        self._attempt(lease, opened, observation), timeout=self.lease_seconds * 0.8
                    )
            except (TimeoutError, asyncio.TimeoutError):
                # Worker cancellation is not an observed HTTP timeout. A
                # reservation remains pending until a known result is ACKed.
                if observation.reservation is not None:
                    return True
                outcome = _Outcome("pending", "network")
        now = self._clock()
        retry_at = now + timedelta(seconds=self.retry_seconds)
        if observation.retry_deadline is not None:
            retry_at = min(retry_at, observation.retry_deadline)
        # A DB failure after HTTP acceptance is intentionally not converted to
        # success. Expiry/restart retries these exact protected body bytes/key.
        await self.outbox.finish_delivery(
            lease,
            now=now,
            state=outcome.state,
            error_code=outcome.error,
            retry_at=retry_at if outcome.state == "pending" else None,
        )
        return True

    async def _record_outcome(
        self, observation: _HttpObservation, outcome: ActivityOutcome
    ) -> None:
        if self.activity is not None and observation.reservation is not None:
            completed = await self.activity.complete_attempt(
                observation.reservation, outcome=outcome, now=self._clock()
            )
            if not completed:
                raise ReportingNotificationError("activity_completion_unconfirmed")

    async def _attempt(
        self, lease: DeliveryLease, opened: OpenedReportingDelivery, observation: _HttpObservation
    ) -> _Outcome:
        binding = lease.delivery.binding
        try:
            current = await self.subscriptions.get_active(
                account_id=binding.account_id,
                subscriber_id=binding.subscriber_id,
                notification_type=binding.notification_type,
            )
        except Exception:
            return _Outcome("pending", "subscription_unavailable")
        if type(current) is ReportingNotificationSubscription:
            try:
                ReportingNotificationSubscription.__post_init__(current)
            except (ValueError, TypeError, AttributeError):
                return _Outcome("suppressed", "subscription_changed")
        if (
            type(current) is not ReportingNotificationSubscription
            or not current.matches(opened.event)
            or current.subscriber_id != binding.subscriber_id
            or current.principal_id != binding.principal_id
            or current.fingerprint != binding.subscription_fingerprint
        ):
            return _Outcome("suppressed", "subscription_changed")
        try:
            sender = await self._sender(opened.subscription)
        except ScopePermanentlyUnknown:
            return _Outcome("quarantined", "permanent_scope")
        except ScopeTransientlyUnavailable:
            return _Outcome("pending", "signing_unavailable")
        try:
            if not await self.outbox.delivery_lease_current(lease, now=self._clock()):
                return _Outcome("pending", "lease_expired")
            try:
                # Invoke the concrete SDK seam. Adopter-provided sender objects,
                # subclasses, overridden send methods, clients and hooks never
                # cross this boundary. URL/DNS validation is inside this seam.
                request = ActivityRequest(opened.subscription.url, len(opened.prepared.body))
                callback_used = False

                async def current_fence() -> bool:
                    nonlocal callback_used
                    if callback_used:
                        raise PreparedWebhookAttemptExpiredError("prepared_attempt_expired")
                    callback_used = True
                    if not await self.outbox.delivery_lease_current(lease, now=self._clock()):
                        return False
                    if self.delivery_window is not None and self.activity is None:
                        raise ReportingNotificationError("activity_requires_reporting_outbox")
                    if self.activity is not None:
                        # Signing, URL/DNS preparation, and the final fence
                        # precede this transaction. No awaitable preparation
                        # remains between reservation and starting peer I/O.
                        try:
                            if self.delivery_window is not None:
                                (
                                    observation.reservation,
                                    observation.retry_deadline,
                                    observation.retry_window_expired,
                                ) = await self.delivery_window.reserve_attempt(
                                    self.activity, lease, request=request, now=self._clock()
                                )
                            else:
                                observation.reservation = await self.activity.reserve_attempt(
                                    lease, request=request, now=self._clock()
                                )
                        except ReportingNotificationError as error:
                            if error.code == "activity_lease_expired":
                                raise PreparedWebhookAttemptExpiredError(
                                    "prepared_attempt_expired"
                                ) from None
                            raise
                        if observation.reservation is None:
                            return False
                    observation.started_ns = time.monotonic_ns()
                    return True

                with protected_transport_logs():
                    result = await WebhookSender.send_prepared(
                        sender, opened.prepared, before_attempt=current_fence
                    )
            except PreparedWebhookAttemptExpiredError:
                return _Outcome(
                    "suppressed" if observation.retry_window_expired else "pending", "lease_expired"
                )
            except ReportingNotificationError:
                # Failed/unknown reservation commit means no HTTP and no ACK.
                raise
            except SSRFValidationError as error:
                return (
                    _Outcome("pending", "network")
                    if error.transient
                    else (_Outcome("quarantined", "invalid_configuration"))
                )
            except (httpx.TimeoutException, TimeoutError):
                await self._record_outcome(observation, ActivityOutcome("timeout"))
                return _Outcome("pending", "network")
            except (httpx.TransportError, OSError):
                await self._record_outcome(observation, ActivityOutcome("connection_error"))
                return _Outcome("pending", "network")
            except (ValueError, TypeError):
                if observation.reservation is not None:
                    raise ReportingNotificationError("activity_result_unknown") from None
                return _Outcome("quarantined", "invalid_payload")
            except Exception:
                # Signing backends can fail transiently; retain no exception
                # prose, traceback, headers, URL, or response/provider body.
                if observation.reservation is not None:
                    raise ReportingNotificationError("activity_result_unknown") from None
                return _Outcome("pending", "signing_unavailable")
            await self._record_outcome(
                observation,
                ActivityOutcome(
                    "success" if result.ok else "failed",
                    result.status_code,
                    max(0, (time.monotonic_ns() - observation.started_ns) // 1_000_000),
                ),
            )
            if result.ok:
                return _Outcome("complete")
            if result.status_code in {408, 425, 429} or 500 <= result.status_code < 600:
                return _Outcome("pending", "retryable_http")
            return _Outcome("quarantined", "permanent_http")
        finally:
            await sender.aclose()

    async def _sender(self, subscription: ReportingNotificationSubscription) -> WebhookSender:
        timeout = min(10.0, self.lease_seconds * 0.5)
        if subscription.authentication is not None:
            auth = subscription.authentication
            if subscription.signing_scope_id is not None:
                raise ScopePermanentlyUnknown from None
            if auth.scheme == "Bearer":
                return WebhookSender.from_bearer_token(
                    auth.credentials,
                    timeout_seconds=timeout,
                    allowed_destination_ports=frozenset({443}),
                )
            return WebhookSender.from_adcp_legacy_hmac(
                auth.credentials.encode("utf-8"),
                key_id="reporting-registration",
                timeout_seconds=timeout,
                allowed_destination_ports=frozenset({443}),
            )
        if self.signing is None or subscription.signing_scope_id is None:
            raise ScopePermanentlyUnknown from None
        try:
            material = await self.signing.resolve(
                account_id=subscription.account_id,
                principal_id=subscription.principal_id,
                signing_scope_id=subscription.signing_scope_id,
            )
        except (ScopePermanentlyUnknown, ScopeTransientlyUnavailable):
            raise
        except Exception:
            raise ScopeTransientlyUnavailable from None
        if type(material) is not ReportingSigningMaterial:
            raise ScopePermanentlyUnknown from None
        try:
            ReportingSigningMaterial.__post_init__(material)
            return WebhookSender(
                private_key=material.private_key,
                key_id=material.key_id,
                alg=material.algorithm,
                timeout_seconds=timeout,
                allowed_destination_ports=frozenset({443}),
            )
        except (ValueError, TypeError, AttributeError):
            raise ScopePermanentlyUnknown from None

At-least-once delivery with immutable bytes and per-attempt signing.

A worker has no general-purpose sender injection. It constructs the exact SDK sender from trusted current key material, with the SDK-owned IP-pinned transport, public HTTPS/443 only, no rewrite hooks, and no redirects.

Cancellation/process death intentionally leaves the lease to expire. Claims never exhaust a retry budget. A receiver must dedupe by authenticated sender and idempotency key: HTTP acceptance and our ACK cannot be one transaction.

Methods

async def advertised_notifications(self,
ledger: ReportingLedgerStore,
*,
account_id: str,
ready_scope: ReportingDeliveryScope | None = None,
activity_projector: ReportingActivityProjector | None = None) ‑> dict[str, str | bool]
Expand source code
async def advertised_notifications(
    self,
    ledger: ReportingLedgerStore,
    *,
    account_id: str,
    ready_scope: ReportingDeliveryScope | None = None,
    activity_projector: ReportingActivityProjector | None = None,
) -> dict[str, str | bool]:
    """Check the installed chain at startup before publishing these fields.

    Merge the result into the producer's capability block while this worker
    is scheduled. Ready support additionally requires a retained Managed
    scope. Activity additionally requires the mounted durable projector;
    status notification projection belongs to the subsequent slice.
    """
    from adcp.reporting.outbox._capabilities import advertised_notifications

    return await advertised_notifications(
        self,
        ledger,
        account_id=account_id,
        ready_scope=ready_scope,
        activity_projector=activity_projector,
    )

Check the installed chain at startup before publishing these fields.

Merge the result into the producer's capability block while this worker is scheduled. Ready support additionally requires a retained Managed scope. Activity additionally requires the mounted durable projector; status notification projection belongs to the subsequent slice.

async def deliver_one(self, *, account_id: str) ‑> bool
Expand source code
async def deliver_one(self, *, account_id: str) -> bool:
    lease = await self.outbox.claim_delivery(
        account_id=account_id, now=self._clock(), lease_seconds=self.lease_seconds
    )
    if lease is None:
        return False
    observation = _HttpObservation()
    if self.delivery_window is not None:
        observation.retry_deadline, observation.retry_window_expired = (
            await self.delivery_window.inspect(lease, now=self._clock())
        )
    try:
        opened = self.cipher.open(lease.delivery)
    except (ReportingNotificationError, ValueError, TypeError):
        outcome = _Outcome("quarantined", "integrity_failure")
    else:
        try:
            if observation.retry_window_expired:
                outcome = _Outcome("suppressed", "lease_expired")
            else:
                outcome = await asyncio.wait_for(
                    self._attempt(lease, opened, observation), timeout=self.lease_seconds * 0.8
                )
        except (TimeoutError, asyncio.TimeoutError):
            # Worker cancellation is not an observed HTTP timeout. A
            # reservation remains pending until a known result is ACKed.
            if observation.reservation is not None:
                return True
            outcome = _Outcome("pending", "network")
    now = self._clock()
    retry_at = now + timedelta(seconds=self.retry_seconds)
    if observation.retry_deadline is not None:
        retry_at = min(retry_at, observation.retry_deadline)
    # A DB failure after HTTP acceptance is intentionally not converted to
    # success. Expiry/restart retries these exact protected body bytes/key.
    await self.outbox.finish_delivery(
        lease,
        now=now,
        state=outcome.state,
        error_code=outcome.error,
        retry_at=retry_at if outcome.state == "pending" else None,
    )
    return True
async def expand_one(self, *, account_id: str) ‑> bool
Expand source code
async def expand_one(self, *, account_id: str) -> bool:
    lease = await self.outbox.claim_expansion(
        account_id=account_id, now=self._clock(), lease_seconds=self.lease_seconds
    )
    if lease is None:
        return False
    try:
        event = decode_event(lease.event)
        if (
            event.account_id != lease.account_id
            or event.notification_id != lease.notification_id
            or event.consumer_namespace != lease.consumer_namespace
        ):
            raise ReportingNotificationError("invalid_event")
    except Exception:
        await self.outbox.finish_expansion(
            lease, now=self._clock(), state="quarantined", error_code="invalid_payload"
        )
        return True

    try:
        # One external snapshot. Its membership is fixed only when the
        # subsequent all-N insertion + complete checkpoint commits.
        resolved = await asyncio.wait_for(
            self.subscriptions.list_active(
                account_id=account_id, notification_type=event.notification_type
            ),
            timeout=self.lease_seconds * 0.8,
        )
    except Exception:
        now = self._clock()
        await self.outbox.finish_expansion(
            lease,
            now=now,
            state="pending",
            error_code="subscription_unavailable",
            retry_at=now + timedelta(seconds=self.retry_seconds),
        )
        return True

    try:
        if not isinstance(resolved, (list, tuple)):
            raise ReportingNotificationError()
        snapshot = tuple(resolved)
        subscribers: set[str] = set()
        for subscription in snapshot:
            if (
                type(subscription) is not ReportingNotificationSubscription
                or subscription.account_id != account_id
                or subscription.subscriber_id in subscribers
            ):
                raise ReportingNotificationError()
            # Reconstruct the closed normalized type rather than trusting
            # an overridden fingerprint/matches property on a subclass.
            ReportingNotificationSubscription.__post_init__(subscription)
            subscribers.add(subscription.subscriber_id)
        deliveries = tuple(
            self.cipher.prepare(event, subscription, lease.emission_generation)
            for subscription in snapshot
            if subscription.matches(event)
        )
    except (ValueError, TypeError, AttributeError):
        await self.outbox.finish_expansion(
            lease, now=self._clock(), state="quarantined", error_code="invalid_configuration"
        )
        return True
    # Empty valid membership completes too. Failure here rolls back the
    # entire snapshot. A restarted worker may then observe new membership;
    # it can never combine committed rows from two snapshots.
    await self.outbox.complete_expansion(lease, deliveries, now=self._clock())
    return True
class ReportingSigningMaterial (private_key: PrivateKey,
key_id: str,
algorithm: str,
advertised_algorithms: frozenset[str])
Expand source code
@dataclass(frozen=True)
class ReportingSigningMaterial:
    """Trusted current key material, not an adopter-supplied HTTP sender."""

    private_key: PrivateKey = field(repr=False)
    key_id: str
    algorithm: str
    advertised_algorithms: frozenset[str]

    def __post_init__(self) -> None:
        object.__setattr__(self, "advertised_algorithms", frozenset(self.advertised_algorithms))
        key_matches = (
            self.algorithm == ALG_ED25519
            and isinstance(self.private_key, ed25519.Ed25519PrivateKey)
        ) or (
            self.algorithm == ALG_ES256
            and isinstance(self.private_key, ec.EllipticCurvePrivateKey)
            and isinstance(self.private_key.curve, ec.SECP256R1)
        )
        if (
            self.algorithm not in ALLOWED_ALGS
            or self.algorithm not in self.advertised_algorithms
            or not self.advertised_algorithms.issubset(ALLOWED_ALGS)
            or not key_matches
            or not self.key_id
        ):
            raise ReportingNotificationError("permanent_scope")

Trusted current key material, not an adopter-supplied HTTP sender.

Instance variables

var advertised_algorithms : frozenset[str]
var algorithm : str
var key_id : str
var private_key : cryptography.hazmat.primitives.asymmetric.ed25519.Ed25519PrivateKey | cryptography.hazmat.primitives.asymmetric.ec.EllipticCurvePrivateKey
class ReportingSigningResolver (*args, **kwargs)
Expand source code
class ReportingSigningResolver(Protocol):
    async def resolve(
        self, *, account_id: str, principal_id: str, signing_scope_id: str
    ) -> ReportingSigningMaterial: ...

Base class for protocol classes.

Protocol classes are defined as::

class Proto(Protocol):
    def meth(self) -> int:
        ...

Such classes are primarily used with static type checkers that recognize structural subtyping (static duck-typing).

For example::

class C:
    def meth(self) -> int:
        return 0

def func(x: Proto) -> int:
    return x.meth()

func(C())  # Passes static type check

See PEP 544 for details. Protocol classes decorated with @typing.runtime_checkable act as simple-minded runtime protocols that check only the presence of given attributes, ignoring their type signatures. Protocol classes can be generic, they are defined as::

class GenProto(Protocol[T]):
    def meth(self) -> T:
        ...

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def resolve(self, *, account_id: str, principal_id: str, signing_scope_id: str) ‑> ReportingSigningMaterial
Expand source code
async def resolve(
    self, *, account_id: str, principal_id: str, signing_scope_id: str
) -> ReportingSigningMaterial: ...
class ReportingStatusDirty (sequence: int,
scope: ReportingStatusScope,
reason: DirtyReason,
changed_at: datetime,
cause_id: str,
cause_generation: int,
before: ReportingStatusEvidence | None = None,
after: ReportingStatusEvidence | None = None)
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")

ReportingStatusDirty(sequence: 'int', scope: 'ReportingStatusScope', reason: 'DirtyReason', changed_at: 'datetime', cause_id: 'str', cause_generation: 'int', before: 'ReportingStatusEvidence | None' = None, after: 'ReportingStatusEvidence | None' = None)

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Instance variables

var after : ReportingStatusEvidence | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
var before : ReportingStatusEvidence | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
var cause_generation : int
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
var cause_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
var changed_at : datetime.datetime
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
var reason : Literal['configuration', 'obligation', 'revision', 'adjustment', 'readability', 'consumer_status', 'issue', 'destination', 'materialization', 'receipt', 'clock']
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
var scope : ReportingStatusScope
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
var sequence : int
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusDirty(_ClosedValue):
    sequence: int
    scope: ReportingStatusScope
    reason: DirtyReason
    changed_at: datetime
    cause_id: str
    cause_generation: int
    before: ReportingStatusEvidence | None = None
    after: ReportingStatusEvidence | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        object.__setattr__(self, "changed_at", aware_utc(self.changed_at))
        if self.sequence < 1 or self.cause_generation < 1:
            raise ReportingNotificationError("invalid_status_evidence")
class ReportingStatusEvidence (record_kind: "Literal['configuration', 'obligation', 'revision', 'adjustment', 'consumer_status', 'issue', 'destination_binding', 'obligation_delivery', 'materialization_attempt', 'materialization', 'materialization_check', 'revision_receipt', 'adjustment_receipt', 'clock']",
record_id: str,
record_version: int = 1,
readable: bool | None = None,
issue_state: "Literal['open', 'acknowledged', 'waived', 'resolved'] | None" = None,
opened_at: datetime | None = None,
retired_at: datetime | None = None,
supersedes_id: str | None = None,
activated_at: datetime | None = None,
deactivated_at: datetime | None = None,
automated_recovery_seconds: float | None = None,
status_retention_days: int | None = None)
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))

Replay evidence for one status mutation, containing no adopter prose.

Immutable records are referenced in their account/consumer namespace. The mutable inputs (readability and issue lifecycle) are snapshotted on both sides, so rapid reversals cannot disappear before a projector runs.

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Instance variables

var activated_at : datetime.datetime | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var automated_recovery_seconds : float | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var deactivated_at : datetime.datetime | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var issue_state : Literal['open', 'acknowledged', 'resolved', 'waived'] | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var opened_at : datetime.datetime | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var readable : bool | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var record_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var record_kind : Literal['configuration', 'obligation', 'revision', 'adjustment', 'consumer_status', 'issue', 'destination_binding', 'obligation_delivery', 'materialization_attempt', 'materialization', 'materialization_check', 'revision_receipt', 'adjustment_receipt', 'clock']
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var record_version : int
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var retired_at : datetime.datetime | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var status_retention_days : int | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
var supersedes_id : str | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusEvidence(_ClosedValue):
    """Replay evidence for one status mutation, containing no adopter prose.

    Immutable records are referenced in their account/consumer namespace. The
    mutable inputs (readability and issue lifecycle) are snapshotted on both
    sides, so rapid reversals cannot disappear before a projector runs.
    """

    record_kind: Literal[
        "configuration",
        "obligation",
        "revision",
        "adjustment",
        "consumer_status",
        "issue",
        "destination_binding",
        "obligation_delivery",
        "materialization_attempt",
        "materialization",
        "materialization_check",
        "revision_receipt",
        "adjustment_receipt",
        "clock",
    ]
    record_id: str
    record_version: int = 1
    readable: bool | None = None
    issue_state: Literal["open", "acknowledged", "waived", "resolved"] | None = None
    opened_at: datetime | None = None
    retired_at: datetime | None = None
    supersedes_id: str | None = None
    activated_at: datetime | None = None
    deactivated_at: datetime | None = None
    automated_recovery_seconds: float | None = None
    status_retention_days: int | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.record_id, maximum=255)
        if self.record_version < 1:
            raise ReportingNotificationError("invalid_status_evidence")
        for name in ("opened_at", "retired_at", "activated_at", "deactivated_at"):
            value = getattr(self, name)
            if value is not None:
                object.__setattr__(self, name, aware_utc(value))
class ReportingStatusNotificationLifecycle (*args, **kwargs)
Expand source code
class ReportingStatusNotificationLifecycle(Protocol):
    """Opt-in lifecycle; existing ledger and outbox implementations need no changes."""

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

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

    async def ready(self) -> bool: ...

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

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

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

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

    async def drain(self, *, max_turns: int = 100) -> int: ...

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

Opt-in lifecycle; existing ledger and outbox implementations need no changes.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def aclose(self) ‑> None
Expand source code
async def aclose(self) -> None: ...
async def baseline(self) ‑> None
Expand source code
async def baseline(self) -> None: ...
async def deliver_once(self, *, account_id: str) ‑> bool
Expand source code
async def deliver_once(self, *, account_id: str) -> bool: ...
async def drain(self, *, max_turns: int = 100) ‑> int
Expand source code
async def drain(self, *, max_turns: int = 100) -> int: ...
async def expand_once(self, *, account_id: str) ‑> bool
Expand source code
async def expand_once(self, *, account_id: str) -> bool: ...
async def migrate(self) ‑> None
Expand source code
async def migrate(self) -> None: ...
async def project_dirty_once(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def project_dirty_once(self, *, account_id: str) -> StatusTurn: ...
async def ready(self) ‑> bool
Expand source code
async def ready(self) -> bool: ...
async def sweep_due_once(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def sweep_due_once(self, *, account_id: str) -> StatusTurn: ...
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 ReportingStatusScope (account_id: str,
generation_key: ReportingConfigurationGenerationKey | None = None,
reporting_obligation_id: str | None = None,
consumer_id: str | None = None,
feed_purpose: FeedPurpose | None = None)
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusScope(_ClosedValue):
    """A typed projection target. Missing detail means invalidate the account.

    Legacy issue keys are opaque. They are never parsed to infer a generation,
    obligation, consumer, or health. Adopters may supply this optional scope to
    the concrete stores' issue methods without changing ReportingLedgerStore.
    """

    account_id: str
    generation_key: ReportingConfigurationGenerationKey | None = None
    reporting_obligation_id: str | None = None
    consumer_id: str | None = None
    feed_purpose: FeedPurpose | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        if self.consumer_id is not None:
            consumer_reference(self.consumer_id)
        if self.reporting_obligation_id is not None:
            reporting_identifier(self.reporting_obligation_id, maximum=255)
        if self.generation_key is not None:
            if self.generation_key.account_id != self.account_id or self.consumer_id not in {
                None,
                self.generation_key.consumer_id,
            }:
                raise ReportingNotificationError("invalid_status_scope")
            object.__setattr__(self, "consumer_id", self.generation_key.consumer_id)

    @property
    def checkpoint_key(self) -> tuple[str, str, str, int, str, str]:
        """Six independent, non-null columns. Scope kinds never share an ID space."""
        if self.generation_key is None:
            raise ReportingNotificationError("invalid_status_scope")
        return (
            self.account_id,
            self.consumer_id or "",
            self.generation_key.delivery_config_id,
            self.generation_key.delivery_config_version,
            "obligation" if self.reporting_obligation_id is not None else "configuration",
            self.reporting_obligation_id or "",
        )

    @classmethod
    def for_obligation(
        cls, obligation: ReportingObligationRecord, consumer_id: str | None = None
    ) -> ReportingStatusScope:
        # Existing ledger fields are strings; validating through the closed
        # adapter refuses a provider-supplied feed label rather than echoing it.
        return decode_status_scope(
            {
                "account_id": obligation.account_id,
                "generation_key": asdict(obligation.generation_key),
                "reporting_obligation_id": obligation.reporting_obligation_id,
                "consumer_id": obligation.consumer_id if consumer_id is None else consumer_id,
                "feed_purpose": obligation.feed_purpose,
            }
        )

A typed projection target. Missing detail means invalidate the account.

Legacy issue keys are opaque. They are never parsed to infer a generation, obligation, consumer, or health. Adopters may supply this optional scope to the concrete stores' issue methods without changing ReportingLedgerStore.

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Static methods

def for_obligation(obligation: ReportingObligationRecord, consumer_id: str | None = None) ‑> ReportingStatusScope

Instance variables

var account_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusScope(_ClosedValue):
    """A typed projection target. Missing detail means invalidate the account.

    Legacy issue keys are opaque. They are never parsed to infer a generation,
    obligation, consumer, or health. Adopters may supply this optional scope to
    the concrete stores' issue methods without changing ReportingLedgerStore.
    """

    account_id: str
    generation_key: ReportingConfigurationGenerationKey | None = None
    reporting_obligation_id: str | None = None
    consumer_id: str | None = None
    feed_purpose: FeedPurpose | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        if self.consumer_id is not None:
            consumer_reference(self.consumer_id)
        if self.reporting_obligation_id is not None:
            reporting_identifier(self.reporting_obligation_id, maximum=255)
        if self.generation_key is not None:
            if self.generation_key.account_id != self.account_id or self.consumer_id not in {
                None,
                self.generation_key.consumer_id,
            }:
                raise ReportingNotificationError("invalid_status_scope")
            object.__setattr__(self, "consumer_id", self.generation_key.consumer_id)

    @property
    def checkpoint_key(self) -> tuple[str, str, str, int, str, str]:
        """Six independent, non-null columns. Scope kinds never share an ID space."""
        if self.generation_key is None:
            raise ReportingNotificationError("invalid_status_scope")
        return (
            self.account_id,
            self.consumer_id or "",
            self.generation_key.delivery_config_id,
            self.generation_key.delivery_config_version,
            "obligation" if self.reporting_obligation_id is not None else "configuration",
            self.reporting_obligation_id or "",
        )

    @classmethod
    def for_obligation(
        cls, obligation: ReportingObligationRecord, consumer_id: str | None = None
    ) -> ReportingStatusScope:
        # Existing ledger fields are strings; validating through the closed
        # adapter refuses a provider-supplied feed label rather than echoing it.
        return decode_status_scope(
            {
                "account_id": obligation.account_id,
                "generation_key": asdict(obligation.generation_key),
                "reporting_obligation_id": obligation.reporting_obligation_id,
                "consumer_id": obligation.consumer_id if consumer_id is None else consumer_id,
                "feed_purpose": obligation.feed_purpose,
            }
        )
prop checkpoint_key : tuple[str, str, str, int, str, str]
Expand source code
@property
def checkpoint_key(self) -> tuple[str, str, str, int, str, str]:
    """Six independent, non-null columns. Scope kinds never share an ID space."""
    if self.generation_key is None:
        raise ReportingNotificationError("invalid_status_scope")
    return (
        self.account_id,
        self.consumer_id or "",
        self.generation_key.delivery_config_id,
        self.generation_key.delivery_config_version,
        "obligation" if self.reporting_obligation_id is not None else "configuration",
        self.reporting_obligation_id or "",
    )

Six independent, non-null columns. Scope kinds never share an ID space.

var consumer_id : str | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusScope(_ClosedValue):
    """A typed projection target. Missing detail means invalidate the account.

    Legacy issue keys are opaque. They are never parsed to infer a generation,
    obligation, consumer, or health. Adopters may supply this optional scope to
    the concrete stores' issue methods without changing ReportingLedgerStore.
    """

    account_id: str
    generation_key: ReportingConfigurationGenerationKey | None = None
    reporting_obligation_id: str | None = None
    consumer_id: str | None = None
    feed_purpose: FeedPurpose | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        if self.consumer_id is not None:
            consumer_reference(self.consumer_id)
        if self.reporting_obligation_id is not None:
            reporting_identifier(self.reporting_obligation_id, maximum=255)
        if self.generation_key is not None:
            if self.generation_key.account_id != self.account_id or self.consumer_id not in {
                None,
                self.generation_key.consumer_id,
            }:
                raise ReportingNotificationError("invalid_status_scope")
            object.__setattr__(self, "consumer_id", self.generation_key.consumer_id)

    @property
    def checkpoint_key(self) -> tuple[str, str, str, int, str, str]:
        """Six independent, non-null columns. Scope kinds never share an ID space."""
        if self.generation_key is None:
            raise ReportingNotificationError("invalid_status_scope")
        return (
            self.account_id,
            self.consumer_id or "",
            self.generation_key.delivery_config_id,
            self.generation_key.delivery_config_version,
            "obligation" if self.reporting_obligation_id is not None else "configuration",
            self.reporting_obligation_id or "",
        )

    @classmethod
    def for_obligation(
        cls, obligation: ReportingObligationRecord, consumer_id: str | None = None
    ) -> ReportingStatusScope:
        # Existing ledger fields are strings; validating through the closed
        # adapter refuses a provider-supplied feed label rather than echoing it.
        return decode_status_scope(
            {
                "account_id": obligation.account_id,
                "generation_key": asdict(obligation.generation_key),
                "reporting_obligation_id": obligation.reporting_obligation_id,
                "consumer_id": obligation.consumer_id if consumer_id is None else consumer_id,
                "feed_purpose": obligation.feed_purpose,
            }
        )
var feed_purpose : Literal['pacing', 'analytics', 'billing'] | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusScope(_ClosedValue):
    """A typed projection target. Missing detail means invalidate the account.

    Legacy issue keys are opaque. They are never parsed to infer a generation,
    obligation, consumer, or health. Adopters may supply this optional scope to
    the concrete stores' issue methods without changing ReportingLedgerStore.
    """

    account_id: str
    generation_key: ReportingConfigurationGenerationKey | None = None
    reporting_obligation_id: str | None = None
    consumer_id: str | None = None
    feed_purpose: FeedPurpose | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        if self.consumer_id is not None:
            consumer_reference(self.consumer_id)
        if self.reporting_obligation_id is not None:
            reporting_identifier(self.reporting_obligation_id, maximum=255)
        if self.generation_key is not None:
            if self.generation_key.account_id != self.account_id or self.consumer_id not in {
                None,
                self.generation_key.consumer_id,
            }:
                raise ReportingNotificationError("invalid_status_scope")
            object.__setattr__(self, "consumer_id", self.generation_key.consumer_id)

    @property
    def checkpoint_key(self) -> tuple[str, str, str, int, str, str]:
        """Six independent, non-null columns. Scope kinds never share an ID space."""
        if self.generation_key is None:
            raise ReportingNotificationError("invalid_status_scope")
        return (
            self.account_id,
            self.consumer_id or "",
            self.generation_key.delivery_config_id,
            self.generation_key.delivery_config_version,
            "obligation" if self.reporting_obligation_id is not None else "configuration",
            self.reporting_obligation_id or "",
        )

    @classmethod
    def for_obligation(
        cls, obligation: ReportingObligationRecord, consumer_id: str | None = None
    ) -> ReportingStatusScope:
        # Existing ledger fields are strings; validating through the closed
        # adapter refuses a provider-supplied feed label rather than echoing it.
        return decode_status_scope(
            {
                "account_id": obligation.account_id,
                "generation_key": asdict(obligation.generation_key),
                "reporting_obligation_id": obligation.reporting_obligation_id,
                "consumer_id": obligation.consumer_id if consumer_id is None else consumer_id,
                "feed_purpose": obligation.feed_purpose,
            }
        )
var generation_key : ReportingConfigurationGenerationKey | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusScope(_ClosedValue):
    """A typed projection target. Missing detail means invalidate the account.

    Legacy issue keys are opaque. They are never parsed to infer a generation,
    obligation, consumer, or health. Adopters may supply this optional scope to
    the concrete stores' issue methods without changing ReportingLedgerStore.
    """

    account_id: str
    generation_key: ReportingConfigurationGenerationKey | None = None
    reporting_obligation_id: str | None = None
    consumer_id: str | None = None
    feed_purpose: FeedPurpose | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        if self.consumer_id is not None:
            consumer_reference(self.consumer_id)
        if self.reporting_obligation_id is not None:
            reporting_identifier(self.reporting_obligation_id, maximum=255)
        if self.generation_key is not None:
            if self.generation_key.account_id != self.account_id or self.consumer_id not in {
                None,
                self.generation_key.consumer_id,
            }:
                raise ReportingNotificationError("invalid_status_scope")
            object.__setattr__(self, "consumer_id", self.generation_key.consumer_id)

    @property
    def checkpoint_key(self) -> tuple[str, str, str, int, str, str]:
        """Six independent, non-null columns. Scope kinds never share an ID space."""
        if self.generation_key is None:
            raise ReportingNotificationError("invalid_status_scope")
        return (
            self.account_id,
            self.consumer_id or "",
            self.generation_key.delivery_config_id,
            self.generation_key.delivery_config_version,
            "obligation" if self.reporting_obligation_id is not None else "configuration",
            self.reporting_obligation_id or "",
        )

    @classmethod
    def for_obligation(
        cls, obligation: ReportingObligationRecord, consumer_id: str | None = None
    ) -> ReportingStatusScope:
        # Existing ledger fields are strings; validating through the closed
        # adapter refuses a provider-supplied feed label rather than echoing it.
        return decode_status_scope(
            {
                "account_id": obligation.account_id,
                "generation_key": asdict(obligation.generation_key),
                "reporting_obligation_id": obligation.reporting_obligation_id,
                "consumer_id": obligation.consumer_id if consumer_id is None else consumer_id,
                "feed_purpose": obligation.feed_purpose,
            }
        )
var reporting_obligation_id : str | None
Expand source code
@dataclass(frozen=True, slots=True)
class ReportingStatusScope(_ClosedValue):
    """A typed projection target. Missing detail means invalidate the account.

    Legacy issue keys are opaque. They are never parsed to infer a generation,
    obligation, consumer, or health. Adopters may supply this optional scope to
    the concrete stores' issue methods without changing ReportingLedgerStore.
    """

    account_id: str
    generation_key: ReportingConfigurationGenerationKey | None = None
    reporting_obligation_id: str | None = None
    consumer_id: str | None = None
    feed_purpose: FeedPurpose | None = None

    def __post_init__(self) -> None:
        _freeze_fields(self)
        principal_reference(self.account_id)
        if self.consumer_id is not None:
            consumer_reference(self.consumer_id)
        if self.reporting_obligation_id is not None:
            reporting_identifier(self.reporting_obligation_id, maximum=255)
        if self.generation_key is not None:
            if self.generation_key.account_id != self.account_id or self.consumer_id not in {
                None,
                self.generation_key.consumer_id,
            }:
                raise ReportingNotificationError("invalid_status_scope")
            object.__setattr__(self, "consumer_id", self.generation_key.consumer_id)

    @property
    def checkpoint_key(self) -> tuple[str, str, str, int, str, str]:
        """Six independent, non-null columns. Scope kinds never share an ID space."""
        if self.generation_key is None:
            raise ReportingNotificationError("invalid_status_scope")
        return (
            self.account_id,
            self.consumer_id or "",
            self.generation_key.delivery_config_id,
            self.generation_key.delivery_config_version,
            "obligation" if self.reporting_obligation_id is not None else "configuration",
            self.reporting_obligation_id or "",
        )

    @classmethod
    def for_obligation(
        cls, obligation: ReportingObligationRecord, consumer_id: str | None = None
    ) -> ReportingStatusScope:
        # Existing ledger fields are strings; validating through the closed
        # adapter refuses a provider-supplied feed label rather than echoing it.
        return decode_status_scope(
            {
                "account_id": obligation.account_id,
                "generation_key": asdict(obligation.generation_key),
                "reporting_obligation_id": obligation.reporting_obligation_id,
                "consumer_id": obligation.consumer_id if consumer_id is None else consumer_id,
                "feed_purpose": obligation.feed_purpose,
            }
        )
class ReportingStatusService (support: ReportingStatusSupport)
Expand source code
class ReportingStatusService:
    """Deterministic single turns for the adopter's supervised scheduler.

    The caller owns scheduling and the database pool. ``drain`` processes work
    due now; future HTTP retries and deadlines remain durable. Stop scheduling,
    await in-flight turns, drain, close this service, then close the owned pool.
    No background task or second HTTP client is created by this composition.
    """

    def __init__(self, support: ReportingStatusSupport) -> None:
        self.support = support
        self._closed = False

    def _open(self, account_id: str | None = None) -> None:
        if self._closed:
            raise ReportingNotificationError("status_service_closed")
        if account_id is not None and account_id not in self.support.account_ids:
            raise ReportingNotificationError("status_account_unavailable")

    async def migrate(self) -> None:
        self._open()
        await self.support.store.create_schema()

    async def baseline(self) -> None:
        self._open()
        for account_id in sorted(set(self.support.account_ids)):
            await self.support.store.baseline(account_id=account_id)

    async def ready(self) -> bool:
        return not self._closed and await self.support.durable()

    async def project_dirty_once(self, *, account_id: str) -> StatusTurn:
        self._open(account_id)
        if self.support.projector is None:
            raise ReportingNotificationError("status_projector_unavailable")
        return await self.support.projector.run_once(account_id=account_id)

    async def rebuild_selector_once(self) -> StatusTurn:
        """Indexed C cutover, including retained accounts outside current registrations."""
        self._open()
        store = self.support.store
        if not isinstance(store, StatusSelectorRebuildStore):
            return StatusTurn(False)  # Existing custom lifecycle protocols stay compatible.
        return await store.rebuild_one()

    async def sweep_due_once(self, *, account_id: str) -> StatusTurn:
        self._open(account_id)
        if self.support.sweeper is None:
            raise ReportingNotificationError("status_sweeper_unavailable")
        return await self.support.sweeper.run_once(account_id=account_id)

    async def expand_once(self, *, account_id: str) -> bool:
        self._open(account_id)
        return await self.support.worker.expand_one(account_id=account_id)

    async def deliver_once(self, *, account_id: str) -> bool:
        self._open(account_id)
        return await self.support.worker.deliver_one(account_id=account_id)

    async def drain(self, *, max_turns: int = 100) -> int:
        self._open()
        if max_turns < 1:
            raise ValueError("max_turns must be positive")
        for turn in range(max_turns):
            worked = (await self.rebuild_selector_once()).did_work
            for account_id in sorted(set(self.support.account_ids)):
                worked = (await self.project_dirty_once(account_id=account_id)).did_work or worked
                worked = (await self.sweep_due_once(account_id=account_id)).did_work or worked
                worked = await self.expand_once(account_id=account_id) or worked
                worked = await self.deliver_once(account_id=account_id) or worked
            if not worked:
                return turn
        raise ReportingNotificationError("status_drain_limit")

    async def aclose(self) -> None:
        self._closed = True

Deterministic single turns for the adopter's supervised scheduler.

The caller owns scheduling and the database pool. drain processes work due now; future HTTP retries and deadlines remain durable. Stop scheduling, await in-flight turns, drain, close this service, then close the owned pool. No background task or second HTTP client is created by this composition.

Methods

async def aclose(self) ‑> None
Expand source code
async def aclose(self) -> None:
    self._closed = True
async def baseline(self) ‑> None
Expand source code
async def baseline(self) -> None:
    self._open()
    for account_id in sorted(set(self.support.account_ids)):
        await self.support.store.baseline(account_id=account_id)
async def deliver_once(self, *, account_id: str) ‑> bool
Expand source code
async def deliver_once(self, *, account_id: str) -> bool:
    self._open(account_id)
    return await self.support.worker.deliver_one(account_id=account_id)
async def drain(self, *, max_turns: int = 100) ‑> int
Expand source code
async def drain(self, *, max_turns: int = 100) -> int:
    self._open()
    if max_turns < 1:
        raise ValueError("max_turns must be positive")
    for turn in range(max_turns):
        worked = (await self.rebuild_selector_once()).did_work
        for account_id in sorted(set(self.support.account_ids)):
            worked = (await self.project_dirty_once(account_id=account_id)).did_work or worked
            worked = (await self.sweep_due_once(account_id=account_id)).did_work or worked
            worked = await self.expand_once(account_id=account_id) or worked
            worked = await self.deliver_once(account_id=account_id) or worked
        if not worked:
            return turn
    raise ReportingNotificationError("status_drain_limit")
async def expand_once(self, *, account_id: str) ‑> bool
Expand source code
async def expand_once(self, *, account_id: str) -> bool:
    self._open(account_id)
    return await self.support.worker.expand_one(account_id=account_id)
async def migrate(self) ‑> None
Expand source code
async def migrate(self) -> None:
    self._open()
    await self.support.store.create_schema()
async def project_dirty_once(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def project_dirty_once(self, *, account_id: str) -> StatusTurn:
    self._open(account_id)
    if self.support.projector is None:
        raise ReportingNotificationError("status_projector_unavailable")
    return await self.support.projector.run_once(account_id=account_id)
async def ready(self) ‑> bool
Expand source code
async def ready(self) -> bool:
    return not self._closed and await self.support.durable()
async def rebuild_selector_once(self) ‑> StatusTurn
Expand source code
async def rebuild_selector_once(self) -> StatusTurn:
    """Indexed C cutover, including retained accounts outside current registrations."""
    self._open()
    store = self.support.store
    if not isinstance(store, StatusSelectorRebuildStore):
        return StatusTurn(False)  # Existing custom lifecycle protocols stay compatible.
    return await store.rebuild_one()

Indexed C cutover, including retained accounts outside current registrations.

async def sweep_due_once(self, *, account_id: str) ‑> StatusTurn
Expand source code
async def sweep_due_once(self, *, account_id: str) -> StatusTurn:
    self._open(account_id)
    if self.support.sweeper is None:
        raise ReportingNotificationError("status_sweeper_unavailable")
    return await self.support.sweeper.run_once(account_id=account_id)
class ReportingStatusSupport (store: StatusNotificationStore,
worker: ReportingNotificationWorker,
projector: ReportingStatusProjector | None = None,
sweeper: ReportingStatusSweeper | None = None,
scheduled: bool = False,
account_ids: tuple[str, ...] = (),
account_surface_complete: bool = False,
handler: ADCPHandler[Any] | None = None,
readiness_subscriptions: tuple[ReportingNotificationSubscription, ...] = ())
Expand source code
@dataclass(frozen=True)
class ReportingStatusSupport:
    """Mount only while the adopter schedules projector, sweep and HTTP turns.

    ``account_ids`` is the server's explicit supported account set, not an
    account-listing service or request-body identity. Global static claims
    require every one to be baselined. The protocol capability request has no
    account selector: callers must declare this enumeration complete for the
    authenticated or deployment surface, never select it from request context.
    """

    store: StatusNotificationStore
    worker: ReportingNotificationWorker
    projector: ReportingStatusProjector | None = None
    sweeper: ReportingStatusSweeper | None = None
    scheduled: bool = False
    account_ids: tuple[str, ...] = ()
    account_surface_complete: bool = False
    handler: ADCPHandler[Any] | None = None
    readiness_subscriptions: tuple[ReportingNotificationSubscription, ...] = ()

    async def queue_durable(self) -> bool:
        from adcp.reporting.outbox.status_memory import InMemoryStatusNotificationStore

        if (
            not self.scheduled
            or self.projector is None
            or self.sweeper is None
            or type(self.worker) is not ReportingNotificationWorker
            or type(self.worker.cipher) is not ReportingEnvelopeCipher
            or type(self.projector) is not ReportingStatusProjector
            or type(self.sweeper) is not ReportingStatusSweeper
            or self.projector.store is not self.store
            or self.sweeper.store is not self.store
        ):
            return False
        if type(self.store) is InMemoryStatusNotificationStore:
            return False
        from adcp.reporting.outbox.status_pg import (
            PgReportingStatusOutbox,
            PgStatusNotificationStore,
        )
        from adcp.reporting.outbox.status_schema import validate_status_schema

        if (
            type(self.store) is not PgStatusNotificationStore
            or not isinstance(self.store, PgStatusNotificationStore)
            or type(self.worker.outbox) is not PgReportingStatusOutbox
            or self.worker.outbox is not self.store.outbox
            or self.store.ledger._pool is not self.store.outbox._pool
            or not self.store.ledger._notifications_enabled
            or (self.worker.activity is not None and self.worker.activity is not self.worker.outbox)
        ):
            return False
        try:
            async with self.store.ledger._pool.connection() as connection:
                await validate_status_schema(connection, activity=self.worker.activity is not None)
        except ReportingNotificationError:
            return False
        return True

    async def account_ready(self, *, account_id: str) -> bool:
        if not await self.queue_durable() or not await self.store.baseline_ready(
            account_id=account_id
        ):
            return False
        from adcp.reporting.outbox.status_pg import PgStatusNotificationStore

        assert isinstance(self.store, PgStatusNotificationStore)
        async with self.store.ledger._pool.connection() as connection:
            managed = await (
                await connection.execute(
                    "SELECT EXISTS(SELECT 1 FROM reporting_reconciliation_records"
                    " WHERE account_id=%s AND record_kind='destination_binding')",
                    (account_id,),
                )
            ).fetchone()
            if managed is not None and managed[0]:
                return False  # Managed expiry/reconciliation health ships with D.
        subscriptions = await self.worker.subscriptions.list_active(
            account_id=account_id, notification_type="reporting.status_changed"
        )
        probes = tuple(s for s in self.readiness_subscriptions if s.account_id == account_id)
        if not probes:
            return False
        identities: set[str] = set()
        for subscription in subscriptions:
            if (
                type(subscription) is not ReportingNotificationSubscription
                or subscription.account_id != account_id
                or subscription.subscriber_id in identities
                or "reporting.status_changed" not in subscription.event_types
                or not (
                    subscription.active and subscription.authorized and subscription.proof_valid
                )
            ):
                raise ReportingNotificationError("status_subscription_unready")
            ReportingNotificationSubscription.__post_init__(subscription)
            identities.add(subscription.subscriber_id)
            sender = await self.worker._sender(subscription)
            await sender.aclose()
        for subscription in probes:
            if (
                type(subscription) is not ReportingNotificationSubscription
                or "reporting.status_changed" not in subscription.event_types
                or not (
                    subscription.active and subscription.authorized and subscription.proof_valid
                )
            ):
                raise ReportingNotificationError("status_subscription_unready")
            ReportingNotificationSubscription.__post_init__(subscription)
            sender = await self.worker._sender(subscription)
            await sender.aclose()
        return True

    async def durable(self) -> bool:
        from adcp.server.mcp_tools import get_tools_for_handler

        if not self.account_surface_complete or not self.account_ids or self.handler is None:
            return False
        projection = getattr(self.handler, "reporting_status_handler", None)
        if (
            type(projection) is not ReportingStatusHandler
            or projection._store is not getattr(self.store, "ledger", None)
            or projection._escalation != getattr(self.store, "escalation", None)
            or not projection._consumer_status_enabled
            or "get_reporting_status"
            not in {t["name"] for t in get_tools_for_handler(self.handler)}
        ):
            return False
        return all([await self.account_ready(account_id=a) for a in sorted(set(self.account_ids))])

    async def advertised_notifications(self) -> dict[str, str]:
        """An explicit global claim, after checking the entire trusted surface."""
        return (
            {
                "status_task": "get_reporting_status",
                "status_notification": "reporting.status_changed",
            }
            if await self.durable()
            else {}
        )

Mount only while the adopter schedules projector, sweep and HTTP turns.

account_ids is the server's explicit supported account set, not an account-listing service or request-body identity. Global static claims require every one to be baselined. The protocol capability request has no account selector: callers must declare this enumeration complete for the authenticated or deployment surface, never select it from request context.

Instance variables

var account_ids : tuple[str, ...]
var account_surface_complete : bool
var handler : ADCPHandler[typing.Any] | None
var projector : ReportingStatusProjector | None
var readiness_subscriptions : tuple[ReportingNotificationSubscription, ...]
var scheduled : bool
var store : StatusNotificationStore
var sweeper : ReportingStatusSweeper | None
var worker : ReportingNotificationWorker

Methods

async def account_ready(self, *, account_id: str) ‑> bool
Expand source code
async def account_ready(self, *, account_id: str) -> bool:
    if not await self.queue_durable() or not await self.store.baseline_ready(
        account_id=account_id
    ):
        return False
    from adcp.reporting.outbox.status_pg import PgStatusNotificationStore

    assert isinstance(self.store, PgStatusNotificationStore)
    async with self.store.ledger._pool.connection() as connection:
        managed = await (
            await connection.execute(
                "SELECT EXISTS(SELECT 1 FROM reporting_reconciliation_records"
                " WHERE account_id=%s AND record_kind='destination_binding')",
                (account_id,),
            )
        ).fetchone()
        if managed is not None and managed[0]:
            return False  # Managed expiry/reconciliation health ships with D.
    subscriptions = await self.worker.subscriptions.list_active(
        account_id=account_id, notification_type="reporting.status_changed"
    )
    probes = tuple(s for s in self.readiness_subscriptions if s.account_id == account_id)
    if not probes:
        return False
    identities: set[str] = set()
    for subscription in subscriptions:
        if (
            type(subscription) is not ReportingNotificationSubscription
            or subscription.account_id != account_id
            or subscription.subscriber_id in identities
            or "reporting.status_changed" not in subscription.event_types
            or not (
                subscription.active and subscription.authorized and subscription.proof_valid
            )
        ):
            raise ReportingNotificationError("status_subscription_unready")
        ReportingNotificationSubscription.__post_init__(subscription)
        identities.add(subscription.subscriber_id)
        sender = await self.worker._sender(subscription)
        await sender.aclose()
    for subscription in probes:
        if (
            type(subscription) is not ReportingNotificationSubscription
            or "reporting.status_changed" not in subscription.event_types
            or not (
                subscription.active and subscription.authorized and subscription.proof_valid
            )
        ):
            raise ReportingNotificationError("status_subscription_unready")
        ReportingNotificationSubscription.__post_init__(subscription)
        sender = await self.worker._sender(subscription)
        await sender.aclose()
    return True
async def advertised_notifications(self) ‑> dict[str, str]
Expand source code
async def advertised_notifications(self) -> dict[str, str]:
    """An explicit global claim, after checking the entire trusted surface."""
    return (
        {
            "status_task": "get_reporting_status",
            "status_notification": "reporting.status_changed",
        }
        if await self.durable()
        else {}
    )

An explicit global claim, after checking the entire trusted surface.

async def durable(self) ‑> bool
Expand source code
async def durable(self) -> bool:
    from adcp.server.mcp_tools import get_tools_for_handler

    if not self.account_surface_complete or not self.account_ids or self.handler is None:
        return False
    projection = getattr(self.handler, "reporting_status_handler", None)
    if (
        type(projection) is not ReportingStatusHandler
        or projection._store is not getattr(self.store, "ledger", None)
        or projection._escalation != getattr(self.store, "escalation", None)
        or not projection._consumer_status_enabled
        or "get_reporting_status"
        not in {t["name"] for t in get_tools_for_handler(self.handler)}
    ):
        return False
    return all([await self.account_ready(account_id=a) for a in sorted(set(self.account_ids))])
async def queue_durable(self) ‑> bool
Expand source code
async def queue_durable(self) -> bool:
    from adcp.reporting.outbox.status_memory import InMemoryStatusNotificationStore

    if (
        not self.scheduled
        or self.projector is None
        or self.sweeper is None
        or type(self.worker) is not ReportingNotificationWorker
        or type(self.worker.cipher) is not ReportingEnvelopeCipher
        or type(self.projector) is not ReportingStatusProjector
        or type(self.sweeper) is not ReportingStatusSweeper
        or self.projector.store is not self.store
        or self.sweeper.store is not self.store
    ):
        return False
    if type(self.store) is InMemoryStatusNotificationStore:
        return False
    from adcp.reporting.outbox.status_pg import (
        PgReportingStatusOutbox,
        PgStatusNotificationStore,
    )
    from adcp.reporting.outbox.status_schema import validate_status_schema

    if (
        type(self.store) is not PgStatusNotificationStore
        or not isinstance(self.store, PgStatusNotificationStore)
        or type(self.worker.outbox) is not PgReportingStatusOutbox
        or self.worker.outbox is not self.store.outbox
        or self.store.ledger._pool is not self.store.outbox._pool
        or not self.store.ledger._notifications_enabled
        or (self.worker.activity is not None and self.worker.activity is not self.worker.outbox)
    ):
        return False
    try:
        async with self.store.ledger._pool.connection() as connection:
            await validate_status_schema(connection, activity=self.worker.activity is not None)
    except ReportingNotificationError:
        return False
    return True
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 ReportingSubscriptionResolver (*args, **kwargs)
Expand source code
class ReportingSubscriptionResolver(Protocol):
    """External state: active-at-expansion, never claimed transactionally atomic.

    Each result must come from one normalized trusted account configuration,
    including current authorization and proof validity. Fanout reads once; the
    worker repeats exact resolution immediately before *every* HTTP attempt.
    """

    async def list_active(
        self, *, account_id: str, notification_type: str
    ) -> tuple[ReportingNotificationSubscription, ...]: ...

    async def get_active(
        self, *, account_id: str, subscriber_id: str, notification_type: str
    ) -> ReportingNotificationSubscription | None: ...

External state: active-at-expansion, never claimed transactionally atomic.

Each result must come from one normalized trusted account configuration, including current authorization and proof validity. Fanout reads once; the worker repeats exact resolution immediately before every HTTP attempt.

Ancestors

  • typing.Protocol
  • typing.Generic

Methods

async def get_active(self, *, account_id: str, subscriber_id: str, notification_type: str) ‑> ReportingNotificationSubscription | None
Expand source code
async def get_active(
    self, *, account_id: str, subscriber_id: str, notification_type: str
) -> ReportingNotificationSubscription | None: ...
async def list_active(self, *, account_id: str, notification_type: str) ‑> tuple[ReportingNotificationSubscription, ...]
Expand source code
async def list_active(
    self, *, account_id: str, notification_type: str
) -> tuple[ReportingNotificationSubscription, ...]: ...
class RevisionPublished (reporting_revision_id: str,
consumer_id: str,
finality: ReportingFinality,
supersedes_reporting_revision_id: str | None = None,
*,
kind: "Literal['revision_published']" = 'revision_published')
Expand source code
@dataclass(frozen=True, slots=True)
class RevisionPublished(_ClosedValue):
    reporting_revision_id: str
    consumer_id: str
    finality: ReportingFinality
    supersedes_reporting_revision_id: str | None = None
    kind: Literal["revision_published"] = field(default="revision_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_revision_id, maximum=255)
        consumer_reference(self.consumer_id)
        if self.supersedes_reporting_revision_id is not None:
            reporting_identifier(self.supersedes_reporting_revision_id, maximum=255)

RevisionPublished(reporting_revision_id: 'str', consumer_id: 'str', finality: 'ReportingFinality', supersedes_reporting_revision_id: 'str | None' = None, *, kind: "Literal['revision_published']" = 'revision_published')

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Instance variables

var consumer_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class RevisionPublished(_ClosedValue):
    reporting_revision_id: str
    consumer_id: str
    finality: ReportingFinality
    supersedes_reporting_revision_id: str | None = None
    kind: Literal["revision_published"] = field(default="revision_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_revision_id, maximum=255)
        consumer_reference(self.consumer_id)
        if self.supersedes_reporting_revision_id is not None:
            reporting_identifier(self.supersedes_reporting_revision_id, maximum=255)
var finality : Literal['snapshot', 'official']
Expand source code
@dataclass(frozen=True, slots=True)
class RevisionPublished(_ClosedValue):
    reporting_revision_id: str
    consumer_id: str
    finality: ReportingFinality
    supersedes_reporting_revision_id: str | None = None
    kind: Literal["revision_published"] = field(default="revision_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_revision_id, maximum=255)
        consumer_reference(self.consumer_id)
        if self.supersedes_reporting_revision_id is not None:
            reporting_identifier(self.supersedes_reporting_revision_id, maximum=255)
var kind : Literal['revision_published']
Expand source code
@dataclass(frozen=True, slots=True)
class RevisionPublished(_ClosedValue):
    reporting_revision_id: str
    consumer_id: str
    finality: ReportingFinality
    supersedes_reporting_revision_id: str | None = None
    kind: Literal["revision_published"] = field(default="revision_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_revision_id, maximum=255)
        consumer_reference(self.consumer_id)
        if self.supersedes_reporting_revision_id is not None:
            reporting_identifier(self.supersedes_reporting_revision_id, maximum=255)
var reporting_revision_id : str
Expand source code
@dataclass(frozen=True, slots=True)
class RevisionPublished(_ClosedValue):
    reporting_revision_id: str
    consumer_id: str
    finality: ReportingFinality
    supersedes_reporting_revision_id: str | None = None
    kind: Literal["revision_published"] = field(default="revision_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_revision_id, maximum=255)
        consumer_reference(self.consumer_id)
        if self.supersedes_reporting_revision_id is not None:
            reporting_identifier(self.supersedes_reporting_revision_id, maximum=255)
var supersedes_reporting_revision_id : str | None
Expand source code
@dataclass(frozen=True, slots=True)
class RevisionPublished(_ClosedValue):
    reporting_revision_id: str
    consumer_id: str
    finality: ReportingFinality
    supersedes_reporting_revision_id: str | None = None
    kind: Literal["revision_published"] = field(default="revision_published", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        reporting_identifier(self.reporting_revision_id, maximum=255)
        consumer_reference(self.consumer_id)
        if self.supersedes_reporting_revision_id is not None:
            reporting_identifier(self.supersedes_reporting_revision_id, maximum=255)
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 StatusChanged (scope: ReportingStatusScope,
health: ReportingHealth,
fingerprint: str,
checkpoint_generation: int,
previous_health: ReportingHealth | None = None,
issue_ids: tuple[str, ...] = (),
*,
kind: "Literal['status_changed']" = 'status_changed')
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)

One locked scope checkpoint generation; private identity stays off the wire.

Ancestors

  • adcp.reporting.ledger.delivery_models._ClosedValue

Instance variables

var checkpoint_generation : int
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)
var fingerprint : str
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)
var health : Literal['healthy', 'waiting', 'delayed', 'action_required', 'complete']
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)
var issue_ids : tuple[str, ...]
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)
var kind : Literal['status_changed']
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)
var previous_health : Literal['healthy', 'waiting', 'delayed', 'action_required', 'complete'] | None
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)
var scope : ReportingStatusScope
Expand source code
@dataclass(frozen=True, slots=True)
class StatusChanged(_ClosedValue):
    """One locked scope checkpoint generation; private identity stays off the wire."""

    scope: ReportingStatusScope
    health: ReportingHealth
    fingerprint: str
    checkpoint_generation: int
    previous_health: ReportingHealth | None = None
    issue_ids: tuple[str, ...] = ()
    kind: Literal["status_changed"] = field(default="status_changed", kw_only=True)

    def __post_init__(self) -> None:
        _freeze_fields(self)
        if (
            self.scope.generation_key is None
            or self.scope.feed_purpose is None
            or self.checkpoint_generation < 1
            or len(self.fingerprint) != 64
            or any(c not in "0123456789abcdef" for c in self.fingerprint)
            or tuple(sorted(set(self.issue_ids))) != self.issue_ids
            or len(self.issue_ids) > 16
            or (self.health in {"delayed", "action_required"} and not self.issue_ids)
        ):
            raise ReportingNotificationError("invalid_status_projection")
        for issue_id in self.issue_ids:
            reporting_identifier(issue_id)
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
class StoredDelivery (binding: DeliveryBinding,
envelope: bytes)
Expand source code
@dataclass(frozen=True)
class StoredDelivery:
    binding: DeliveryBinding
    envelope: bytes = field(repr=False)

StoredDelivery(binding: 'DeliveryBinding', envelope: 'bytes')

Instance variables

var binding : DeliveryBinding
var envelope : bytes
class WebhookAttempt (binding: DeliveryBinding,
attempt: int,
lease_token: str,
reservation_token: str,
fired_at: datetime,
request: ActivityRequest,
outcome: ActivityOutcome | None = None,
completed_at: datetime | None = None)
Expand source code
@dataclass(frozen=True)
class WebhookAttempt:
    binding: DeliveryBinding
    attempt: int
    lease_token: str = field(repr=False)
    reservation_token: str = field(repr=False)
    fired_at: datetime
    request: ActivityRequest
    outcome: ActivityOutcome | None = None
    completed_at: datetime | None = None

    def to_wire(self) -> dict[str, Any]:
        from adcp.validation.schema_loader import get_named_validator

        outcome = self.outcome
        status = outcome.status if outcome else "pending"
        row: dict[str, Any] = {
            "idempotency_key": self.binding.idempotency_key,
            "subscriber_id": self.binding.subscriber_id,
            "notification_type": self.binding.notification_type,
            "attempt": self.attempt,
            "fired_at": self.fired_at.isoformat(),
            "completed_at": self.completed_at.isoformat() if self.completed_at else None,
            "status": status,
            "url": sanitize_activity_url(self.request.url),
            "http_status_code": outcome.http_status_code if outcome else None,
            "response_time_ms": outcome.response_time_ms if outcome else None,
            "payload_size_bytes": self.request.payload_size_bytes,
            "error_message": {
                "pending": None,
                "success": None,
                "failed": "HTTP non-success response",
                "timeout": "HTTP attempt timed out",
                "connection_error": "Connection failed",
            }[status],
        }
        # Reporting notifications normally have an ID; optional wire fields
        # must be absent, rather than null, when the source has none.
        if self.binding.notification_id:
            row["notification_id"] = self.binding.notification_id
        validator = get_named_validator("core/webhook-activity-record.json")
        if validator is None or not validator.is_valid(row):
            raise ReportingNotificationError("invalid_activity_record")
        return row

WebhookAttempt(binding: 'DeliveryBinding', attempt: 'int', lease_token: 'str', reservation_token: 'str', fired_at: 'datetime', request: 'ActivityRequest', outcome: 'ActivityOutcome | None' = None, completed_at: 'datetime | None' = None)

Instance variables

var attempt : int
var binding : DeliveryBinding
var completed_at : datetime.datetime | None
var fired_at : datetime.datetime
var lease_token : str
var outcome : ActivityOutcome | None
var request : ActivityRequest
var reservation_token : str

Methods

def to_wire(self) ‑> dict[str, typing.Any]
Expand source code
def to_wire(self) -> dict[str, Any]:
    from adcp.validation.schema_loader import get_named_validator

    outcome = self.outcome
    status = outcome.status if outcome else "pending"
    row: dict[str, Any] = {
        "idempotency_key": self.binding.idempotency_key,
        "subscriber_id": self.binding.subscriber_id,
        "notification_type": self.binding.notification_type,
        "attempt": self.attempt,
        "fired_at": self.fired_at.isoformat(),
        "completed_at": self.completed_at.isoformat() if self.completed_at else None,
        "status": status,
        "url": sanitize_activity_url(self.request.url),
        "http_status_code": outcome.http_status_code if outcome else None,
        "response_time_ms": outcome.response_time_ms if outcome else None,
        "payload_size_bytes": self.request.payload_size_bytes,
        "error_message": {
            "pending": None,
            "success": None,
            "failed": "HTTP non-success response",
            "timeout": "HTTP attempt timed out",
            "connection_error": "Connection failed",
        }[status],
    }
    # Reporting notifications normally have an ID; optional wire fields
    # must be absent, rather than null, when the source has none.
    if self.binding.notification_id:
        row["notification_id"] = self.binding.notification_id
    validator = get_named_validator("core/webhook-activity-record.json")
    if validator is None or not validator.is_valid(row):
        raise ReportingNotificationError("invalid_activity_record")
    return row