Module adcp.reporting.materializer.capture
Immutable finish boundaries and a real, isolated notification enqueue path.
B2.1 does not expose a delivery worker or an activation certificate. A/B/C workers cannot see this queue. Pre-activation events stay quarantined forever. B2.4 must prove the mounted projection and tier components before admitting new events in the original verified-finish transaction; it cannot release old ones.
Functions
def decode_materializer_boundary(value: dict[str, Any]) ‑> ReportingMaterializerBoundary-
Expand source code
def decode_materializer_boundary(value: dict[str, Any]) -> ReportingMaterializerBoundary: result = None try: if type(value) is dict and type(value.get("version")) is int and value["version"] == 1: result = ReportingMaterializerBoundary( ReportingDeliveryPrincipal(value["account_id"], value["consumer_id"]), value["sequence"], value["account_sequence"], value["reporting_materialization_id"], datetime.fromisoformat(value["as_of"]), _core_adapter().validate_python(value["core"]), tuple(decode_record(r) for r in value["reconciliation"]), ) except (ValueError, TypeError, KeyError, ValidationError): # Convert malformed persisted input to the closed protocol error below. pass if result is None or result.to_storage() != value: raise ReportingNotificationError("materializer_boundary_invalid") return result async def enqueue_materializer_event_on(connection: Any, event: ReportingDomainEvent) ‑> None-
Expand source code
async def enqueue_materializer_event_on(connection: Any, event: ReportingDomainEvent) -> None: from adcp.reporting.outbox.pg import enqueue_event if event.notification_type != "reporting.delivery_ready": raise ReportingNotificationError("materializer_event_invalid") await enqueue_event(_MaterializerQueueConnection(connection), event) def private_snapshot(snapshot: ReportingStatusSnapshot, caller: ReportingDeliveryPrincipal) ‑> ReportingStatusSnapshot-
Expand source code
def private_snapshot( snapshot: ReportingStatusSnapshot, caller: ReportingDeliveryPrincipal ) -> ReportingStatusSnapshot: """Internal boundaries also exclude other consumers' private statements.""" if snapshot.account_id != caller.account_id: raise ReportingNotificationError("private_snapshot_unavailable") configurations = tuple( c for c in snapshot.configurations if c.consumer_id == caller.consumer_id ) keys = {c.generation_key for c in configurations} obligations = tuple(o for o in snapshot.obligations if o.generation_key in keys) obligation_ids = {o.reporting_obligation_id for o in obligations} revisions = tuple(r for r in snapshot.revisions if r.reporting_obligation_id in obligation_ids) revision_ids = {r.reporting_revision_id for r in revisions} adjustments = tuple( a for a in snapshot.adjustments if a.adjusts_reporting_revision_id in revision_ids ) statuses = tuple( s for s in snapshot.statuses if s.consumer_id == caller.consumer_id and s.generation_key in keys ) issue_scopes = dict(snapshot.issue_scopes) issues = tuple( i for i in snapshot.lifecycles if i.consumer_id == caller.consumer_id or ( i.consumer_id is None and i.issue_id in issue_scopes and issue_scopes[i.issue_id].generation_key in keys ) ) issue_ids = {i.issue_id for i in issues} visible_ids = ( obligation_ids | revision_ids | {a.reporting_adjustment_id for a in adjustments} | {s.reporting_status_id for s in statuses} ) return replace( snapshot, configurations=configurations, obligations=obligations, revisions=revisions, adjustments=adjustments, statuses=statuses, lifecycles=issues, consumer_ids=(caller.consumer_id,), issue_scopes=tuple((i, scope) for i, scope in snapshot.issue_scopes if i in issue_ids), changes=tuple( c for c in snapshot.changes if c[3] == caller.consumer_id and c[2] in visible_ids ), )Internal boundaries also exclude other consumers' private statements.
Classes
class MaterializerNotificationState (events: dict[tuple[str, str, str], ReportingDomainEvent] = <factory>,
causes: dict[tuple[str, str, str, str, str, int], str] = <factory>,
expansions: dict[tuple[str, str, str, int], _Work] = <factory>,
deliveries: dict[tuple[str, str, str], tuple[StoredDelivery, _Work]] = <factory>,
emissions: set[tuple[str, str, str, int, str]] = <factory>,
dirty: list[ReportingStatusDirty] = <factory>,
boundaries: list[StatusBoundary] = <factory>,
checkpoints: dict[tuple[str, str], int] = <factory>,
issue_scopes: dict[tuple[str, str], ReportingStatusScope] = <factory>,
activity_heads: dict[tuple[str, str, str, str, str], int] = <factory>,
activity: dict[tuple[str, str, str, str, str, int], WebhookAttempt] = <factory>)-
Expand source code
class MaterializerNotificationState(NotificationState): """Pre-activation events remain permanently quarantined, including on restart.""" def enqueue(self, event: ReportingDomainEvent) -> None: super().enqueue(event) for (account, consumer, notification_id, _), work in self.expansions.items(): if (account, consumer, notification_id) == ( event.account_id, event.consumer_namespace, event.notification_id, ): work.state = "quarantined"Pre-activation events remain permanently quarantined, including on restart.
Ancestors
Inherited members
class ReportingMaterializerBoundary (caller: ReportingDeliveryPrincipal,
sequence: int,
account_sequence: int,
reporting_materialization_id: str,
as_of: datetime,
core: ReportingStatusSnapshot,
reconciliation: tuple[ReportingDeliveryRecord, ...])-
Expand source code
@dataclass(frozen=True) class ReportingMaterializerBoundary: """Private frozen projection inputs; never returned as a public wire blob.""" caller: ReportingDeliveryPrincipal sequence: int account_sequence: int reporting_materialization_id: str as_of: datetime core: ReportingStatusSnapshot = field(repr=False) reconciliation: tuple[ReportingDeliveryRecord, ...] = field(repr=False) def __post_init__(self) -> None: outcomes = tuple( record for record in self.reconciliation if isinstance(record, ReportingMaterializationRecord) and record.reporting_materialization_id == self.reporting_materialization_id ) if ( type(self.sequence) is not int or self.sequence < 1 or type(self.account_sequence) is not int or self.account_sequence < self.sequence or self.core.account_id != self.caller.account_id or self.core.as_of != self.as_of or self.core.consumer_ids != (self.caller.consumer_id,) or any(s.consumer_id != self.caller.consumer_id for s in self.core.statuses) or any( i.consumer_id not in {None, self.caller.consumer_id} for i in self.core.lifecycles ) or any(principal(r) != self.caller for r in self.reconciliation) or len(outcomes) != 1 or outcomes[0].completed_at != self.as_of ): raise ReportingNotificationError("materializer_boundary_invalid") reporting_identifier(self.reporting_materialization_id, maximum=255) object.__setattr__(self, "as_of", aware_utc(self.as_of)) def to_storage(self) -> dict[str, Any]: return { "version": 1, "account_id": self.caller.account_id, "consumer_id": self.caller.consumer_id, "sequence": self.sequence, "account_sequence": self.account_sequence, "reporting_materialization_id": self.reporting_materialization_id, "as_of": self.as_of.isoformat(), "core": _core_adapter().dump_python(self.core, mode="json"), "reconciliation": [payload(r) for r in self.reconciliation], }Private frozen projection inputs; never returned as a public wire blob.
Instance variables
var account_sequence : intvar as_of : datetime.datetimevar caller : ReportingDeliveryPrincipalvar core : ReportingStatusSnapshotvar reconciliation : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...]var reporting_materialization_id : strvar sequence : int
Methods
def to_storage(self) ‑> dict[str, typing.Any]-
Expand source code
def to_storage(self) -> dict[str, Any]: return { "version": 1, "account_id": self.caller.account_id, "consumer_id": self.caller.consumer_id, "sequence": self.sequence, "account_sequence": self.account_sequence, "reporting_materialization_id": self.reporting_materialization_id, "as_of": self.as_of.isoformat(), "core": _core_adapter().dump_python(self.core, mode="json"), "reconciliation": [payload(r) for r in self.reconciliation], }