Module adcp.reporting.outbox.models
Persistence contract for independently leased fanout and HTTP work.
These are additive protocols. Existing ReportingLedgerStore implementations do not acquire any new required method or producer callback.
Classes
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 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 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