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 : int
var as_of : datetime.datetime
var caller : ReportingDeliveryPrincipal
var core : ReportingStatusSnapshot
var reconciliation : tuple[ReportingDestinationBinding | ReportingObligationDeliveryRecord | ReportingMaterializationAttempt | ReportingMaterializationRecord | ReportingMaterializationCheck | ReportingRevisionReceiptRecord | ReportingAdjustmentReceiptRecord, ...]
var reporting_materialization_id : str
var 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],
    }