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 | Nonevar response_time_ms : int | Nonevar 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 : intvar 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: strOnly 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 : strvar auth_mode : strvar body_sha256 : strvar cause_generation : intvar cause_id : strvar cause_kind : strvar consumer_namespace : strvar delivery_id : strvar destination_sha256 : strvar emission_generation : intvar envelope_version : intvar idempotency_key : strvar key_version : strvar notification_id : strvar notification_type : strvar principal_id : strvar signing_scope_id : str | Nonevar subscriber_id : strvar 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: intDeliveryLease(delivery: 'StoredDelivery', token: 'str', expires_at: 'datetime', attempt_count: 'int')
Instance variables
var attempt_count : intvar delivery : StoredDeliveryvar expires_at : datetime.datetimevar 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 | NoneDeliveryStatus(delivery: 'StoredDelivery', state: 'WorkState', attempt_count: 'int', due_at: 'datetime', error_code: 'ErrorCode | None')
Instance variables
var attempt_count : intvar delivery : StoredDeliveryvar due_at : datetime.datetimevar error_code : Literal['network', 'retryable_http', 'permanent_http', 'signing_unavailable', 'permanent_scope', 'subscription_unavailable', 'subscription_changed', 'invalid_configuration', 'invalid_payload', 'integrity_failure', 'lease_expired'] | Nonevar 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 : strvar attempt_count : intvar consumer_namespace : strvar emission_generation : intvar event : dict[str, typing.Any]var expires_at : datetime.datetimevar notification_id : strvar 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 TrueView 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.outboxView 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 countGlobally 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.
nowarguments implement the shared memory protocol. In PostgreSQL they never override the database clock; only the explicit constructorclockseam 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
- PgReportingOutbox
- adcp.reporting.outbox._activity_pg._PgReportingActivity
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
- ReportingActivityReader
- typing.Protocol
- typing.Generic
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. Mountprojectorseparately asaccount_activity=to enable list-accounts enrichment. Neither relationship notification claims nor memory components supply evidence of durable reporting activity.Instance variables
var ledger : ReportingLedgerStorevar projector : ReportingActivityProjector | Nonevar status : ReportingStatusSupport | Nonevar 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 = NoneDiscard 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 NoneAES-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 NoneAuthenticate 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 : strvar 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 checkSee 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 : strvar active : boolprop 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 | Nonevar configuration_revision : strvar 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 : strvar proof_of_control_ref : strvar proof_valid : boolvar signing_scope_id : str | Nonevar subscriber_id : strvar 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 NoneAt-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 : strvar key_id : strvar 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 checkSee 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 = TrueDeterministic single turns for the adopter's supervised scheduler.
The caller owns scheduling and the database pool.
drainprocesses 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_idsis 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 : boolvar handler : ADCPHandler[typing.Any] | Nonevar projector : ReportingStatusProjector | Nonevar readiness_subscriptions : tuple[ReportingNotificationSubscription, ...]var scheduled : boolvar store : StatusNotificationStorevar sweeper : ReportingStatusSweeper | Nonevar 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 : floatvar 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: ReportingStatusSnapshotStatusBoundary(transaction_id: 'str', through: 'int', dirty: 'tuple[ReportingStatusDirty, …]', snapshot: 'ReportingStatusSnapshot')
Instance variables
var dirty : tuple[ReportingStatusDirty, ...]var snapshot : ReportingStatusSnapshotvar through : intvar 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_VERSIONStatusCheckpoint(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 : boolvar fingerprint : strvar generation : intvar lease_expires_at : datetime.datetime | Nonevar lease_token : str | Nonevar next_due_at : datetime.datetime | Nonevar publishable : boolvar scope : ReportingStatusScopevar selector_semantics_version : intvar selector_writer_floor : intvar 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: datetimeStatusDueLease(scope: 'ReportingStatusScope', token: 'str', expires_at: 'datetime', due_at: 'datetime')
Instance variables
var due_at : datetime.datetimevar expires_at : datetime.datetimevar scope : ReportingStatusScopevar 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 = 0StatusTurn(did_work: 'bool', events: 'int' = 0)
Instance variables
var did_work : boolvar 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 : DeliveryBindingvar 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 rowWebhookAttempt(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 : intvar binding : DeliveryBindingvar completed_at : datetime.datetime | Nonevar fired_at : datetime.datetimevar lease_token : strvar outcome : ActivityOutcome | Nonevar request : ActivityRequestvar 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