Module adcp.reporting.production

Production reporting composition with explicit, frozen provider contracts.

Mount one support's handler, start its owned lifecycle, then activate accounts after migration and drain. The memory implementation is for conformance and does not claim durability. PostgreSQL uses the shared optional driver guard.

Sub-modules

adcp.reporting.production.configuration

A typed admission boundary for the application's existing account task.

adcp.reporting.production.contracts

Provider declarations and trusted, immutable source generation bindings.

adcp.reporting.production.delivery_window

Immutable first-attempt deadlines for the three production delivery queues.

adcp.reporting.production.handler

Authenticated production routes with private, revision-bound exact reads.

adcp.reporting.production.memory

Reference production transactions for shared state-machine conformance.

adcp.reporting.production.notifications

Optional owned notification delivery, independent of complete polling.

adcp.reporting.production.offerings

Atomic offerings bind public promises to the actual producer and verifier.

adcp.reporting.production.pg

Production admission selects new participants in the single B2.1 transaction.

adcp.reporting.production.schema

Separate production participants; no inherited mandatory manifest is widened.

adcp.reporting.production.service

One live SDK composition for admitted reporting work and truthful discovery.

adcp.reporting.production.source_registry

Fixed execution profiles and durable, secret-free service admission facts …

Functions

def production_notification_workers(store: InMemoryReportingProductionStore | PgReportingProductionStore,
projection: InMemoryReportingStatusProjection | PgReportingStatusProjection,
*,
subscriptions: ReportingSubscriptionResolver,
cipher: ReportingEnvelopeCipher,
signing: ReportingProductionSigning) ‑> tuple[ReportingNotificationWorker, ...]
Expand source code
def production_notification_workers(
    store: InMemoryReportingProductionStore | PgReportingProductionStore,
    projection: InMemoryReportingStatusProjection | PgReportingStatusProjection,
    *,
    subscriptions: ReportingSubscriptionResolver,
    cipher: ReportingEnvelopeCipher,
    signing: ReportingProductionSigning,
) -> tuple[ReportingNotificationWorker, ...]:
    """Build all three real SDK queue workers for ``ReportingProductionSupport``.

    The support schedules these workers itself. Omitting them keeps polling
    and enabled atomic logical enqueue available, without advertising push.
    Retained epoch-zero readiness queues are never among these participants.
    """
    from adcp.reporting.outbox.pg import PgReportingOutbox
    from adcp.reporting.production.pg import PgReportingProductionOutbox, PgReportingProductionStore

    if (
        projection.ledger is not store
        or not projection.policy["notifications_enabled"]
        or type(signing) is not ReportingProductionSigning
    ):
        raise ReportingNotificationError("notification_chain_unready")
    ReportingProductionSigning.__post_init__(signing)
    signed_subscriptions = _SignedSubscriptions(subscriptions)
    outboxes: tuple[Any, ...]
    if type(store) is InMemoryReportingProductionStore:
        outboxes = (
            InMemoryReportingOutbox(store),
            projection.outbox,
            InMemoryReportingProductionOutbox(store),
        )
    elif type(store) is PgReportingProductionStore:
        outboxes = (
            PgReportingOutbox(pool=store._pool, clock=store._clock),
            projection.outbox,
            PgReportingProductionOutbox(pool=store._pool, clock=store._clock),
        )
    else:
        raise ReportingNotificationError("notification_chain_unready")
    return tuple(
        ReportingNotificationWorker(
            outbox=outbox,
            subscriptions=signed_subscriptions,
            cipher=cipher,
            signing=signing,
            activity=outbox,
            clock=store._clock,
            delivery_window=ProductionDeliveryWindow(store, queue),
        )
        for outbox, queue in zip(outboxes, ("core", "status", "ready"))
    )

Build all three real SDK queue workers for ReportingProductionSupport.

The support schedules these workers itself. Omitting them keeps polling and enabled atomic logical enqueue available, without advertising push. Retained epoch-zero readiness queues are never among these participants.

Classes

class InMemoryReportingProductionOutbox (store: InMemoryReportingProductionStore)
Expand source code
class InMemoryReportingProductionOutbox(InMemoryReportingOutbox):
    """The ordinary SDK expansion/delivery state machine, in a separate collection."""

    _store: InMemoryReportingProductionStore

    def __init__(self, store: InMemoryReportingProductionStore) -> None:
        if store._production_outbox is None:
            raise ReportingNotificationError("notifications_disabled")
        self._store = store

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

The ordinary SDK expansion/delivery state machine, in a separate collection.

Ancestors

Inherited members

class InMemoryReportingProductionStore (**kwargs: Any)
Expand source code
class InMemoryReportingProductionStore(InMemoryReportingProjectionStore):
    """Matches production atomicity and epochs, without claiming durable storage."""

    _materializer_writer_epoch = 2
    _production_support: weakref.ReferenceType[ReportingProductionSupport] | None = None

    def __init__(self, **kwargs: Any) -> None:
        super().__init__(**kwargs)
        self._production_accounts: dict[str, _Admission] = {}
        self._production_delivery_windows: dict[
            tuple[str, str], tuple[str, str, datetime, datetime]
        ] = {}
        self._production_generations: dict[ReportingConfigurationGenerationKey, str] = {}
        self._production_source_bindings: dict[
            ReportingConfigurationGenerationKey, ReportingProductionSourceBinding
        ] = {}
        self._production_destination_bindings: dict[
            tuple[ReportingConfigurationGenerationKey, str], bytes
        ] = {}
        self._production_outbox = (
            NotificationState() if self._notification_state is not None else None
        )
        self._production_status_heads: dict[ReportingDeliveryPrincipal, int] = {}
        self._production_boundaries: list[ReportingMaterializerBoundary] = []
        self._production_closed: dict[ReportingConfigurationGenerationKey, datetime] = {}
        self._production_source_turns: dict[ReportingConfigurationGenerationKey, int] = {}
        self._production_source_work: dict[str, _SourceWork] = {}

    def _owner(self) -> ReportingProductionSupport:
        from adcp.reporting.production.service import production_owner

        return production_owner(self)

    async def admit_production_configuration(
        self,
        configuration: ReportingConfiguration,
        binding: ReportingDestinationBinding,
        *,
        offering_id: str,
        service_context: ReportingProductionSourceContext | None = None,
    ) -> None:
        offering = self._owner()._configuration_offering(
            configuration, binding, offering_id=offering_id
        )
        async with self._mutation():
            await self.put_configuration(configuration)
            await self.put_destination_binding(binding)
            self._enroll(
                configuration, offering._producer_key, binding, service_context=service_context
            )

    def _enroll(
        self,
        configuration: ReportingConfiguration,
        producer_key: str,
        destination: ReportingDestinationBinding,
        *,
        service_context: ReportingProductionSourceContext | None = None,
    ) -> None:
        previous_context = self._production_source_bindings.get(configuration.generation_key)
        binding = self._owner()._source_binding(
            configuration,
            producer_key,
            service_context=service_context,
            document=previous_context.document() if previous_context is not None else None,
        )
        previous = self._production_generations.setdefault(
            configuration.generation_key, producer_key
        )
        previous_binding = self._production_source_bindings.setdefault(
            configuration.generation_key, binding
        )
        if previous != producer_key or previous_binding != binding:
            raise ReportingNotificationError("reporting_production_source_conflict")
        owner = self._owner()
        offering = owner._configuration_offering(
            configuration, destination, producer_key=producer_key
        )
        document = canonical_json_utf8_v1(owner._destination_binding(destination, offering).wire())
        original = self._production_destination_bindings.setdefault(
            (configuration.generation_key, destination.consumer_id), document
        )
        if original != document:
            raise ReportingNotificationError("reporting_production_destination_conflict")

    def _destination_document(self, binding: ReportingDestinationBinding) -> dict[str, Any] | None:
        raw = self._production_destination_bindings.get(
            (binding.generation_key, binding.consumer_id)
        )
        return dict(json.loads(raw)) if raw is not None else None

    def _configuration_lease_eligible(self, configuration: ReportingConfiguration) -> bool:
        if configuration.quarantined:
            return False
        try:
            self._check_source_generation(configuration)
        except Exception:
            return False
        return True

    def _check_source_generation(self, configuration: ReportingConfiguration) -> None:
        producer_key = self._production_generations.get(configuration.generation_key)
        binding = self._production_source_bindings.get(configuration.generation_key)
        if (
            configuration.account_id not in self._production_accounts
            or producer_key not in self._owner()._producer_keys()
            or binding is None
            or self._configurations.get(configuration.generation_key) != configuration
        ):
            raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer generation is unavailable")
        self._owner()._check_source_binding(configuration, producer_key, binding.document())

    def _wake_obligation(self, account_id: str, obligation_id: str) -> None:
        super()._wake_obligation(account_id, obligation_id)
        obligation = self._obligations.get(obligation_id)
        if (
            obligation is not None
            and obligation.account_id == account_id
            and obligation.generation_key in self._production_generations
        ):
            previous = self._production_source_work.get(obligation_id)
            self._production_source_work[obligation_id] = _SourceWork(
                obligation, previous.turn if previous else 0
            )

    async def producer_closed_through(
        self, configuration: ReportingConfiguration
    ) -> datetime | None:
        async with self._lock:
            self._check_source_generation(configuration)
            return self._production_closed.get(configuration.generation_key)

    async def producer_constituents(
        self, configuration: ReportingConfiguration, obligation: ReportingObligationRecord
    ) -> tuple[ReportingConstituent, ...]:
        async with self._lock:
            self._check_source_generation(configuration)
            if obligation.generation_key != configuration.generation_key or set(
                obligation.media_buy_ids
            ) != set(configuration.media_buy_ids):
                raise LedgerConflictError("HISTORY_UNAVAILABLE", "source denominator differs")
            return self._production_source_bindings[configuration.generation_key].constituents()

    async def commit_producer_period(
        self,
        configuration: ReportingConfiguration,
        obligation: ReportingObligationRecord,
        *,
        previous_end: datetime | None,
    ) -> ReportingObligationRecord:
        check_next_period(configuration, obligation, previous_end)
        async with self._mutation():
            self._check_source_generation(configuration)
            current = self._production_closed.get(configuration.generation_key)
            if current != previous_end and (current is None or current < obligation.period.end):
                raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer progress changed")
            stored = await self.commit_obligation(obligation)
            self._production_source_work.setdefault(
                stored.reporting_obligation_id, _SourceWork(stored)
            )
            self._production_closed[configuration.generation_key] = max(
                obligation.period.end, current or obligation.period.end
            )
            return stored

    async def next_producer_obligations(
        self, configuration: ReportingConfiguration, *, now: datetime, limit: int
    ) -> tuple[str, ...]:
        if type(limit) is not int or not 1 <= limit <= 64:
            raise ValueError("production acquisition limit must be in 1..64")
        async with self._lock:
            self._check_source_generation(configuration)
            candidates = sorted(
                (work.turn, work.obligation.reporting_obligation_id)
                for work in self._production_source_work.values()
                if work.obligation.generation_key == configuration.generation_key
                and work.obligation.period.end <= now
                and work.state == "pending"
            )[:limit]
            turn = self._production_source_turns.get(configuration.generation_key, 0)
            for offset, (_, identifier) in enumerate(candidates, 1):
                work = self._production_source_work[identifier]
                self._production_source_work[identifier] = _SourceWork(
                    work.obligation, turn + offset
                )
            self._production_source_turns[configuration.generation_key] = turn + len(candidates)
            return tuple(identifier for _, identifier in candidates)

    async def finish_producer_acquisition(
        self, configuration: ReportingConfiguration, *, reporting_obligation_id: str
    ) -> None:
        async with self._lock:
            self._check_source_generation(configuration)
            work = self._production_source_work[reporting_obligation_id]
            if work.obligation.generation_key != configuration.generation_key:
                raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer generation differs")
            revisions = await self.list_revisions(
                account_id=configuration.account_id,
                reporting_obligation_id=reporting_obligation_id,
            )
            self._production_source_work[reporting_obligation_id] = _SourceWork(
                work.obligation, work.turn, acquisition_state(work.obligation, revisions)
            )

    async def _activate_production(self, *, account_id: str) -> bool:
        owner = self._owner()
        async with self._mutation():
            projection = self._projection_accounts.get(account_id)
            if (
                projection is None
                or not projection.ready
                or projection.policy != owner.projection.policy
            ):
                raise ReportingNotificationError("status_projection_activation_required")
            policy = owner._admission_policy()
            current = self._production_accounts.get(account_id)
            if current is not None:
                if current.policy != policy:
                    raise ReportingNotificationError("reporting_production_policy_conflict")
                return False
            self._production_accounts[account_id] = _Admission(self._clock(), policy)
            for _, _, record in self._retained_delivery_records():
                if (
                    isinstance(record, ReportingDestinationBinding)
                    and record.principal.account_id == account_id
                ):
                    configuration = self._configurations[record.generation_key]
                    try:
                        offering = owner._configuration_offering(configuration, record)
                    # Unsupported or unavailable bindings remain unadmitted.
                    except Exception:  # nosec B112
                        continue
                    self._enroll(configuration, offering._producer_key, record)
            return True

    def _check_admission(self, lease: ReportingMaterializerLease) -> None:
        if lease.admission_epoch != 2:
            return
        admission = self._production_accounts.get(lease.scope.principal.account_id)
        if admission is None or lease.attempt.created_at < admission.activated_at:
            raise failure("BINDING_MISMATCH")
        self._owner()._check_context(
            self._context(lease.scope),
            lease.request.verification_key,
            admission.policy,
            self._production_generations.get(lease.scope.generation_key),
            (
                self._production_source_bindings[lease.scope.generation_key].document()
                if lease.scope.generation_key in self._production_source_bindings
                else None
            ),
            self._destination_document(self._context(lease.scope).binding),
        )

    def _new_admission_epoch(
        self, context: MaterializerContext, key: ReportingVerificationKey
    ) -> int:
        admission = self._production_accounts.get(context.scope.principal.account_id)
        if admission is None:
            raise ReportingNotificationError("reporting_production_activation_required")
        self._owner()._check_context(
            context,
            key,
            admission.policy,
            self._production_generations.get(context.configuration.generation_key),
            (
                self._production_source_bindings[context.configuration.generation_key].document()
                if context.configuration.generation_key in self._production_source_bindings
                else None
            ),
            self._destination_document(context.binding),
        )
        return 2

    def _materializer_candidate_enabled(self, account_id: str) -> bool:
        return account_id in self._production_accounts

    def _held(self, lease: ReportingMaterializerLease) -> _Work | None:
        work = super()._held(lease)
        if work is not None:
            self._check_admission(lease)
        return work

    def _materializer_capture_collections(
        self, outcome: ReportingMaterializationRecord
    ) -> tuple[dict[ReportingDeliveryPrincipal, int], list[ReportingMaterializerBoundary]]:
        work = self._materializer_work.get(
            (
                outcome.scope.principal.account_id,
                outcome.scope.consumer_id,
                outcome.reporting_materialization_id,
            )
        )
        if work is not None and work.admission_epoch == 2:
            return self._production_status_heads, self._production_boundaries
        return super()._materializer_capture_collections(outcome)

    def _enqueue_materializer(
        self, event: ReportingDomainEvent, lease: ReportingMaterializerLease
    ) -> None:
        if lease.admission_epoch == 0:
            return super()._enqueue_materializer(event, lease)
        self._check_admission(lease)
        if self._production_outbox is None:
            raise failure("BINDING_MISMATCH")
        self._production_outbox.enqueue(event)
        if (
            self._production_outbox.events.get(
                (event.account_id, event.consumer_namespace, event.notification_id)
            )
            != event
        ):
            raise failure("BINDING_MISMATCH")

    async def read_production_boundaries(
        self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
    ) -> tuple[ReportingMaterializerBoundary, ...]:
        if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100:
            raise ValueError("production boundary reads require bounded positions")
        async with self._lock:
            return tuple(
                b for b in self._production_boundaries if b.caller == caller and b.sequence > after
            )[:limit]

Matches production atomicity and epochs, without claiming durable storage.

Ancestors

Methods

async def admit_production_configuration(self,
configuration: ReportingConfiguration,
binding: ReportingDestinationBinding,
*,
offering_id: str,
service_context: ReportingProductionSourceContext | None = None) ‑> None
Expand source code
async def admit_production_configuration(
    self,
    configuration: ReportingConfiguration,
    binding: ReportingDestinationBinding,
    *,
    offering_id: str,
    service_context: ReportingProductionSourceContext | None = None,
) -> None:
    offering = self._owner()._configuration_offering(
        configuration, binding, offering_id=offering_id
    )
    async with self._mutation():
        await self.put_configuration(configuration)
        await self.put_destination_binding(binding)
        self._enroll(
            configuration, offering._producer_key, binding, service_context=service_context
        )
async def commit_producer_period(self,
configuration: ReportingConfiguration,
obligation: ReportingObligationRecord,
*,
previous_end: datetime | None) ‑> ReportingObligationRecord
Expand source code
async def commit_producer_period(
    self,
    configuration: ReportingConfiguration,
    obligation: ReportingObligationRecord,
    *,
    previous_end: datetime | None,
) -> ReportingObligationRecord:
    check_next_period(configuration, obligation, previous_end)
    async with self._mutation():
        self._check_source_generation(configuration)
        current = self._production_closed.get(configuration.generation_key)
        if current != previous_end and (current is None or current < obligation.period.end):
            raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer progress changed")
        stored = await self.commit_obligation(obligation)
        self._production_source_work.setdefault(
            stored.reporting_obligation_id, _SourceWork(stored)
        )
        self._production_closed[configuration.generation_key] = max(
            obligation.period.end, current or obligation.period.end
        )
        return stored
async def finish_producer_acquisition(self, configuration: ReportingConfiguration, *, reporting_obligation_id: str) ‑> None
Expand source code
async def finish_producer_acquisition(
    self, configuration: ReportingConfiguration, *, reporting_obligation_id: str
) -> None:
    async with self._lock:
        self._check_source_generation(configuration)
        work = self._production_source_work[reporting_obligation_id]
        if work.obligation.generation_key != configuration.generation_key:
            raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer generation differs")
        revisions = await self.list_revisions(
            account_id=configuration.account_id,
            reporting_obligation_id=reporting_obligation_id,
        )
        self._production_source_work[reporting_obligation_id] = _SourceWork(
            work.obligation, work.turn, acquisition_state(work.obligation, revisions)
        )
async def next_producer_obligations(self, configuration: ReportingConfiguration, *, now: datetime, limit: int) ‑> tuple[str, ...]
Expand source code
async def next_producer_obligations(
    self, configuration: ReportingConfiguration, *, now: datetime, limit: int
) -> tuple[str, ...]:
    if type(limit) is not int or not 1 <= limit <= 64:
        raise ValueError("production acquisition limit must be in 1..64")
    async with self._lock:
        self._check_source_generation(configuration)
        candidates = sorted(
            (work.turn, work.obligation.reporting_obligation_id)
            for work in self._production_source_work.values()
            if work.obligation.generation_key == configuration.generation_key
            and work.obligation.period.end <= now
            and work.state == "pending"
        )[:limit]
        turn = self._production_source_turns.get(configuration.generation_key, 0)
        for offset, (_, identifier) in enumerate(candidates, 1):
            work = self._production_source_work[identifier]
            self._production_source_work[identifier] = _SourceWork(
                work.obligation, turn + offset
            )
        self._production_source_turns[configuration.generation_key] = turn + len(candidates)
        return tuple(identifier for _, identifier in candidates)
async def producer_closed_through(self, configuration: ReportingConfiguration) ‑> datetime.datetime | None
Expand source code
async def producer_closed_through(
    self, configuration: ReportingConfiguration
) -> datetime | None:
    async with self._lock:
        self._check_source_generation(configuration)
        return self._production_closed.get(configuration.generation_key)
async def producer_constituents(self,
configuration: ReportingConfiguration,
obligation: ReportingObligationRecord) ‑> tuple[ProductConstituentV1 | PackageItemConstituentV1 | MediaBuyConstituentV1, ...]
Expand source code
async def producer_constituents(
    self, configuration: ReportingConfiguration, obligation: ReportingObligationRecord
) -> tuple[ReportingConstituent, ...]:
    async with self._lock:
        self._check_source_generation(configuration)
        if obligation.generation_key != configuration.generation_key or set(
            obligation.media_buy_ids
        ) != set(configuration.media_buy_ids):
            raise LedgerConflictError("HISTORY_UNAVAILABLE", "source denominator differs")
        return self._production_source_bindings[configuration.generation_key].constituents()
async def read_production_boundaries(self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100) ‑> tuple[ReportingMaterializerBoundary, ...]
Expand source code
async def read_production_boundaries(
    self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
) -> tuple[ReportingMaterializerBoundary, ...]:
    if type(after) is not int or after < 0 or type(limit) is not int or not 1 <= limit <= 100:
        raise ValueError("production boundary reads require bounded positions")
    async with self._lock:
        return tuple(
            b for b in self._production_boundaries if b.caller == caller and b.sequence > after
        )[:limit]

Inherited members

class PgReportingProductionOutbox (*, pool: AsyncConnectionPool, clock: Callable[[], datetime] | None = None)
Expand source code
class PgReportingProductionOutbox(PgReportingStatusOutbox):
    """Generic crash-safe delivery/fanout and private activity on the admitted queue."""

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

    async def create_schema(self) -> None:
        await PgReportingProductionStore(pool=self._pool, notifications=True).create_schema()

Generic crash-safe delivery/fanout and private activity on the admitted queue.

Ancestors

Methods

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

Inherited members

class PgReportingProductionStore (*,
pool: AsyncConnectionPool,
clock: Callable[[], datetime] | None = None,
notifications: bool = False)
Expand source code
class PgReportingProductionStore(PgReportingProjectionStore):
    """The production store retains all ordinary legacy reader/writer APIs.

    New work requires the live SDK composition and its mounted applicable
    routes. Pending epoch-zero work uses its original tables, identity and
    permanently quarantined enqueue. Installation alone admits nothing.
    """

    _production_support: weakref.ReferenceType[ReportingProductionSupport] | None = None
    _production_lease_samples: dict[tuple[str, ...], _ProducerSample] | None = None

    def _owner(self) -> ReportingProductionSupport:
        from adcp.reporting.production.service import production_owner

        return production_owner(self)

    @asynccontextmanager
    async def _connection(self) -> AsyncIterator[Any]:
        admission = _ADMISSION_CONNECTION.get()
        if admission is not None and admission[:2] == (id(self), asyncio.current_task()):
            yield admission[2]
            return
        async with super()._connection() as connection:
            if _EPOCH.get() == (id(self), 2):
                yield _ProductionConnection(connection)
            else:
                yield connection

    async def admit_production_configuration(
        self,
        configuration: ReportingConfiguration,
        binding: ReportingDestinationBinding,
        *,
        offering_id: str,
        service_context: ReportingProductionSourceContext | None = None,
    ) -> None:
        offering = self._owner()._configuration_offering(
            configuration, binding, offering_id=offering_id
        )
        async with self._connection() as connection, connection.transaction():
            await self._lock_account(connection, configuration.account_id)
            token = _ADMISSION_CONNECTION.set((id(self), asyncio.current_task(), connection))
            try:
                await self.put_configuration(configuration)
                await self.put_destination_binding(binding)
                await self._enroll_on(
                    connection,
                    configuration,
                    offering._producer_key,
                    binding,
                    service_context=service_context,
                )
            finally:
                _ADMISSION_CONNECTION.reset(token)

    async def _enroll_on(
        self,
        connection: Any,
        configuration: ReportingConfiguration,
        producer_key: str,
        destination: ReportingDestinationBinding,
        *,
        service_context: ReportingProductionSourceContext | None = None,
    ) -> None:
        key = configuration.generation_key
        identity = (
            key.account_id,
            key.consumer_id,
            key.delivery_config_id,
            key.delivery_config_version,
        )
        previous_context = await (
            await connection.execute(
                "SELECT source_binding FROM reporting_production_generations"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s",
                identity,
            )
        ).fetchone()
        binding = (
            self._owner()
            ._source_binding(
                configuration,
                producer_key,
                service_context=service_context,
                document=previous_context[0] if previous_context is not None else None,
            )
            .document()
        )
        await connection.execute(
            "INSERT INTO reporting_production_generations VALUES(%s,%s,%s,%s,%s,%s::jsonb)"
            " ON CONFLICT DO NOTHING",
            (*identity, producer_key, json.dumps(binding)),
        )
        row = await (
            await connection.execute(
                "SELECT producer_key,source_binding FROM reporting_production_generations"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s",
                identity,
            )
        ).fetchone()
        if row is None or row != (producer_key, binding):
            raise ReportingNotificationError("reporting_production_source_conflict")
        owner = self._owner()
        offering = owner._configuration_offering(
            configuration, destination, producer_key=producer_key
        )
        document = owner._destination_binding(destination, offering).wire()
        destination_identity = (
            key.account_id,
            destination.consumer_id,
            key.delivery_config_id,
            key.delivery_config_version,
        )
        await connection.execute(
            "INSERT INTO reporting_production_destination_bindings VALUES(%s,%s,%s,%s,%s::jsonb)"
            " ON CONFLICT DO NOTHING",
            (*destination_identity, json.dumps(document)),
        )
        original = await (
            await connection.execute(
                "SELECT method FROM reporting_production_destination_bindings"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s",
                destination_identity,
            )
        ).fetchone()
        if original is None or original[0] != document:
            raise ReportingNotificationError("reporting_production_destination_conflict")

    async def lease_period_close(
        self, *, worker_id: str, now: datetime, lease_seconds: float
    ) -> LeasedConfiguration | None:
        keys = self._owner()._producer_keys()
        if not keys:
            return None
        moment = _utc(now)
        expires = moment + timedelta(seconds=lease_seconds)
        if expires <= moment:
            raise ValueError("producer lease duration must be positive")
        if self._production_lease_samples is None:
            self._production_lease_samples = {}
        after = self._production_lease_samples.get(keys)
        following: _ProducerSample | None = None
        result = None
        async with self._connection() as connection, connection.transaction():
            await self._bound_lock_waits(connection)
            # Discover a bounded set without locking configuration rows. The
            # inherited configuration trigger takes the account lock, so that
            # lock must precede the row lock here, just as it does in activation.
            # A busy account cannot advance a durable rank while another
            # transaction holds its lock. Continue a read-only sample instead
            # of repeatedly trying the same prefix. At most one nonempty
            # 32-row window is examined; an empty tail may wrap once.
            for _ in range(2):
                continuation = (
                    " AND (greatest(coalesce(t.lease_turn,0),coalesce(p.probe_turn,0)),"
                    " coalesce(c.lease_expires_at,'-infinity'::timestamptz),"
                    " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_config_version)"
                    " > (%s,coalesce(%s::timestamptz,'-infinity'::timestamptz),%s,%s,%s,%s)"
                    if after is not None
                    else ""
                )
                query = (
                    "SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.delivery"
                    "_config_version,"  # nosec B608
                    " greatest(coalesce(t.lease_turn,0),coalesce(p.probe_turn,0)),"
                    " c.lease_expires_at"
                    " FROM reporting_production_generations g"
                    " JOIN reporting_production_accounts a ON a.account_id=g.account_id"
                    " JOIN reporting_configurations c"
                    " ON (c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_c"
                    "onfig_version)="
                    " (g.account_id,g.consumer_id,g.delivery_config_id,g.delivery_config_version)"
                    " LEFT JOIN adcp_reporting_configuration_lease_turns t"
                    " ON (t.account_id,t.consumer_id,t.delivery_config_id,t.delivery_c"
                    "onfig_version)="
                    " (g.account_id,g.consumer_id,g.delivery_config_id,g.delivery_config_version)"
                    " LEFT JOIN reporting_production_source_probe_turns p"
                    " ON (p.account_id,p.consumer_id,p.delivery_config_id,p.delivery_c"
                    "onfig_version)="
                    " (g.account_id,g.consumer_id,g.delivery_config_id,g.delivery_config_version)"
                    " WHERE g.producer_key=ANY(%s)"
                    " AND NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lease"
                    "_expires_at<=%s)"
                    + continuation
                    + " ORDER BY greatest(coalesce(t.lease_turn,0),coalesce(p.probe_turn,0)),"
                    " c.lease_expires_at NULLS FIRST,"
                    " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_confi"
                    "g_version LIMIT 32"
                )
                # Only the fixed SDK continuation above changes this SQL;
                # every cursor value and source identity remains a parameter.
                rows = await (
                    await connection.execute(query, (list(keys), moment, *(after or ())))
                ).fetchall()
                if rows or after is None:
                    break
                after = None
            for candidate in rows:
                row = candidate[:4]
                following = (candidate[4], candidate[5], *row)
                locked = await (
                    await connection.execute(
                        "SELECT pg_try_advisory_xact_lock(hashtext('adcp.reporting:' || %s))",
                        (row[0],),
                    )
                ).fetchone()
                if not locked[0]:
                    continue
                current = await (
                    await connection.execute(_SOURCE_CONFIGURATION, (*row, list(keys)))
                ).fetchone()
                if current is None:
                    continue
                configuration = _configuration_from_row(current[:18])
                try:
                    self._owner()._check_source_binding(configuration, current[18], current[19])
                except Exception:
                    # A permanently revoked generation must not occupy the
                    # first bounded window forever. This is a probe, not a
                    # lease: preserve ordinary lease ranks and all source and
                    # external identities. Its durable rank shares the same
                    # ordering clock as successful configuration acquisitions.
                    await connection.execute(
                        "INSERT INTO reporting_production_source_probe_turns"
                        " (account_id,consumer_id,delivery_config_id,delivery_config_versi"
                        "on,probe_turn)"
                        " VALUES(%s,%s,%s,%s,nextval('adcp_reporting_configuration_lease_t"
                        "urn_seq'))"
                        " ON CONFLICT(account_id,consumer_id,delivery_config_id,delivery_c"
                        "onfig_version)"
                        " DO UPDATE SET probe_turn="
                        "nextval('adcp_reporting_configuration_lease_turn_seq')",
                        tuple(row),
                    )
                    continue
                acquired = await (
                    await connection.execute(
                        "UPDATE reporting_configurations SET lease_worker_id=%s,lease_expires_at=%s"
                        " WHERE (account_id,consumer_id,delivery_config_id,delivery_config"
                        "_version) = ("
                        " SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.deliver"
                        "y_config_version"
                        " FROM reporting_configurations c WHERE c.account_id=%s"
                        " AND c.consumer_id=%s AND c.delivery_config_id=%s AND c.delivery_"
                        "config_version=%s"
                        " AND NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lease"
                        "_expires_at<=%s)"
                        " FOR UPDATE OF c SKIP LOCKED)"
                        " RETURNING account_id",
                        (worker_id, expires, *row, moment),
                    )
                ).fetchone()
                if acquired is None:
                    continue
                await self._retain_materializer_generation_on_lease_change(connection, tuple(row))
                await connection.execute(
                    "INSERT INTO adcp_reporting_configuration_lease_turns"
                    " (account_id,consumer_id,delivery_config_id,delivery_config_versi"
                    "on,lease_turn)"
                    " VALUES(%s,%s,%s,%s,nextval('adcp_reporting_configuration_lease_turn_seq'))"
                    " ON CONFLICT(account_id,consumer_id,delivery_config_id,delivery_c"
                    "onfig_version)"
                    " DO UPDATE SET lease_turn="
                    "nextval('adcp_reporting_configuration_lease_turn_seq')",
                    tuple(row),
                )
                result = LeasedConfiguration(row[0], row[1], row[2], row[3], expires)
                break
        # Hints are per store and selected producer keys, not durable work or
        # leases. Publish a hint only after commit; a failed mutation retries
        # the same window. Successful acquisition returns to the durable
        # turn-primary order. A fresh store starts there too. Concurrent hints
        # may cause a bounded revisit, but cannot authorize or fence any work.
        if result is None and following is not None:
            self._production_lease_samples[keys] = following
        else:
            self._production_lease_samples.pop(keys, None)
        return result

    async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) -> None:
        async with self._connection() as connection, connection.transaction():
            await self._lock_account(connection, lease.account_id)
            identity = (
                lease.account_id,
                lease.consumer_id,
                lease.delivery_config_id,
                lease.delivery_config_version,
            )
            released = await (
                await connection.execute(
                    "UPDATE reporting_configurations SET lease_worker_id=NULL,lease_expires_at=NULL"
                    " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                    " AND delivery_config_version=%s"
                    " AND lease_worker_id=%s AND lease_expires_at=%s RETURNING account_id",
                    (*identity, worker_id, lease.lease_expires_at),
                )
            ).fetchone()
            if released is not None:
                await self._retain_materializer_generation_on_lease_change(connection, identity)

    async def _retain_materializer_generation_on_lease_change(
        self, connection: Any, identity: tuple[str, str, str, int]
    ) -> None:
        # The immutable inherited trigger treats *every* configuration UPDATE
        # as source invalidation, including lease-only bookkeeping. These two
        # SDK statements change only lease fields, under the account lock.
        # Cancel just their one trigger increment in the same transaction;
        # the wakeup remains harmless. No intervening source write can be
        # hidden, no committed generation goes backwards, and pending work's
        # original generation/epoch/external identity is never rewritten.
        # A real configuration, revision or readability change still executes
        # the original trigger without this correction and fences old work.
        await connection.execute(
            "UPDATE reporting_materializer_candidates SET generation=generation-1"
            " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
            " AND delivery_config_version=%s",
            identity,
        )

    @asynccontextmanager
    async def _source_connection(
        self, configuration: ReportingConfiguration
    ) -> AsyncIterator[tuple[Any, tuple[str, str, str, int]]]:
        keys = self._owner()._producer_keys()
        key = configuration.generation_key
        identity = (
            key.account_id,
            key.consumer_id,
            key.delivery_config_id,
            key.delivery_config_version,
        )
        async with self._connection() as connection, connection.transaction():
            await self._lock_account(connection, key.account_id)
            row = await (
                await connection.execute(
                    _SOURCE_CONFIGURATION,
                    (*identity, list(keys)),
                )
            ).fetchone()
            if row is None or _configuration_from_row(row[:18]) != configuration:
                raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer generation unavailable")
            self._owner()._check_source_binding(configuration, row[18], row[19])
            await connection.execute(
                "INSERT INTO reporting_production_source_progress"
                " (account_id,consumer_id,delivery_config_id,delivery_config_versi"
                "on) VALUES(%s,%s,%s,%s)"
                " ON CONFLICT DO NOTHING",
                identity,
            )
            token = _ADMISSION_CONNECTION.set((id(self), asyncio.current_task(), connection))
            try:
                yield connection, identity
            finally:
                _ADMISSION_CONNECTION.reset(token)

    async def producer_constituents(
        self, configuration: ReportingConfiguration, obligation: ReportingObligationRecord
    ) -> tuple[ReportingConstituent, ...]:
        async with self._source_connection(configuration) as (connection, identity):
            if obligation.generation_key != configuration.generation_key or set(
                obligation.media_buy_ids
            ) != set(configuration.media_buy_ids):
                raise LedgerConflictError("HISTORY_UNAVAILABLE", "source denominator differs")
            row = await (
                await connection.execute(
                    "SELECT producer_key,source_binding FROM reporting_production_generations"
                    " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                    " AND delivery_config_version=%s",
                    identity,
                )
            ).fetchone()
            binding = self._owner()._check_source_binding(configuration, row[0], row[1])
            return binding.constituents()

    async def producer_closed_through(
        self, configuration: ReportingConfiguration
    ) -> datetime | None:
        async with self._source_connection(configuration) as (connection, identity):
            row = await (
                await connection.execute(
                    "SELECT closed_through FROM reporting_production_source_progress"
                    " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                    " AND delivery_config_version=%s",
                    identity,
                )
            ).fetchone()
            return row[0] if row is not None else None

    async def commit_producer_period(
        self,
        configuration: ReportingConfiguration,
        obligation: ReportingObligationRecord,
        *,
        previous_end: datetime | None,
    ) -> ReportingObligationRecord:
        check_next_period(configuration, obligation, previous_end)
        async with self._source_connection(configuration) as (connection, identity):
            row = await (
                await connection.execute(
                    "SELECT closed_through FROM reporting_production_source_progress"
                    " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                    " AND delivery_config_version=%s",
                    identity,
                )
            ).fetchone()
            current = row[0]
            if current != previous_end and (current is None or current < obligation.period.end):
                raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer progress changed")
            stored = await self.commit_obligation(obligation)
            await connection.execute(
                "INSERT INTO reporting_production_source_work"
                " (account_id,consumer_id,delivery_config_id,delivery_config_versi"
                "on,reporting_obligation_id,"
                " period_end) VALUES(%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING",
                (*identity, stored.reporting_obligation_id, stored.period.end),
            )
            await connection.execute(
                "UPDATE reporting_production_source_progress"
                " SET closed_through=greatest(closed_through,%s)"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s",
                (stored.period.end, *identity),
            )
            return stored

    async def next_producer_obligations(
        self, configuration: ReportingConfiguration, *, now: datetime, limit: int
    ) -> tuple[str, ...]:
        if type(limit) is not int or not 1 <= limit <= 64:
            raise ValueError("production acquisition limit must be in 1..64")
        async with self._source_connection(configuration) as (connection, identity):
            rows = await (
                await connection.execute(
                    "SELECT reporting_obligation_id FROM reporting_production_source_work"
                    " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                    " AND delivery_config_version=%s"
                    " AND state='pending' AND period_end<=%s"
                    " ORDER BY acquisition_turn,reporting_obligation_id LIMIT %s FOR UPDATE",
                    (*identity, now, limit),
                )
            ).fetchall()
            if not rows:
                return ()
            head = await (
                await connection.execute(
                    "UPDATE reporting_production_source_progress"
                    " SET acquisition_turn=acquisition_turn+%s"
                    " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                    " AND delivery_config_version=%s"
                    " RETURNING acquisition_turn",
                    (len(rows), *identity),
                )
            ).fetchone()
            for offset, row in enumerate(rows, 1):
                await connection.execute(
                    "UPDATE reporting_production_source_work SET acquisition_turn=%s"
                    " WHERE account_id=%s AND reporting_obligation_id=%s",
                    (head[0] - len(rows) + offset, identity[0], row[0]),
                )
            return tuple(row[0] for row in rows)

    async def finish_producer_acquisition(
        self, configuration: ReportingConfiguration, *, reporting_obligation_id: str
    ) -> None:
        async with self._source_connection(configuration) as (connection, identity):
            obligation = await self.get_obligation(
                account_id=identity[0], reporting_obligation_id=reporting_obligation_id
            )
            if obligation is None or obligation.generation_key != configuration.generation_key:
                raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer generation differs")
            revisions = await self.list_revisions(
                account_id=identity[0], reporting_obligation_id=reporting_obligation_id
            )
            await connection.execute(
                "UPDATE reporting_production_source_work SET state=%s"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s"
                " AND reporting_obligation_id=%s",
                (acquisition_state(obligation, revisions), *identity, reporting_obligation_id),
            )

    @asynccontextmanager
    async def _lease_epoch(self, lease: ReportingMaterializerLease) -> AsyncIterator[None]:
        token = _EPOCH.set((id(self), lease.admission_epoch))
        try:
            yield
        finally:
            _EPOCH.reset(token)

    async def create_schema(self) -> None:
        async with self._connection() as connection, connection.transaction():
            await self._create_schema_on(connection)
            root = files("adcp.reporting.ledger")
            for name in (
                "reporting_materializer.sql",
                "reporting_receipt_ingestion.sql",
                "reporting_feed.sql",
                "reporting_status_notifications.sql",
                "reporting_status_selector_version.sql",
                "reporting_projection.sql",
                "reporting_projection_notifications.sql",
                "reporting_projection_feed.sql",
                "reporting_production.sql",
            ):
                await connection.execute(root.joinpath(name).read_text())
            await connection.execute(
                files("adcp.reporting.production").joinpath("service_context.sql").read_text()
            )

    async def materializer_ready(self) -> bool:
        owner = self._owner()
        if not await owner._schema_ready():
            raise ReportingNotificationError("reporting_production_schema_unready")
        owner._assert_components()
        return True

    async def _activate_production(self, *, account_id: str) -> bool:
        owner = self._owner()
        await self.materializer_ready()
        async with self._connection() as connection, connection.transaction():
            await self._lock_account(connection, account_id)
            owner._assert_components()
            projection = await (
                await connection.execute(
                    "SELECT policy FROM reporting_projection_accounts"
                    " WHERE account_id=%s AND current_input IS NOT NULL",
                    (account_id,),
                )
            ).fetchone()
            if projection is None or projection[0] != owner.projection.policy:
                raise ReportingNotificationError("status_projection_activation_required")
            policy = owner._admission_policy()
            current = await (
                await connection.execute(
                    "SELECT policy FROM reporting_production_accounts WHERE account_id=%s",
                    (account_id,),
                )
            ).fetchone()
            if current is not None:
                if current[0] != policy:
                    raise ReportingNotificationError("reporting_production_policy_conflict")
                return False
            await connection.execute(
                "INSERT INTO reporting_production_accounts(account_id,policy) VALUES(%s,%s::jsonb)",
                (account_id, json.dumps(policy)),
            )
            configurations = {}
            bindings = await (
                await connection.execute(
                    "SELECT payload FROM reporting_reconciliation_records"
                    " WHERE account_id=%s AND namespace='destination_binding'",
                    (account_id,),
                )
            ).fetchall()
            for (document,) in bindings:
                binding = decode_record(document)
                if not isinstance(binding, ReportingDestinationBinding):
                    raise ReportingNotificationError("reporting_production_history_corrupt")
                owned = await self._list_configurations_on(
                    connection, account_id=account_id, consumer_id=binding.consumer_id
                )
                configurations.update({c.generation_key: c for c in owned if not c.quarantined})
                if binding.generation_key not in configurations:
                    continue
                configuration = configurations[binding.generation_key]
                try:
                    offering = owner._configuration_offering(configuration, binding)
                # Unsupported or unavailable bindings remain unadmitted.
                except Exception:  # nosec B112
                    continue
                await self._enroll_on(connection, configuration, offering._producer_key, binding)
            return True

    async def _materializer_context_on(
        self, connection: Any, scope: ReportingDeliveryScope
    ) -> MaterializerContext:
        context = await super()._materializer_context_on(connection, scope)
        if isinstance(connection, _ProductionConnection):
            owner = self._owner()
            row = await (
                await connection.execute(
                    "SELECT a.policy,g.producer_key,g.source_binding,d.method"
                    " FROM reporting_production_accounts a"
                    " LEFT JOIN reporting_production_generations g"
                    " ON g.account_id=a.account_id AND g.consumer_id=%s AND g.delivery_config_id=%s"
                    " AND g.delivery_config_version=%s"
                    " LEFT JOIN reporting_production_destination_bindings d"
                    " ON (d.account_id,d.consumer_id,d.delivery_config_id,d.delivery_c"
                    "onfig_version)="
                    " (g.account_id,g.consumer_id,g.delivery_config_id,g.delivery_config_version)"
                    " AND d.consumer_id=%s WHERE a.account_id=%s",
                    (
                        scope.consumer_id,
                        scope.generation_key.delivery_config_id,
                        scope.generation_key.delivery_config_version,
                        scope.consumer_id,
                        scope.principal.account_id,
                    ),
                )
            ).fetchone()
            if row is None:
                raise ReportingNotificationError("reporting_production_activation_required")
            owner._check_context(
                context,
                key_for(context.binding, context.obligation, owner.keys),
                row[0],
                row[1],
                row[2],
                row[3],
            )
        return context

    async def _claim_account_on(
        self,
        connection: Any,
        account_id: str,
        keys: tuple[ReportingVerificationKey, ...],
        lease_seconds: int,
    ) -> ReportingMaterializerLease | ReportingMaterializerTurn:
        activated = await (
            await connection.execute(
                "SELECT 1 FROM reporting_production_accounts WHERE account_id=%s", (account_id,)
            )
        ).fetchone()
        if activated is None:
            return ReportingMaterializerTurn("idle")
        old = await (
            await connection.execute(
                "SELECT to_jsonb(w) FROM reporting_materializer_work w WHERE account_id=%s"
                " AND state='pending' AND due_at<=clock_timestamp()"
                " AND (lease_until IS NULL OR lease_until<=clock_timestamp())"
                " ORDER BY due_at,reporting_materialization_id LIMIT 1 FOR UPDATE",
                (account_id,),
            )
        ).fetchone()
        if old is not None:
            return await super()._lease_on(connection, old[0], keys, lease_seconds)
        return await super()._claim_account_on(
            _ProductionConnection(connection), account_id, keys, lease_seconds
        )

    async def _schedule_account_on(self, connection: Any, account_id: str) -> None:
        if isinstance(connection, _ProductionConnection):
            connection = connection.connection
        # _OWNED contains only fixed SDK tables/columns; account values are bound.
        await connection.execute(
            "UPDATE reporting_materializer_accounts SET due_at=(SELECT min(due) FROM ("  # nosec B608
            " SELECT due_at AS due FROM reporting_materializer_work"
            " WHERE account_id=%s AND state='pending'"
            " UNION ALL SELECT due_at FROM reporting_production_work"
            " WHERE account_id=%s AND state='pending'"
            " UNION ALL SELECT c.due_at FROM reporting_materializer_candidates c"
            " WHERE c.account_id=%s AND c.due_at IS NOT NULL AND NOT EXISTS (SELECT 1 FROM "
            + _OWNED
            + " w WHERE w.account_id=c.account_id AND w.consumer_id=c.consumer_id"  # nosec B608
            " AND w.delivery_config_id=c.delivery_config_id"
            " AND w.delivery_config_version=c.delivery_config_version"
            " AND w.reporting_obligation_id=c.reporting_obligation_id AND w.state='pending')"
            " UNION ALL SELECT clock_timestamp() FROM reporting_materializer_discovery"
            " WHERE account_id=%s AND NOT complete) ready) WHERE account_id=%s",
            (account_id,) * 5,
        )

    async def renew_materialization(
        self, lease: ReportingMaterializerLease, *, lease_seconds: int = 30
    ) -> bool:
        async with self._lease_epoch(lease):
            return await super().renew_materialization(lease, lease_seconds=lease_seconds)

    async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None:
        async with self._lease_epoch(lease):
            await super().authorize_materialization(lease)

    async def finish_materialization(
        self,
        lease: ReportingMaterializerLease,
        *,
        prepared: ReportingPreparedRevision | None = None,
        verified: ReportingVerifiedDestination | None = None,
        error: ReportingWriterFailure | None = None,
    ) -> ReportingMaterializerTurn:
        if lease.admission_epoch == 2:
            self._owner()._assert_components()
        async with self._lease_epoch(lease):
            return await super().finish_materialization(
                lease, prepared=prepared, verified=verified, error=error
            )

    async def read_production_boundaries(
        self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
    ) -> tuple[ReportingMaterializerBoundary, ...]:
        token = _EPOCH.set((id(self), 2))
        try:
            return await super().read_materializer_boundaries(
                caller=caller, after=after, limit=limit
            )
        finally:
            _EPOCH.reset(token)

The production store retains all ordinary legacy reader/writer APIs.

New work requires the live SDK composition and its mounted applicable routes. Pending epoch-zero work uses its original tables, identity and permanently quarantined enqueue. Installation alone admits nothing.

Ancestors

Methods

async def admit_production_configuration(self,
configuration: ReportingConfiguration,
binding: ReportingDestinationBinding,
*,
offering_id: str,
service_context: ReportingProductionSourceContext | None = None) ‑> None
Expand source code
async def admit_production_configuration(
    self,
    configuration: ReportingConfiguration,
    binding: ReportingDestinationBinding,
    *,
    offering_id: str,
    service_context: ReportingProductionSourceContext | None = None,
) -> None:
    offering = self._owner()._configuration_offering(
        configuration, binding, offering_id=offering_id
    )
    async with self._connection() as connection, connection.transaction():
        await self._lock_account(connection, configuration.account_id)
        token = _ADMISSION_CONNECTION.set((id(self), asyncio.current_task(), connection))
        try:
            await self.put_configuration(configuration)
            await self.put_destination_binding(binding)
            await self._enroll_on(
                connection,
                configuration,
                offering._producer_key,
                binding,
                service_context=service_context,
            )
        finally:
            _ADMISSION_CONNECTION.reset(token)
async def authorize_materialization(self, lease: ReportingMaterializerLease) ‑> None
Expand source code
async def authorize_materialization(self, lease: ReportingMaterializerLease) -> None:
    async with self._lease_epoch(lease):
        await super().authorize_materialization(lease)
async def commit_producer_period(self,
configuration: ReportingConfiguration,
obligation: ReportingObligationRecord,
*,
previous_end: datetime | None) ‑> ReportingObligationRecord
Expand source code
async def commit_producer_period(
    self,
    configuration: ReportingConfiguration,
    obligation: ReportingObligationRecord,
    *,
    previous_end: datetime | None,
) -> ReportingObligationRecord:
    check_next_period(configuration, obligation, previous_end)
    async with self._source_connection(configuration) as (connection, identity):
        row = await (
            await connection.execute(
                "SELECT closed_through FROM reporting_production_source_progress"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s",
                identity,
            )
        ).fetchone()
        current = row[0]
        if current != previous_end and (current is None or current < obligation.period.end):
            raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer progress changed")
        stored = await self.commit_obligation(obligation)
        await connection.execute(
            "INSERT INTO reporting_production_source_work"
            " (account_id,consumer_id,delivery_config_id,delivery_config_versi"
            "on,reporting_obligation_id,"
            " period_end) VALUES(%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING",
            (*identity, stored.reporting_obligation_id, stored.period.end),
        )
        await connection.execute(
            "UPDATE reporting_production_source_progress"
            " SET closed_through=greatest(closed_through,%s)"
            " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
            " AND delivery_config_version=%s",
            (stored.period.end, *identity),
        )
        return stored
async def finish_materialization(self,
lease: ReportingMaterializerLease,
*,
prepared: ReportingPreparedRevision | None = None,
verified: ReportingVerifiedDestination | None = None,
error: ReportingWriterFailure | None = None) ‑> ReportingMaterializerTurn
Expand source code
async def finish_materialization(
    self,
    lease: ReportingMaterializerLease,
    *,
    prepared: ReportingPreparedRevision | None = None,
    verified: ReportingVerifiedDestination | None = None,
    error: ReportingWriterFailure | None = None,
) -> ReportingMaterializerTurn:
    if lease.admission_epoch == 2:
        self._owner()._assert_components()
    async with self._lease_epoch(lease):
        return await super().finish_materialization(
            lease, prepared=prepared, verified=verified, error=error
        )
async def finish_producer_acquisition(self, configuration: ReportingConfiguration, *, reporting_obligation_id: str) ‑> None
Expand source code
async def finish_producer_acquisition(
    self, configuration: ReportingConfiguration, *, reporting_obligation_id: str
) -> None:
    async with self._source_connection(configuration) as (connection, identity):
        obligation = await self.get_obligation(
            account_id=identity[0], reporting_obligation_id=reporting_obligation_id
        )
        if obligation is None or obligation.generation_key != configuration.generation_key:
            raise LedgerConflictError("HISTORY_UNAVAILABLE", "producer generation differs")
        revisions = await self.list_revisions(
            account_id=identity[0], reporting_obligation_id=reporting_obligation_id
        )
        await connection.execute(
            "UPDATE reporting_production_source_work SET state=%s"
            " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
            " AND delivery_config_version=%s"
            " AND reporting_obligation_id=%s",
            (acquisition_state(obligation, revisions), *identity, reporting_obligation_id),
        )
async def lease_period_close(self, *, worker_id: str, now: datetime, lease_seconds: float) ‑> LeasedConfiguration | None
Expand source code
async def lease_period_close(
    self, *, worker_id: str, now: datetime, lease_seconds: float
) -> LeasedConfiguration | None:
    keys = self._owner()._producer_keys()
    if not keys:
        return None
    moment = _utc(now)
    expires = moment + timedelta(seconds=lease_seconds)
    if expires <= moment:
        raise ValueError("producer lease duration must be positive")
    if self._production_lease_samples is None:
        self._production_lease_samples = {}
    after = self._production_lease_samples.get(keys)
    following: _ProducerSample | None = None
    result = None
    async with self._connection() as connection, connection.transaction():
        await self._bound_lock_waits(connection)
        # Discover a bounded set without locking configuration rows. The
        # inherited configuration trigger takes the account lock, so that
        # lock must precede the row lock here, just as it does in activation.
        # A busy account cannot advance a durable rank while another
        # transaction holds its lock. Continue a read-only sample instead
        # of repeatedly trying the same prefix. At most one nonempty
        # 32-row window is examined; an empty tail may wrap once.
        for _ in range(2):
            continuation = (
                " AND (greatest(coalesce(t.lease_turn,0),coalesce(p.probe_turn,0)),"
                " coalesce(c.lease_expires_at,'-infinity'::timestamptz),"
                " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_config_version)"
                " > (%s,coalesce(%s::timestamptz,'-infinity'::timestamptz),%s,%s,%s,%s)"
                if after is not None
                else ""
            )
            query = (
                "SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.delivery"
                "_config_version,"  # nosec B608
                " greatest(coalesce(t.lease_turn,0),coalesce(p.probe_turn,0)),"
                " c.lease_expires_at"
                " FROM reporting_production_generations g"
                " JOIN reporting_production_accounts a ON a.account_id=g.account_id"
                " JOIN reporting_configurations c"
                " ON (c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_c"
                "onfig_version)="
                " (g.account_id,g.consumer_id,g.delivery_config_id,g.delivery_config_version)"
                " LEFT JOIN adcp_reporting_configuration_lease_turns t"
                " ON (t.account_id,t.consumer_id,t.delivery_config_id,t.delivery_c"
                "onfig_version)="
                " (g.account_id,g.consumer_id,g.delivery_config_id,g.delivery_config_version)"
                " LEFT JOIN reporting_production_source_probe_turns p"
                " ON (p.account_id,p.consumer_id,p.delivery_config_id,p.delivery_c"
                "onfig_version)="
                " (g.account_id,g.consumer_id,g.delivery_config_id,g.delivery_config_version)"
                " WHERE g.producer_key=ANY(%s)"
                " AND NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lease"
                "_expires_at<=%s)"
                + continuation
                + " ORDER BY greatest(coalesce(t.lease_turn,0),coalesce(p.probe_turn,0)),"
                " c.lease_expires_at NULLS FIRST,"
                " c.account_id,c.consumer_id,c.delivery_config_id,c.delivery_confi"
                "g_version LIMIT 32"
            )
            # Only the fixed SDK continuation above changes this SQL;
            # every cursor value and source identity remains a parameter.
            rows = await (
                await connection.execute(query, (list(keys), moment, *(after or ())))
            ).fetchall()
            if rows or after is None:
                break
            after = None
        for candidate in rows:
            row = candidate[:4]
            following = (candidate[4], candidate[5], *row)
            locked = await (
                await connection.execute(
                    "SELECT pg_try_advisory_xact_lock(hashtext('adcp.reporting:' || %s))",
                    (row[0],),
                )
            ).fetchone()
            if not locked[0]:
                continue
            current = await (
                await connection.execute(_SOURCE_CONFIGURATION, (*row, list(keys)))
            ).fetchone()
            if current is None:
                continue
            configuration = _configuration_from_row(current[:18])
            try:
                self._owner()._check_source_binding(configuration, current[18], current[19])
            except Exception:
                # A permanently revoked generation must not occupy the
                # first bounded window forever. This is a probe, not a
                # lease: preserve ordinary lease ranks and all source and
                # external identities. Its durable rank shares the same
                # ordering clock as successful configuration acquisitions.
                await connection.execute(
                    "INSERT INTO reporting_production_source_probe_turns"
                    " (account_id,consumer_id,delivery_config_id,delivery_config_versi"
                    "on,probe_turn)"
                    " VALUES(%s,%s,%s,%s,nextval('adcp_reporting_configuration_lease_t"
                    "urn_seq'))"
                    " ON CONFLICT(account_id,consumer_id,delivery_config_id,delivery_c"
                    "onfig_version)"
                    " DO UPDATE SET probe_turn="
                    "nextval('adcp_reporting_configuration_lease_turn_seq')",
                    tuple(row),
                )
                continue
            acquired = await (
                await connection.execute(
                    "UPDATE reporting_configurations SET lease_worker_id=%s,lease_expires_at=%s"
                    " WHERE (account_id,consumer_id,delivery_config_id,delivery_config"
                    "_version) = ("
                    " SELECT c.account_id,c.consumer_id,c.delivery_config_id,c.deliver"
                    "y_config_version"
                    " FROM reporting_configurations c WHERE c.account_id=%s"
                    " AND c.consumer_id=%s AND c.delivery_config_id=%s AND c.delivery_"
                    "config_version=%s"
                    " AND NOT c.quarantined AND (c.lease_expires_at IS NULL OR c.lease"
                    "_expires_at<=%s)"
                    " FOR UPDATE OF c SKIP LOCKED)"
                    " RETURNING account_id",
                    (worker_id, expires, *row, moment),
                )
            ).fetchone()
            if acquired is None:
                continue
            await self._retain_materializer_generation_on_lease_change(connection, tuple(row))
            await connection.execute(
                "INSERT INTO adcp_reporting_configuration_lease_turns"
                " (account_id,consumer_id,delivery_config_id,delivery_config_versi"
                "on,lease_turn)"
                " VALUES(%s,%s,%s,%s,nextval('adcp_reporting_configuration_lease_turn_seq'))"
                " ON CONFLICT(account_id,consumer_id,delivery_config_id,delivery_c"
                "onfig_version)"
                " DO UPDATE SET lease_turn="
                "nextval('adcp_reporting_configuration_lease_turn_seq')",
                tuple(row),
            )
            result = LeasedConfiguration(row[0], row[1], row[2], row[3], expires)
            break
    # Hints are per store and selected producer keys, not durable work or
    # leases. Publish a hint only after commit; a failed mutation retries
    # the same window. Successful acquisition returns to the durable
    # turn-primary order. A fresh store starts there too. Concurrent hints
    # may cause a bounded revisit, but cannot authorize or fence any work.
    if result is None and following is not None:
        self._production_lease_samples[keys] = following
    else:
        self._production_lease_samples.pop(keys, None)
    return result
async def materializer_ready(self) ‑> bool
Expand source code
async def materializer_ready(self) -> bool:
    owner = self._owner()
    if not await owner._schema_ready():
        raise ReportingNotificationError("reporting_production_schema_unready")
    owner._assert_components()
    return True
async def next_producer_obligations(self, configuration: ReportingConfiguration, *, now: datetime, limit: int) ‑> tuple[str, ...]
Expand source code
async def next_producer_obligations(
    self, configuration: ReportingConfiguration, *, now: datetime, limit: int
) -> tuple[str, ...]:
    if type(limit) is not int or not 1 <= limit <= 64:
        raise ValueError("production acquisition limit must be in 1..64")
    async with self._source_connection(configuration) as (connection, identity):
        rows = await (
            await connection.execute(
                "SELECT reporting_obligation_id FROM reporting_production_source_work"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s"
                " AND state='pending' AND period_end<=%s"
                " ORDER BY acquisition_turn,reporting_obligation_id LIMIT %s FOR UPDATE",
                (*identity, now, limit),
            )
        ).fetchall()
        if not rows:
            return ()
        head = await (
            await connection.execute(
                "UPDATE reporting_production_source_progress"
                " SET acquisition_turn=acquisition_turn+%s"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s"
                " RETURNING acquisition_turn",
                (len(rows), *identity),
            )
        ).fetchone()
        for offset, row in enumerate(rows, 1):
            await connection.execute(
                "UPDATE reporting_production_source_work SET acquisition_turn=%s"
                " WHERE account_id=%s AND reporting_obligation_id=%s",
                (head[0] - len(rows) + offset, identity[0], row[0]),
            )
        return tuple(row[0] for row in rows)
async def producer_closed_through(self, configuration: ReportingConfiguration) ‑> datetime.datetime | None
Expand source code
async def producer_closed_through(
    self, configuration: ReportingConfiguration
) -> datetime | None:
    async with self._source_connection(configuration) as (connection, identity):
        row = await (
            await connection.execute(
                "SELECT closed_through FROM reporting_production_source_progress"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s",
                identity,
            )
        ).fetchone()
        return row[0] if row is not None else None
async def producer_constituents(self,
configuration: ReportingConfiguration,
obligation: ReportingObligationRecord) ‑> tuple[ProductConstituentV1 | PackageItemConstituentV1 | MediaBuyConstituentV1, ...]
Expand source code
async def producer_constituents(
    self, configuration: ReportingConfiguration, obligation: ReportingObligationRecord
) -> tuple[ReportingConstituent, ...]:
    async with self._source_connection(configuration) as (connection, identity):
        if obligation.generation_key != configuration.generation_key or set(
            obligation.media_buy_ids
        ) != set(configuration.media_buy_ids):
            raise LedgerConflictError("HISTORY_UNAVAILABLE", "source denominator differs")
        row = await (
            await connection.execute(
                "SELECT producer_key,source_binding FROM reporting_production_generations"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s",
                identity,
            )
        ).fetchone()
        binding = self._owner()._check_source_binding(configuration, row[0], row[1])
        return binding.constituents()
async def read_production_boundaries(self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100) ‑> tuple[ReportingMaterializerBoundary, ...]
Expand source code
async def read_production_boundaries(
    self, *, caller: ReportingDeliveryPrincipal, after: int = 0, limit: int = 100
) -> tuple[ReportingMaterializerBoundary, ...]:
    token = _EPOCH.set((id(self), 2))
    try:
        return await super().read_materializer_boundaries(
            caller=caller, after=after, limit=limit
        )
    finally:
        _EPOCH.reset(token)
async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) ‑> None
Expand source code
async def release_period_close(self, lease: LeasedConfiguration, *, worker_id: str) -> None:
    async with self._connection() as connection, connection.transaction():
        await self._lock_account(connection, lease.account_id)
        identity = (
            lease.account_id,
            lease.consumer_id,
            lease.delivery_config_id,
            lease.delivery_config_version,
        )
        released = await (
            await connection.execute(
                "UPDATE reporting_configurations SET lease_worker_id=NULL,lease_expires_at=NULL"
                " WHERE account_id=%s AND consumer_id=%s AND delivery_config_id=%s"
                " AND delivery_config_version=%s"
                " AND lease_worker_id=%s AND lease_expires_at=%s RETURNING account_id",
                (*identity, worker_id, lease.lease_expires_at),
            )
        ).fetchone()
        if released is not None:
            await self._retain_materializer_generation_on_lease_change(connection, identity)
async def renew_materialization(self, lease: ReportingMaterializerLease, *, lease_seconds: int = 30) ‑> bool
Expand source code
async def renew_materialization(
    self, lease: ReportingMaterializerLease, *, lease_seconds: int = 30
) -> bool:
    async with self._lease_epoch(lease):
        return await super().renew_materialization(lease, lease_seconds=lease_seconds)

Inherited members

class ReportingConfigurationAdmission (offering_id: str,
configuration: ReportingConfiguration,
binding: ReportingDestinationBinding,
*,
configuration_wire: Mapping[str, Any])
Expand source code
@dataclass(frozen=True)
class ReportingConfigurationAdmission:
    """Resolved by trusted account/provider code, never decoded from buyer JSON.

    The bound references are opaque; credentials stay in destination sessions.
    The SDK validates the complete frozen tuple before admitting configuration.
    """

    offering_id: str
    configuration: ReportingConfiguration
    binding: ReportingDestinationBinding
    configuration_wire: Mapping[str, Any] = field(kw_only=True, repr=False, compare=False)
    _wire: bytes = field(init=False, repr=False)

    def __post_init__(self) -> None:
        from adcp.validation.schema_loader import get_named_validator

        wire = canonical_json_utf8_v1(dict(self.configuration_wire))
        validator = get_named_validator("core/reporting-delivery-config.json")
        if validator is None or next(validator.iter_errors(json.loads(wire)), None) is not None:
            raise ValueError("configuration must satisfy the complete public contract")
        object.__setattr__(self, "_wire", wire)

    def wire(self) -> dict[str, Any]:
        return dict(json.loads(self._wire))

    def check(self, support: ReportingProductionSupport) -> None:
        from adcp.reporting.ledger.store import reject_reserved_authoritative_party

        config, binding, raw = self.configuration, self.binding, self.wire()
        reject_reserved_authoritative_party(config)
        if raw.get("authoritative_party", "seller") != "seller":
            raise LedgerConflictError("UNSUPPORTED_FEATURE", "consumer authority is reserved")
        offering = support._configuration_offering(config, binding, offering_id=self.offering_id)
        schedule = offering.configuration_schedule(config)
        supplied_schedule = dict(raw["schedule"])
        if "period_anchor" in supplied_schedule:
            supplied_schedule["period_anchor"] = aware_timestamp(
                supplied_schedule["period_anchor"]
            ).isoformat()
            schedule["period_anchor"] = aware_timestamp(schedule["period_anchor"]).isoformat()
        scope = raw["scope"]
        # This producer freezes an explicit full media-buy denominator. An
        # application's dynamic all-buy or partial-coverage implementation must
        # not be represented as this installed exact-generation contract.
        if (
            raw["delivery_config_id"] != config.delivery_config_id
            or raw["delivery_config_version"] != config.delivery_config_version
            or raw["offering_id"] != self.offering_id
            or raw["feed_purpose"] != config.feed_purpose
            or raw["report_definition_id"] != config.report_definition_id
            or raw["reporting_profile"] != config.reporting_profile
            or raw["required_finality"] != config.required_finality
            or raw["reconciliation_mode"] != binding.reconciliation_mode
            or set(scope) != {"media_buy_ids"}
            or set(scope["media_buy_ids"]) != set(config.media_buy_ids)
            or raw["coverage_requirement"] != "full"
            or supplied_schedule != schedule
            or raw.get("method") != support._destination_binding(binding, offering).wire()
            or raw["active"] != (config.activated_at is not None and config.deactivated_at is None)
            or (
                "revocation_effective_at" in raw
                and aware_timestamp(raw["revocation_effective_at"]) != config.deactivated_at
            )
        ):
            raise failure("BINDING_MISMATCH")

Resolved by trusted account/provider code, never decoded from buyer JSON.

The bound references are opaque; credentials stay in destination sessions. The SDK validates the complete frozen tuple before admitting configuration.

Instance variables

var binding : ReportingDestinationBinding
var configuration : ReportingConfiguration
var configuration_wire : Mapping[str, typing.Any]
var offering_id : str

Methods

def check(self,
support: ReportingProductionSupport) ‑> None
Expand source code
def check(self, support: ReportingProductionSupport) -> None:
    from adcp.reporting.ledger.store import reject_reserved_authoritative_party

    config, binding, raw = self.configuration, self.binding, self.wire()
    reject_reserved_authoritative_party(config)
    if raw.get("authoritative_party", "seller") != "seller":
        raise LedgerConflictError("UNSUPPORTED_FEATURE", "consumer authority is reserved")
    offering = support._configuration_offering(config, binding, offering_id=self.offering_id)
    schedule = offering.configuration_schedule(config)
    supplied_schedule = dict(raw["schedule"])
    if "period_anchor" in supplied_schedule:
        supplied_schedule["period_anchor"] = aware_timestamp(
            supplied_schedule["period_anchor"]
        ).isoformat()
        schedule["period_anchor"] = aware_timestamp(schedule["period_anchor"]).isoformat()
    scope = raw["scope"]
    # This producer freezes an explicit full media-buy denominator. An
    # application's dynamic all-buy or partial-coverage implementation must
    # not be represented as this installed exact-generation contract.
    if (
        raw["delivery_config_id"] != config.delivery_config_id
        or raw["delivery_config_version"] != config.delivery_config_version
        or raw["offering_id"] != self.offering_id
        or raw["feed_purpose"] != config.feed_purpose
        or raw["report_definition_id"] != config.report_definition_id
        or raw["reporting_profile"] != config.reporting_profile
        or raw["required_finality"] != config.required_finality
        or raw["reconciliation_mode"] != binding.reconciliation_mode
        or set(scope) != {"media_buy_ids"}
        or set(scope["media_buy_ids"]) != set(config.media_buy_ids)
        or raw["coverage_requirement"] != "full"
        or supplied_schedule != schedule
        or raw.get("method") != support._destination_binding(binding, offering).wire()
        or raw["active"] != (config.activated_at is not None and config.deactivated_at is None)
        or (
            "revocation_effective_at" in raw
            and aware_timestamp(raw["revocation_effective_at"]) != config.deactivated_at
        )
    ):
        raise failure("BINDING_MISMATCH")
def wire(self) ‑> dict[str, typing.Any]
Expand source code
def wire(self) -> dict[str, Any]:
    return dict(json.loads(self._wire))
class ReportingProductionConfigurationTask (handle: ConfigurationTask, account: AccountCapabilities)
Expand source code
@dataclass(frozen=True)
class ReportingProductionConfigurationTask:
    """Compose the account implementation with enforced SDK reporting admission.

    ``handle`` remains the application's authenticated, caller-owned desired
    state account task. It resolves provider grants and opaque bindings, calls
    the supplied ``admit`` for every accepted reporting generation, and returns
    the normal account response. The wrapper verifies returned ready states
    against the actual stored records and that call's validated admissions.
    A task cannot return a successful unsupported promise without a matching
    admitted configuration, including on exact replay.
    """

    handle: ConfigurationTask = field(repr=False)
    account: AccountCapabilities = field(repr=False, compare=False)
    _account_wire: bytes = field(init=False, repr=False)

    def __post_init__(self) -> None:
        from adcp.validation.schema_loader import get_named_validator

        raw = self.account.model_dump(mode="json", exclude_none=True, exclude_unset=True)
        validator = get_named_validator("protocol/get-adcp-capabilities-response.json")
        if (
            validator is None
            or next(
                validator.evolve(schema=validator.schema["properties"]["account"]).iter_errors(raw),
                None,
            )
            is not None
            or raw.get("account_financials") is True
            or any(
                raw.get(name, {}).get("supported") is True
                for name in ("notifications", "change_feed", "identity_updates")
            )
        ):
            raise ValueError(
                "configuration task requires its actual account capabilities and mounted operations"
            )
        object.__setattr__(self, "_account_wire", canonical_json_utf8_v1(raw))

    def account_capabilities(self) -> dict[str, Any]:
        return dict(json.loads(self._account_wire))

    async def execute(
        self,
        support: ReportingProductionSupport,
        request: dict[str, Any],
        context: ToolContext | None,
    ) -> dict[str, Any]:
        from adcp.validation.schema_loader import get_named_validator

        validator = get_named_validator("account/sync-accounts-request.json")
        if validator is None or next(validator.iter_errors(request), None) is not None:
            raise LedgerConflictError("INVALID_REQUEST", "account request is invalid")
        for entry in request["accounts"]:
            configurations = entry.get("reporting_delivery_configs", ())
            keys = [(c["delivery_config_id"], c["delivery_config_version"]) for c in configurations]
            if len(set(keys)) != len(keys):
                raise LedgerConflictError(
                    "INVALID_REQUEST", "configuration generations must be unique"
                )
        admitted: dict[tuple[str, str, str, int], ReportingConfigurationAdmission] = {}

        async def admit(value: ReportingConfigurationAdmission) -> None:
            if request.get("dry_run") is True:
                raise LedgerConflictError(
                    "INVALID_REQUEST", "dry runs cannot admit reporting state"
                )
            if type(value) is not ReportingConfigurationAdmission or context is None:
                raise LedgerConflictError("UNAUTHORIZED", "configuration authority is unavailable")
            who = await support.handler._authorize(
                {"account": {"account_id": value.configuration.account_id}}, context
            )
            if who != value.binding.principal:
                raise LedgerConflictError("UNAUTHORIZED", "configuration authority is unavailable")
            matched = False
            for entry in request["accounts"]:
                if "reporting_delivery_configs" not in entry or "account" not in entry:
                    continue
                try:
                    requested_caller = await support.handler._authorize(
                        {"account": entry["account"]}, context
                    )
                # An unresolved account cannot authorize this admission.
                except Exception:  # nosec B112
                    continue
                if requested_caller != who:
                    continue
                desired = entry["reporting_delivery_configs"]
                matched = value.wire() in desired or (
                    value.configuration.deactivated_at is not None
                    and not any(
                        (c["delivery_config_id"], c["delivery_config_version"])
                        == (
                            value.configuration.delivery_config_id,
                            value.configuration.delivery_config_version,
                        )
                        for c in desired
                    )
                )
                if matched:
                    break
            if not matched:
                raise LedgerConflictError(
                    "INVALID_REQUEST", "configuration differs from requested account state"
                )
            support.validate_configuration(value)
            await support._check_notifications(who.account_id)
            key = (
                who.account_id,
                who.consumer_id,
                value.configuration.delivery_config_id,
                value.configuration.delivery_config_version,
            )
            if key in admitted and admitted[key] != value:
                raise LedgerConflictError(
                    "CONFIGURATION_GENERATION_IMMUTABLE", "configuration identity conflicts"
                )
            await support._admit_configuration(value)
            # The caller already entered the migrated, drained production
            # lifecycle. Complete this account's versioned baseline before
            # echoing ready; adopters need no account-enumeration worker or
            # hand-maintained readiness flag. A failed activation remains
            # unready and can resume under the same durable identities.
            await support.activate(account_id=who.account_id)
            admitted[key] = value

        response = await self.handle(dict(request), context, admit)
        if not isinstance(response, Mapping):
            raise LedgerConflictError(
                "INVALID_REQUEST", "configuration task returned an invalid response"
            )
        if response.get("errors") or request.get("dry_run") is True:
            return dict(response)
        for account in response.get("accounts", ()):
            if account.get("action") == "failed":
                continue
            response_owner = await support.handler._authorize(
                {"account": {"account_id": account.get("account_id")}}, context
            )
            for state in account.get("reporting_delivery_configs", ()):
                if state.get("state") not in {"ready", "inactive"}:
                    continue
                config = state.get("configuration", {})
                key = (
                    response_owner.account_id,
                    response_owner.consumer_id,
                    config.get("delivery_config_id"),
                    config.get("delivery_config_version"),
                )
                value = admitted.get(key)
                if value is None:
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED", "configuration requires SDK admission"
                    )
                validator = get_named_validator("core/reporting-delivery-config-state.json")
                if (
                    validator is None
                    or next(validator.iter_errors(state), None) is not None
                    or config != value.wire()
                ):
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED",
                        "configuration state differs from admission",
                    )
                expected_state = (
                    "inactive" if value.configuration.deactivated_at is not None else "ready"
                )
                if state["state"] != expected_state:
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED", "configuration lifecycle differs"
                    )
                for name in ("activated_at", "deactivated_at"):
                    actual_time = getattr(value.configuration, name)
                    supplied = state.get(name)
                    if supplied is not None and aware_timestamp(supplied) != actual_time:
                        raise LedgerConflictError(
                            "REPORTING_CONFIGURATION_UNADMITTED", "configuration lifecycle differs"
                        )
                coverage = state.get("current_coverage")
                if state["state"] == "ready" and (
                    not isinstance(coverage, dict)
                    or coverage["status"] != "full"
                    or set(coverage["media_buy_ids"]) != set(value.configuration.media_buy_ids)
                    or set(coverage["fully_covered_media_buy_ids"])
                    != set(value.configuration.media_buy_ids)
                    or any(
                        coverage[k]
                        for k in (
                            "partially_covered_media_buy_ids",
                            "unsupported_media_buy_ids",
                            "unknown_media_buy_ids",
                            "unsupported_package_ids",
                            "unknown_package_ids",
                            "limitations",
                        )
                    )
                    or set(coverage["package_ids"]) != set(coverage["covered_package_ids"])
                ):
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED", "configuration coverage differs"
                    )
                if state["state"] == "inactive":
                    from adcp.reporting.ledger.models import derive_period, first_ordinal_after

                    configuration = value.configuration
                    assert configuration.deactivated_at is not None
                    ordinal = first_ordinal_after(
                        configuration.schedule,
                        account_timezone=configuration.account_timezone,
                        activated_at=configuration.deactivated_at,
                    )
                    cutoff = derive_period(
                        configuration.schedule,
                        account_timezone=configuration.account_timezone,
                        ordinal=ordinal,
                    ).start
                    if aware_timestamp(state["publication_stopped_at"]) != cutoff:
                        raise LedgerConflictError(
                            "REPORTING_CONFIGURATION_UNADMITTED", "configuration cutoff differs"
                        )
                who = ReportingDeliveryPrincipal(
                    value.configuration.account_id, value.binding.consumer_id
                )
                actual = await support.store.get_destination_binding(
                    caller=who, generation_key=value.configuration.generation_key
                )
                configs = await support.store.list_configurations(caller=who)
                support.validate_configuration(value)
                if actual != value.binding or value.configuration not in configs:
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED",
                        "configuration admission is unavailable",
                    )
                if state.get("destination_ref") != actual.destination_ref:
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED", "configuration destination differs"
                    )
        return dict(response)

Compose the account implementation with enforced SDK reporting admission.

handle remains the application's authenticated, caller-owned desired state account task. It resolves provider grants and opaque bindings, calls the supplied admit for every accepted reporting generation, and returns the normal account response. The wrapper verifies returned ready states against the actual stored records and that call's validated admissions. A task cannot return a successful unsupported promise without a matching admitted configuration, including on exact replay.

Instance variables

var account : Account
var handle : Callable[[dict[str, typing.Any], ToolContext | None, Callable[[ReportingConfigurationAdmission], Awaitable[None]]], Awaitable[dict[str, typing.Any]]]

Methods

def account_capabilities(self) ‑> dict[str, typing.Any]
Expand source code
def account_capabilities(self) -> dict[str, Any]:
    return dict(json.loads(self._account_wire))
async def execute(self,
support: ReportingProductionSupport,
request: dict[str, Any],
context: ToolContext | None) ‑> dict[str, Any]
Expand source code
async def execute(
    self,
    support: ReportingProductionSupport,
    request: dict[str, Any],
    context: ToolContext | None,
) -> dict[str, Any]:
    from adcp.validation.schema_loader import get_named_validator

    validator = get_named_validator("account/sync-accounts-request.json")
    if validator is None or next(validator.iter_errors(request), None) is not None:
        raise LedgerConflictError("INVALID_REQUEST", "account request is invalid")
    for entry in request["accounts"]:
        configurations = entry.get("reporting_delivery_configs", ())
        keys = [(c["delivery_config_id"], c["delivery_config_version"]) for c in configurations]
        if len(set(keys)) != len(keys):
            raise LedgerConflictError(
                "INVALID_REQUEST", "configuration generations must be unique"
            )
    admitted: dict[tuple[str, str, str, int], ReportingConfigurationAdmission] = {}

    async def admit(value: ReportingConfigurationAdmission) -> None:
        if request.get("dry_run") is True:
            raise LedgerConflictError(
                "INVALID_REQUEST", "dry runs cannot admit reporting state"
            )
        if type(value) is not ReportingConfigurationAdmission or context is None:
            raise LedgerConflictError("UNAUTHORIZED", "configuration authority is unavailable")
        who = await support.handler._authorize(
            {"account": {"account_id": value.configuration.account_id}}, context
        )
        if who != value.binding.principal:
            raise LedgerConflictError("UNAUTHORIZED", "configuration authority is unavailable")
        matched = False
        for entry in request["accounts"]:
            if "reporting_delivery_configs" not in entry or "account" not in entry:
                continue
            try:
                requested_caller = await support.handler._authorize(
                    {"account": entry["account"]}, context
                )
            # An unresolved account cannot authorize this admission.
            except Exception:  # nosec B112
                continue
            if requested_caller != who:
                continue
            desired = entry["reporting_delivery_configs"]
            matched = value.wire() in desired or (
                value.configuration.deactivated_at is not None
                and not any(
                    (c["delivery_config_id"], c["delivery_config_version"])
                    == (
                        value.configuration.delivery_config_id,
                        value.configuration.delivery_config_version,
                    )
                    for c in desired
                )
            )
            if matched:
                break
        if not matched:
            raise LedgerConflictError(
                "INVALID_REQUEST", "configuration differs from requested account state"
            )
        support.validate_configuration(value)
        await support._check_notifications(who.account_id)
        key = (
            who.account_id,
            who.consumer_id,
            value.configuration.delivery_config_id,
            value.configuration.delivery_config_version,
        )
        if key in admitted and admitted[key] != value:
            raise LedgerConflictError(
                "CONFIGURATION_GENERATION_IMMUTABLE", "configuration identity conflicts"
            )
        await support._admit_configuration(value)
        # The caller already entered the migrated, drained production
        # lifecycle. Complete this account's versioned baseline before
        # echoing ready; adopters need no account-enumeration worker or
        # hand-maintained readiness flag. A failed activation remains
        # unready and can resume under the same durable identities.
        await support.activate(account_id=who.account_id)
        admitted[key] = value

    response = await self.handle(dict(request), context, admit)
    if not isinstance(response, Mapping):
        raise LedgerConflictError(
            "INVALID_REQUEST", "configuration task returned an invalid response"
        )
    if response.get("errors") or request.get("dry_run") is True:
        return dict(response)
    for account in response.get("accounts", ()):
        if account.get("action") == "failed":
            continue
        response_owner = await support.handler._authorize(
            {"account": {"account_id": account.get("account_id")}}, context
        )
        for state in account.get("reporting_delivery_configs", ()):
            if state.get("state") not in {"ready", "inactive"}:
                continue
            config = state.get("configuration", {})
            key = (
                response_owner.account_id,
                response_owner.consumer_id,
                config.get("delivery_config_id"),
                config.get("delivery_config_version"),
            )
            value = admitted.get(key)
            if value is None:
                raise LedgerConflictError(
                    "REPORTING_CONFIGURATION_UNADMITTED", "configuration requires SDK admission"
                )
            validator = get_named_validator("core/reporting-delivery-config-state.json")
            if (
                validator is None
                or next(validator.iter_errors(state), None) is not None
                or config != value.wire()
            ):
                raise LedgerConflictError(
                    "REPORTING_CONFIGURATION_UNADMITTED",
                    "configuration state differs from admission",
                )
            expected_state = (
                "inactive" if value.configuration.deactivated_at is not None else "ready"
            )
            if state["state"] != expected_state:
                raise LedgerConflictError(
                    "REPORTING_CONFIGURATION_UNADMITTED", "configuration lifecycle differs"
                )
            for name in ("activated_at", "deactivated_at"):
                actual_time = getattr(value.configuration, name)
                supplied = state.get(name)
                if supplied is not None and aware_timestamp(supplied) != actual_time:
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED", "configuration lifecycle differs"
                    )
            coverage = state.get("current_coverage")
            if state["state"] == "ready" and (
                not isinstance(coverage, dict)
                or coverage["status"] != "full"
                or set(coverage["media_buy_ids"]) != set(value.configuration.media_buy_ids)
                or set(coverage["fully_covered_media_buy_ids"])
                != set(value.configuration.media_buy_ids)
                or any(
                    coverage[k]
                    for k in (
                        "partially_covered_media_buy_ids",
                        "unsupported_media_buy_ids",
                        "unknown_media_buy_ids",
                        "unsupported_package_ids",
                        "unknown_package_ids",
                        "limitations",
                    )
                )
                or set(coverage["package_ids"]) != set(coverage["covered_package_ids"])
            ):
                raise LedgerConflictError(
                    "REPORTING_CONFIGURATION_UNADMITTED", "configuration coverage differs"
                )
            if state["state"] == "inactive":
                from adcp.reporting.ledger.models import derive_period, first_ordinal_after

                configuration = value.configuration
                assert configuration.deactivated_at is not None
                ordinal = first_ordinal_after(
                    configuration.schedule,
                    account_timezone=configuration.account_timezone,
                    activated_at=configuration.deactivated_at,
                )
                cutoff = derive_period(
                    configuration.schedule,
                    account_timezone=configuration.account_timezone,
                    ordinal=ordinal,
                ).start
                if aware_timestamp(state["publication_stopped_at"]) != cutoff:
                    raise LedgerConflictError(
                        "REPORTING_CONFIGURATION_UNADMITTED", "configuration cutoff differs"
                    )
            who = ReportingDeliveryPrincipal(
                value.configuration.account_id, value.binding.consumer_id
            )
            actual = await support.store.get_destination_binding(
                caller=who, generation_key=value.configuration.generation_key
            )
            configs = await support.store.list_configurations(caller=who)
            support.validate_configuration(value)
            if actual != value.binding or value.configuration not in configs:
                raise LedgerConflictError(
                    "REPORTING_CONFIGURATION_UNADMITTED",
                    "configuration admission is unavailable",
                )
            if state.get("destination_ref") != actual.destination_ref:
                raise LedgerConflictError(
                    "REPORTING_CONFIGURATION_UNADMITTED", "configuration destination differs"
                )
    return dict(response)
class ReportingProductionDestination (*args, **kwargs)
Expand source code
@runtime_checkable
class ReportingProductionDestination(
    ReportingDestinationWriter, ReportingDestinationResolver, Protocol
):
    """An actual provider's bounded retention and revocation commitments.

    These promises accompany the writer's exact capabilities. Every write and
    independent readback still opens a separately authorized SDK session. A
    service cannot make the development writer eligible by decorating it.
    """

    @property
    def resource_retention_days(self) -> int: ...

    @property
    def authorization_revocation_seconds(self) -> int: ...

    @property
    def delivery_methods(self) -> tuple[ReportingProductionMethod, ...]: ...

    def configuration_binding(
        self, binding: ReportingDestinationBinding
    ) -> ReportingProductionDestinationBinding | None: ...

An actual provider's bounded retention and revocation commitments.

These promises accompany the writer's exact capabilities. Every write and independent readback still opens a separately authorized SDK session. A service cannot make the development writer eligible by decorating it.

Ancestors

Instance variables

prop authorization_revocation_seconds : int
Expand source code
@property
def authorization_revocation_seconds(self) -> int: ...
prop delivery_methods : tuple[ReportingProductionMethod, ...]
Expand source code
@property
def delivery_methods(self) -> tuple[ReportingProductionMethod, ...]: ...
prop resource_retention_days : int
Expand source code
@property
def resource_retention_days(self) -> int: ...

Methods

def configuration_binding(self, binding: ReportingDestinationBinding) ‑> ReportingProductionDestinationBinding | None
Expand source code
def configuration_binding(
    self, binding: ReportingDestinationBinding
) -> ReportingProductionDestinationBinding | None: ...
class ReportingProductionDestinationBinding (binding: ReportingDestinationBinding,
method: ReportingProductionMethod,
configuration: Mapping[str, Any])
Expand source code
@dataclass(frozen=True)
class ReportingProductionDestinationBinding:
    """The provider's authorized, immutable destination contract for a caller.

    ``configuration`` is the complete secret-free public delivery method,
    including the selected destination. The provider resolves it from its
    trusted binding, not from a buyer's claim. Credentials remain in the
    independently authorized write/readback sessions. The SDK freezes these
    bytes at admission and compares them again before materialization.
    """

    binding: ReportingDestinationBinding = field(repr=False)
    method: ReportingProductionMethod
    configuration: Mapping[str, Any] = field(repr=False, compare=False)
    _wire: bytes = field(init=False, repr=False)

    def __post_init__(self) -> None:
        from adcp.reporting.evidence import consumer_reference, resource_location
        from adcp.validation.schema_loader import get_named_validator

        if (
            type(self.binding) is not ReportingDestinationBinding
            or type(self.method) is not ReportingProductionMethod
        ):
            raise ValueError("destination requires the exact provider binding and method")
        raw = json.loads(canonical_json_utf8_v1(dict(self.configuration)))
        validator = get_named_validator("core/reporting-delivery-method.json")
        offered = self.method.wire()
        if (
            validator is None
            or next(validator.iter_errors(raw), None) is not None
            or any(raw.get(k) != offered.get(k) for k in ("pattern", "transport", "orchestration"))
            or (raw["pattern"] == "file_transfer" and raw["format"] != offered.get("format"))
            or raw["destination"]["mode"] not in offered["destination_modes"]
        ):
            raise ValueError("destination must match the provider's complete method")
        destination = raw["destination"]
        if destination["mode"] == "existing":
            if destination["destination_ref"] != self.binding.destination_ref:
                raise ValueError("destination must match the immutable binding")
        elif destination.get("provider") != offered.get("provider") or destination.get(
            "access_mode"
        ) != offered.get("access_mode"):
            raise ValueError("destination must match the provider's complete method")
        if "location" in destination:
            resource_location(destination["location"])
        if "recipient" in destination:
            consumer_reference(destination["recipient"]["identity"])
        object.__setattr__(self, "_wire", canonical_json_utf8_v1(raw))

    def wire(self) -> dict[str, Any]:
        return dict(json.loads(self._wire))

The provider's authorized, immutable destination contract for a caller.

adcp.reporting.production.configuration is the complete secret-free public delivery method, including the selected destination. The provider resolves it from its trusted binding, not from a buyer's claim. Credentials remain in the independently authorized write/readback sessions. The SDK freezes these bytes at admission and compares them again before materialization.

Instance variables

var binding : ReportingDestinationBinding
var configuration : Mapping[str, typing.Any]
var method : ReportingProductionMethod

Methods

def wire(self) ‑> dict[str, typing.Any]
Expand source code
def wire(self) -> dict[str, Any]:
    return dict(json.loads(self._wire))
class ReportingProductionHandler (production: ReportingProductionSupport,
*,
resolve_account: ReceiptAccountResolver,
buyer_agents: BuyerAgentRegistry | None = None,
adcp_version: str | None = None)
Expand source code
class ReportingProductionHandler(ReportingReceiptHandler):
    """Mount the same instance on MCP/A2A; the support owns its lifecycle."""

    advertised_tools = _REPORTING_TASKS | _APPLICATION_TASKS

    def __init__(
        self,
        production: ReportingProductionSupport,
        *,
        resolve_account: ReceiptAccountResolver,
        buyer_agents: BuyerAgentRegistry | None = None,
        adcp_version: str | None = None,
    ) -> None:
        self.production = production
        self._application: ADCPHandler[Any] | None = None
        self._application_tools: frozenset[str] = frozenset()
        super().__init__(
            production.store,
            resolve_account=resolve_account,
            buyer_agents=buyer_agents,
            consumer_status_enabled=production.projection.consumer_status_enabled,
            adcp_version=adcp_version,
        )

    def bind_application(self, application: ADCPHandler[Any]) -> None:
        """Delegate ordinary tasks before mounting this exact SDK handler.

        Reporting/configuration collisions are refused. Supply the typed B2
        configuration task separately; aggregate delivery and base capabilities
        may coexist and are dispatched explicitly by their protected methods.
        """
        if application is self or not isinstance(application, ADCPHandler):
            raise ValueError("application must be a separate ADCPHandler")
        if self._application is application:
            return
        if (
            self._application is not None
            or self.production._mounts
            or self.production._task is not None
        ):
            raise ValueError("application delegation must be fixed before mounting")
        names = {t["name"] for t in get_tools_for_handler(application, _include_schemas=False)}
        if names & (_REPORTING_TASKS - {"get_adcp_capabilities", "get_media_buy_delivery"}):
            raise ValueError("application reporting handlers conflict with production ownership")
        version = getattr(application, "get_adcp_version", None)
        if callable(version) and version() != self.get_adcp_version():
            raise ValueError("application and production protocol versions differ")
        self._application, self._application_tools = application, frozenset(names)

    async def _delegate(self, task: str, params: Any, context: ToolContext | None) -> Any:
        if self._application is None or task not in self._application_tools:
            return self._not_supported(task)
        result = getattr(self._application, _APPLICATION_METHODS.get(task, task))(params, context)
        return await result if inspect.isawaitable(result) else result

    def advertised_tools_for_instance(self) -> set[str]:
        names = {
            "get_adcp_capabilities",
            "get_reporting_status",
            "get_media_buy_delivery",
            "sync_accounts",
        }
        if any(o.reconciled for o in self.production.offerings):
            names.add("sync_reporting_receipts")
        if self._feed_consumer_status_enabled:
            names.add("sync_reporting_status")
        return names | set(self._application_tools & _APPLICATION_TASKS)

    @_admitted
    async def get_reporting_status(
        self,
        params: GetReportingStatusRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any] | NotImplementedResponse:
        return await super().get_reporting_status(params, context)

    @_admitted
    async def sync_reporting_receipts(
        self,
        params: SyncReportingReceiptsRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any]:
        if not any(o.reconciled for o in self.production.offerings):
            raise _task_error(
                "sync_reporting_receipts", "NOT_SUPPORTED", "receipt task is unavailable"
            )
        return await super().sync_reporting_receipts(params, context)

    async def _authorize(
        self, request: dict[str, Any], context: ToolContext | None
    ) -> ReportingDeliveryPrincipal:
        try:
            if context is None or not isinstance(request.get("account"), dict):
                raise ReportingReceiptError("UNAUTHORIZED")
            consumer = await _consumer(context, self._receipt_registry)
            account = await self._receipt_account_resolver(
                dict(request["account"]), context, consumer
            )
            if isinstance(context, RequestContext) and context.account.id != account:
                raise ReportingReceiptError("UNAUTHORIZED")
            return ReportingDeliveryPrincipal(account, consumer)
        except (ReliableReportingConfigurationError, ReliableReportingUnavailableError):
            raise
        except Exception:
            raise ReportingReceiptError("UNAUTHORIZED") from None

    async def get_adcp_capabilities(
        self,
        params: GetAdcpCapabilitiesRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any]:
        response = capabilities_response(
            ["media_buy"],
            sandbox=False,
            idempotency={"supported": False},
            adcp_version=self.production._protocol_version,
            supported_versions=[self.production._protocol_version],
        )
        if self._application is not None and "get_adcp_capabilities" in self._application_tools:
            protocol = response["adcp"]
            result = await self._delegate("get_adcp_capabilities", params, context)
            base = {} if isinstance(result, NotImplementedResponse) else _request(result)
            if (base.get("media_buy") or {}).get("reporting_delivery") is not None:
                raise ValueError(
                    "application reporting capabilities conflict with production ownership"
                )
            response.update(base)
            response["adcp_version"] = self.production._protocol_version
            response["adcp"] = {
                **response["adcp"],
                "major_versions": protocol["major_versions"],
                "supported_versions": protocol["supported_versions"],
            }
            response["supported_protocols"] = list(
                dict.fromkeys([*base.get("supported_protocols", ()), "media_buy"])
            )
        response["account"] = self.production.configuration_task.account_capabilities()
        declaration = self.production._declared_capabilities()
        reporting = declaration["reporting_delivery"]
        if reporting:
            response["media_buy"] = {
                **response.get("media_buy", {}),
                "reporting_delivery": reporting,
            }
            response["experimental_features"] = list(
                dict.fromkeys(
                    [*response.get("experimental_features", ()), "media_buy.reporting_delivery"]
                )
            )
            if any(
                reporting.get(k)
                for k in ("ledger_notification", "status_notification", "readiness_notification")
            ):
                response["webhook_signing"] = declaration["webhook_signing"]
                response["identity"] = declaration["identity"]
        return response

    @_admitted
    async def sync_accounts(
        self,
        params: SyncAccountsRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any]:
        try:
            self.production._assert_components()
            return await self.production.configuration_task.execute(
                self.production, _request(params), context
            )
        except LedgerConflictError as error:
            code, message = error.code, str(error)
        except ReportingReceiptError as error:
            code, message = error.code, str(error)
        except (ReliableReportingConfigurationError, ReliableReportingUnavailableError):
            raise
        except Exception:
            code, message = (
                "REPORTING_CONFIGURATION_UNAVAILABLE",
                "reporting configuration is unavailable",
            )
        raise _task_error("sync_accounts", code, message)

    @_admitted
    async def sync_reporting_status(
        self,
        params: SyncReportingStatusRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any] | NotImplementedResponse:
        if not self._feed_consumer_status_enabled:
            return self._not_supported("sync_reporting_status")
        try:
            request = _request(params)
            caller = await self._authorize(request, context)
            result = await ConsumerStatusIngest(self.production.store, enabled=True).handle(
                request,
                account_id=caller.account_id,
                consumer_id=caller.consumer_id,
            )
            if await self._authorize(request, context) != caller:
                raise ReportingReceiptError("UNAUTHORIZED")
            return result
        except (ReportingReceiptError, LedgerConflictError) as error:
            code, message = error.code, str(error)
        except (ReliableReportingConfigurationError, ReliableReportingUnavailableError):
            raise
        except Exception:
            code, message = "REPORTING_STATUS_UNAVAILABLE", "reporting status is unavailable"
        raise _task_error("sync_reporting_status", code, message)

    @_admitted
    async def get_media_buy_delivery(
        self,
        params: GetMediaBuyDeliveryRequest | dict[str, Any],
        context: ToolContext | None = None,
    ) -> dict[str, Any] | NotImplementedResponse:
        request = _request(params)
        if "reporting_revision_id" not in request:
            return cast(
                dict[str, Any] | NotImplementedResponse,
                await self._delegate("get_media_buy_delivery", params, context),
            )
        try:
            from adcp.validation.schema_loader import get_named_validator

            validator = get_named_validator("media-buy/get-media-buy-delivery-request.json")
            if validator is None or next(validator.iter_errors(request), None) is not None:
                raise LedgerConflictError("INVALID_REQUEST", "exact revision request is invalid")
            caller = await self._authorize(request, context)
            revision_id = request["reporting_revision_id"]
            status_request = {
                "account": request["account"],
                "view": "revision",
                "reporting_revision_id": revision_id,
            }
            # This is the actual private mounted status path, including exact
            # ownership/visibility and complete captured reconciliation checks.
            projected = await self.get_reporting_status(status_request, context)
            if not isinstance(projected, dict) or "revision" not in projected:
                raise LedgerConflictError("LOOKUP_UNAVAILABLE", "no such revision is available")
            store = self.production.store
            revision = await store.get_revision(
                account_id=caller.account_id, reporting_revision_id=revision_id
            )
            if revision is None or not revision.readable:
                raise LedgerConflictError("LOOKUP_UNAVAILABLE", "no such revision is available")
            pagination = request.get("pagination") or {}
            limit = pagination.get("max_results", 50)
            if type(limit) is not int or not 1 <= limit <= 100:
                raise LedgerConflictError("INVALID_REQUEST", "exact revision page size is invalid")
            binding_hash = hashlib.sha256(
                canonical_json_utf8_v1(
                    {
                        "account": caller.account_id,
                        "consumer": caller.consumer_id,
                        "revision": revision_id,
                        "digest": revision.revision_content_sha256,
                        "count": revision.row_count,
                        "ownership": 2,
                        "limit": limit,
                    }
                )
            ).hexdigest()
            token = pagination.get("cursor")
            offset = 0
            if token is not None:
                if not isinstance(token, str) or not token.startswith("rpr2.") or len(token) > 2048:
                    raise LedgerConflictError(
                        "INVALID_CURSOR", "exact revision cursor is unavailable"
                    )
                cursor = decode_cursor(token[5:])
                if cursor.get("ownership") != 2:
                    raise ADCPTaskError(
                        "get_media_buy_delivery",
                        [
                            {
                                "code": "INVALID_CHECKPOINT",
                                "message": "Restart revision pagination without the legacy cursor.",
                                "recovery": "correctable",
                            }
                        ],
                    )
                if (
                    set(cursor) != {"v", "h", "p", "ownership"}
                    or cursor["v"] != 2
                    or cursor["ownership"] != 2
                    or cursor["h"] != binding_hash
                    or type(cursor["p"]) is not int
                    or not 0 <= cursor["p"] <= revision.row_count
                ):
                    raise LedgerConflictError(
                        "INVALID_CURSOR", "exact revision cursor is unavailable"
                    )
                offset = cursor["p"]
            page = await store.read_revision_rows(
                account_id=caller.account_id,
                reporting_revision_id=revision_id,
                limit=limit,
                cursor=(
                    encode_cursor({"ownership": 2, "revision": revision_id, "offset": offset})
                    if offset
                    else None
                ),
            )
            if page.total_count != revision.row_count or page.reporting_revision_id != revision_id:
                raise ReportingNotificationError("reporting_revision_content_unavailable")
            if await self._authorize(request, context) != caller:
                raise ReportingReceiptError("UNAUTHORIZED")
            after = await store.get_revision(
                account_id=caller.account_id, reporting_revision_id=revision_id
            )
            if after != revision:
                raise LedgerConflictError("LOOKUP_UNAVAILABLE", "no such revision is available")
            wire = projected["revision"]
            totals = (
                [r.to_wire() for r in revision.managed_control_totals]
                if revision.managed_control_totals is not None
                else [{"name": n, "value": v} for n, v in revision.control_totals]
            )
            position: dict[str, Any] = {"total_count": page.total_count, "has_more": page.has_more}
            if page.has_more:
                position["cursor"] = "rpr2." + encode_cursor(
                    {"v": 2, "ownership": 2, "h": binding_hash, "p": offset + len(page.rows)}
                )
            result = {
                "status": "completed",
                "reporting_period": {
                    "start": wire["period"]["start"],
                    "end": wire["period"]["end"],
                },
                "media_buy_deliveries": [],
                "reporting_revision": wire,
                "reporting_revision_binding": {
                    "reporting_revision_id": revision_id,
                    "row_count": revision.row_count,
                    "control_totals": totals,
                    "content_sha256": revision.revision_content_sha256,
                },
                "reporting_rows": list(page.rows),
                "pagination": position,
            }
            if "ext" in projected:
                result["ext"] = projected["ext"]
            return result
        except ADCPTaskError:
            raise
        except (ReportingReceiptError, LedgerConflictError) as error:
            code, message = error.code, str(error)
        except (ReliableReportingConfigurationError, ReliableReportingUnavailableError):
            raise
        except Exception:
            code, message = "REPORTING_CONTENT_UNAVAILABLE", "reporting content is unavailable"
        raise _task_error("get_media_buy_delivery", code, message)

Mount the same instance on MCP/A2A; the support owns its lifecycle.

Ancestors

Class variables

var advertised_tools

Methods

def advertised_tools_for_instance(self) ‑> set[str]
Expand source code
def advertised_tools_for_instance(self) -> set[str]:
    names = {
        "get_adcp_capabilities",
        "get_reporting_status",
        "get_media_buy_delivery",
        "sync_accounts",
    }
    if any(o.reconciled for o in self.production.offerings):
        names.add("sync_reporting_receipts")
    if self._feed_consumer_status_enabled:
        names.add("sync_reporting_status")
    return names | set(self._application_tools & _APPLICATION_TASKS)
def bind_application(self, application: ADCPHandler[Any]) ‑> None
Expand source code
def bind_application(self, application: ADCPHandler[Any]) -> None:
    """Delegate ordinary tasks before mounting this exact SDK handler.

    Reporting/configuration collisions are refused. Supply the typed B2
    configuration task separately; aggregate delivery and base capabilities
    may coexist and are dispatched explicitly by their protected methods.
    """
    if application is self or not isinstance(application, ADCPHandler):
        raise ValueError("application must be a separate ADCPHandler")
    if self._application is application:
        return
    if (
        self._application is not None
        or self.production._mounts
        or self.production._task is not None
    ):
        raise ValueError("application delegation must be fixed before mounting")
    names = {t["name"] for t in get_tools_for_handler(application, _include_schemas=False)}
    if names & (_REPORTING_TASKS - {"get_adcp_capabilities", "get_media_buy_delivery"}):
        raise ValueError("application reporting handlers conflict with production ownership")
    version = getattr(application, "get_adcp_version", None)
    if callable(version) and version() != self.get_adcp_version():
        raise ValueError("application and production protocol versions differ")
    self._application, self._application_tools = application, frozenset(names)

Delegate ordinary tasks before mounting this exact SDK handler.

Reporting/configuration collisions are refused. Supply the typed B2 configuration task separately; aggregate delivery and base capabilities may coexist and are dispatched explicitly by their protected methods.

Inherited members

class ReportingProductionMethod (capability: ReportingWriterCapability, method: Mapping[str, Any])
Expand source code
@dataclass(frozen=True)
class ReportingProductionMethod:
    """The actual provider's complete public method for one writer capability.

    This includes the provider, destination modes, access/producer identity and
    reader requirements. A method advertised by an offering must equal this
    declaration. The copied bytes cannot change through an adopter's mapping.
    Runtime writer, resolver, verifier and authorization checks remain required.
    """

    capability: ReportingWriterCapability
    method: Mapping[str, Any] = field(repr=False, compare=False)
    _wire: bytes = field(init=False, repr=False)

    def __post_init__(self) -> None:
        from adcp.validation.schema_loader import get_named_validator

        raw = json.loads(canonical_json_utf8_v1(dict(self.method)))
        validator = get_named_validator("core/reporting-delivery-offering.json")
        if (
            validator is None
            or next(
                validator.evolve(schema=validator.schema["properties"]["method"]).iter_errors(raw),
                None,
            )
            is not None
            or raw.get("orchestration") != "producer_managed"
            or (raw.get("pattern"), raw.get("transport"), raw.get("format"))
            != (self.capability.method, self.capability.transport, self.capability.format)
        ):
            raise ValueError("production method must match the provider's exact writer capability")
        object.__setattr__(self, "_wire", canonical_json_utf8_v1(raw))

    def wire(self) -> dict[str, Any]:
        return dict(json.loads(self._wire))

The actual provider's complete public method for one writer capability.

This includes the provider, destination modes, access/producer identity and reader requirements. A method advertised by an offering must equal this declaration. The copied bytes cannot change through an adopter's mapping. Runtime writer, resolver, verifier and authorization checks remain required.

Instance variables

var capability : ReportingWriterCapability
var method : Mapping[str, typing.Any]

Methods

def wire(self) ‑> dict[str, typing.Any]
Expand source code
def wire(self) -> dict[str, Any]:
    return dict(json.loads(self._wire))
class ReportingProductionOffering (offering: ReportingDeliveryOffering,
producer: ReportingProducer,
verification_key: ReportingVerificationKey,
source_offering_id: str)
Expand source code
@dataclass(frozen=True)
class ReportingProductionOffering:
    """One public offering and the exact running components that implement it.

    The wire model is copied into immutable bytes so mutating an adopter's
    Pydantic model after construction cannot change the advertised contract.
    Capability availability is still checked against live components.
    """

    offering: ReportingDeliveryOffering = field(repr=False, compare=False)
    producer: ReportingProducer = field(repr=False, compare=False)
    verification_key: ReportingVerificationKey
    source_offering_id: str
    _wire: bytes = field(init=False, repr=False)
    _producer_key: str = field(init=False, repr=False)
    _source_identity: int = field(init=False, repr=False)
    _reader_identity: int = field(init=False, repr=False)

    def __post_init__(self) -> None:
        checked = ReportingDeliveryOffering.model_validate(
            self.offering.model_dump(mode="json", exclude_none=True)
        )
        raw = checked.model_dump(mode="json", exclude_none=True)
        from adcp.validation.schema_loader import get_named_validator

        validator = get_named_validator("core/reporting-delivery-offering.json")
        if validator is None or next(validator.iter_errors(raw), None) is not None:
            raise ValueError("offering must satisfy the complete public contract")
        profile, definition, key = (
            raw["reporting_profile"],
            self.verification_key.definition,
            self.verification_key,
        )
        method = raw.get("method")
        if (
            method is None
            or method.get("orchestration") != "producer_managed"
            or (method["pattern"], method["transport"], method.get("format"))
            != (key.capability.method, key.capability.transport, key.capability.format)
            or (raw["report_definition_id"], profile["id"])
            != (key.report_definition_id, key.reporting_profile)
            or (raw["report_definition_uri"], raw["report_definition_sha256"].lower())
            != (definition.report_definition_uri, definition.report_definition_sha256)
            or (
                profile["version"],
                profile["schema_uri"],
                profile["schema_sha256"].lower(),
                profile["schema_dialect"],
                profile["schema_ref_policy"],
            )
            != (
                definition.schema_version,
                definition.schema_uri,
                definition.schema_sha256,
                definition.schema_dialect,
                definition.schema_ref_policy,
            )
            or not raw["supported_finality"]
            or len(set(raw["supported_finality"])) != len(raw["supported_finality"])
        ):
            raise ValueError("offering must match its installed producer and exact verifier")
        if raw["reconciliation_mode"] == "consumer_receipt":
            if (
                raw["supported_finality"] != ["official"]
                or key.capability.verification_profile != "canonical_digest"
                or (
                    profile.get("canonicalization_id"),
                    profile.get("canonicalization_uri"),
                    str(profile.get("canonicalization_sha256", "")).lower(),
                )
                != (
                    key.canonicalization.canonicalization_id,
                    key.canonicalization.canonicalization_uri,
                    key.canonicalization.canonicalization_sha256,
                )
            ):
                raise ValueError("reconciled offerings require official canonical receipt evidence")
        elif raw["feed_purpose"] == "billing":
            raise ValueError("billing offerings require consumer receipts")
        object.__setattr__(self, "_wire", canonical_json_utf8_v1(raw))
        object.__setattr__(self, "_source_identity", id(self.producer._source))
        object.__setattr__(self, "_reader_identity", id(self.producer._object_reader))
        object.__setattr__(
            self,
            "_producer_key",
            hashlib.sha256(
                canonical_json_utf8_v1(
                    {
                        "offering": raw,
                        "source_offering_id": self.source_offering_id,
                        "publication_namespace": self.producer._offerings.publication_namespace,
                        "source_scope": dict(self.producer._offerings.source_scope),
                    }
                )
            ).hexdigest(),
        )
        self.check_source()

    @property
    def offering_id(self) -> str:
        return str(self.wire()["offering_id"])

    @property
    def reconciled(self) -> bool:
        return bool(self.wire()["reconciliation_mode"] == "consumer_receipt")

    def wire(self) -> dict[str, Any]:
        return dict(json.loads(self._wire))

    def check_source(self, *, effective: bool = False) -> ReportingSourceCapabilitiesV1:
        source = self.producer._source
        if (
            id(source) != self._source_identity
            or not isinstance(source, ReportingProductionSource)
            or id(self.producer._object_reader) != self._reader_identity
            or not isinstance(self.producer._object_reader, ReportingSourceStagedObjectReader)
        ):
            raise failure("BINDING_MISMATCH")
        capabilities = ReportingSourceCapabilitiesV1.model_validate(
            source.capabilities.model_dump(mode="json", exclude_none=True)
        )
        if effective and capabilities.scope != "effective_account":
            raise failure("BINDING_MISMATCH")
        raw = self.wire()
        source_offering = capabilities.offering(self.source_offering_id)
        contract, key = source_offering.contract, self.verification_key
        source_values = contract.model_dump(mode="json")
        expected = {
            "report_definition_id": key.report_definition_id,
            "reporting_profile": key.reporting_profile,
        }
        expected.update(
            {
                name: getattr(key.definition, name)
                for name in (
                    "report_definition_uri",
                    "report_definition_sha256",
                    "schema_version",
                    "schema_uri",
                    "schema_sha256",
                    "schema_dialect",
                    "schema_ref_policy",
                )
            }
        )
        finality = raw["supported_finality"]
        configured = self.producer._offerings
        if (
            any(source_values.get(name) != value for name, value in expected.items())
            or source_offering.publication_namespace != configured.publication_namespace
            or dict(configured.source_scope) != capabilities.source_scope
            or not source_offering.provider_execution.supports_cancellation
            or (
                "official" in finality
                and (
                    not isinstance(source_offering, AuthoritativeOfferingV1)
                    or configured.official_offering_id != self.source_offering_id
                    or (
                        self.reconciled
                        and source_offering.correction_policy != "immutable_correction"
                    )
                )
            )
            or (
                "snapshot" in finality
                and (
                    not isinstance(source_offering, ProvisionalSnapshotOfferingV1)
                    or configured.snapshot_offering_id != self.source_offering_id
                )
            )
        ):
            raise failure("UNSUPPORTED_VERIFICATION")
        schedule = raw["schedule"]
        duration = iso_duration_milliseconds_v1(schedule["period_duration"])
        if not (
            iso_duration_milliseconds_v1(source_offering.windowing.minimum_window)
            <= duration
            <= iso_duration_milliseconds_v1(source_offering.windowing.maximum_window)
        ) or iso_duration_milliseconds_v1(schedule["delivery_sla"]) < iso_duration_milliseconds_v1(
            source_offering.worst_case_availability_lag
        ):
            raise failure("UNSUPPORTED_VERIFICATION")
        if isinstance(
            source_offering, ProvisionalSnapshotOfferingV1
        ) and duration < iso_duration_milliseconds_v1(source_offering.fastest_safe_cadence):
            raise failure("UNSUPPORTED_VERIFICATION")
        if (
            schedule["alignment"] == "source_timezone"
            and schedule.get("period_timezone") != source_offering.source_timezone
        ):
            raise failure("UNSUPPORTED_VERIFICATION")
        return capabilities

    def source_binding(
        self, configuration: ReportingConfiguration
    ) -> ReportingProductionSourceBinding:
        require_account_work(configuration.account_id)
        try:
            capabilities = self.check_source(effective=True)
            source = self.producer._source
            assert isinstance(source, ReportingProductionSource)
            binding = source.configuration_binding(configuration)
            if binding is None:
                # Discovery can probe several generations before selecting a
                # dispatch. Refuse this candidate without revoking a turn that
                # may select another one; the producer's dispatch/publication
                # checks own the account-wide stop after a live denial.
                raise _SourceAuthorizationRevokedError()
            if type(binding) is not ReportingProductionSourceBinding:
                raise failure("BINDING_MISMATCH")
            binding.check(configuration, capabilities, self.source_offering_id)
            return binding
        except _SourceAuthorizationRevokedError:
            raise
        except Exception:
            raise failure("BINDING_MISMATCH") from None

    def configuration_schedule(self, configuration: ReportingConfiguration) -> dict[str, Any]:
        """Project a proven legacy clock into the public schedule vocabulary.

        Existing captured clocks and their closed storage decoder stay intact.
        New production admission requires the normative public phase; an old
        explicit anchor is compatible only when it produces the same periods.
        Calendar-month/year clocks unsupported by that decoder are refused.
        """
        offered = self.wire()["schedule"]
        schedule = configuration.schedule
        alignment = offered["alignment"]
        expected_alignment = (
            alignment if alignment in {"utc", "account_timezone"} else "custom_timezone"
        )
        zone, duration, anchor = _schedule_clock(schedule, configuration.account_timezone)
        if (
            schedule.alignment != expected_alignment
            or schedule.period_duration != offered["period_duration"]
            or schedule.delivery_sla != offered["delivery_sla"]
            or (alignment != "billing_cycle" and (anchor - datetime(1970, 1, 1)) % duration)
            or (alignment == "billing_cycle" and schedule.period_anchor is None)
        ):
            raise failure("BINDING_MISMATCH")
        result: dict[str, Any] = {
            "period_duration": schedule.period_duration,
            "delivery_sla": schedule.delivery_sla,
            "alignment": alignment,
        }
        if alignment in {"source_timezone", "billing_cycle"}:
            result["period_timezone"] = zone.key
        if alignment == "billing_cycle":
            assert schedule.period_anchor is not None
            result["period_anchor"] = schedule.period_anchor.isoformat()
        if (
            offered.get("period_timezone_policy") == "fixed"
            or offered.get("period_anchor_policy") == "fixed"
        ) and result.get("period_timezone") != offered.get("period_timezone"):
            raise failure("BINDING_MISMATCH")
        if offered.get(
            "period_anchor_policy"
        ) == "fixed" and schedule.period_anchor != aware_timestamp(offered["period_anchor"]):
            raise failure("BINDING_MISMATCH")
        if (
            alignment == "source_timezone"
            and zone.key
            != self.check_source(effective=True).offering(self.source_offering_id).source_timezone
        ):
            raise failure("BINDING_MISMATCH")
        return result

    def check_configuration(
        self,
        configuration: ReportingConfiguration,
        binding: ReportingDestinationBinding,
        *,
        retention_days: int,
    ) -> None:
        capabilities = self.check_source(effective=True)
        self.source_binding(configuration)
        self.configuration_schedule(configuration)
        raw, key = self.wire(), self.verification_key
        schedule, offered = configuration.schedule, raw["schedule"]
        if (
            binding.generation_key != configuration.generation_key
            or (
                configuration.report_definition_id,
                configuration.reporting_profile,
                configuration.feed_purpose,
                configuration.required_finality,
            )
            != (
                key.report_definition_id,
                key.reporting_profile,
                raw["feed_purpose"],
                raw["supported_finality"][0],
            )
            or not _same_definition(key, configuration.definition)
            or configuration.authoritative_party != "seller"
            or binding.reconciliation_mode != raw["reconciliation_mode"]
            or binding.feed_purpose != raw["feed_purpose"]
            or (binding.method, binding.transport, binding.format, binding.verification_profile)
            != (
                key.capability.method,
                key.capability.transport,
                key.capability.format,
                key.capability.verification_profile,
            )
            or binding.resource_retention_days != retention_days
            or binding.reader_compatibility != tuple(raw["method"].get("reader_compatibility", ()))
            or schedule.period_duration != offered["period_duration"]
            or iso_duration_to_timedelta(schedule.delivery_sla)
            != iso_duration_to_timedelta(offered["delivery_sla"])
            or iso_duration_milliseconds_v1(schedule.delivery_sla)
            < iso_duration_milliseconds_v1(
                capabilities.offering(self.source_offering_id).worst_case_availability_lag
            )
            or (
                offered.get("period_anchor_policy") == "fixed"
                and (
                    schedule.period_anchor is None
                    or schedule.period_anchor != aware_timestamp(str(offered.get("period_anchor")))
                )
            )
            or (
                offered.get("period_timezone_policy") == "fixed"
                and schedule.period_timezone != offered.get("period_timezone")
            )
        ):
            raise failure("BINDING_MISMATCH")

One public offering and the exact running components that implement it.

The wire model is copied into immutable bytes so mutating an adopter's Pydantic model after construction cannot change the advertised contract. Capability availability is still checked against live components.

Instance variables

var offering : ReportingDeliveryOffering
prop offering_id : str
Expand source code
@property
def offering_id(self) -> str:
    return str(self.wire()["offering_id"])
var producer : ReportingProducer
prop reconciled : bool
Expand source code
@property
def reconciled(self) -> bool:
    return bool(self.wire()["reconciliation_mode"] == "consumer_receipt")
var source_offering_id : str
var verification_key : ReportingVerificationKey

Methods

def check_configuration(self,
configuration: ReportingConfiguration,
binding: ReportingDestinationBinding,
*,
retention_days: int) ‑> None
Expand source code
def check_configuration(
    self,
    configuration: ReportingConfiguration,
    binding: ReportingDestinationBinding,
    *,
    retention_days: int,
) -> None:
    capabilities = self.check_source(effective=True)
    self.source_binding(configuration)
    self.configuration_schedule(configuration)
    raw, key = self.wire(), self.verification_key
    schedule, offered = configuration.schedule, raw["schedule"]
    if (
        binding.generation_key != configuration.generation_key
        or (
            configuration.report_definition_id,
            configuration.reporting_profile,
            configuration.feed_purpose,
            configuration.required_finality,
        )
        != (
            key.report_definition_id,
            key.reporting_profile,
            raw["feed_purpose"],
            raw["supported_finality"][0],
        )
        or not _same_definition(key, configuration.definition)
        or configuration.authoritative_party != "seller"
        or binding.reconciliation_mode != raw["reconciliation_mode"]
        or binding.feed_purpose != raw["feed_purpose"]
        or (binding.method, binding.transport, binding.format, binding.verification_profile)
        != (
            key.capability.method,
            key.capability.transport,
            key.capability.format,
            key.capability.verification_profile,
        )
        or binding.resource_retention_days != retention_days
        or binding.reader_compatibility != tuple(raw["method"].get("reader_compatibility", ()))
        or schedule.period_duration != offered["period_duration"]
        or iso_duration_to_timedelta(schedule.delivery_sla)
        != iso_duration_to_timedelta(offered["delivery_sla"])
        or iso_duration_milliseconds_v1(schedule.delivery_sla)
        < iso_duration_milliseconds_v1(
            capabilities.offering(self.source_offering_id).worst_case_availability_lag
        )
        or (
            offered.get("period_anchor_policy") == "fixed"
            and (
                schedule.period_anchor is None
                or schedule.period_anchor != aware_timestamp(str(offered.get("period_anchor")))
            )
        )
        or (
            offered.get("period_timezone_policy") == "fixed"
            and schedule.period_timezone != offered.get("period_timezone")
        )
    ):
        raise failure("BINDING_MISMATCH")
def check_source(self, *, effective: bool = False) ‑> ReportingSourceCapabilitiesV1
Expand source code
def check_source(self, *, effective: bool = False) -> ReportingSourceCapabilitiesV1:
    source = self.producer._source
    if (
        id(source) != self._source_identity
        or not isinstance(source, ReportingProductionSource)
        or id(self.producer._object_reader) != self._reader_identity
        or not isinstance(self.producer._object_reader, ReportingSourceStagedObjectReader)
    ):
        raise failure("BINDING_MISMATCH")
    capabilities = ReportingSourceCapabilitiesV1.model_validate(
        source.capabilities.model_dump(mode="json", exclude_none=True)
    )
    if effective and capabilities.scope != "effective_account":
        raise failure("BINDING_MISMATCH")
    raw = self.wire()
    source_offering = capabilities.offering(self.source_offering_id)
    contract, key = source_offering.contract, self.verification_key
    source_values = contract.model_dump(mode="json")
    expected = {
        "report_definition_id": key.report_definition_id,
        "reporting_profile": key.reporting_profile,
    }
    expected.update(
        {
            name: getattr(key.definition, name)
            for name in (
                "report_definition_uri",
                "report_definition_sha256",
                "schema_version",
                "schema_uri",
                "schema_sha256",
                "schema_dialect",
                "schema_ref_policy",
            )
        }
    )
    finality = raw["supported_finality"]
    configured = self.producer._offerings
    if (
        any(source_values.get(name) != value for name, value in expected.items())
        or source_offering.publication_namespace != configured.publication_namespace
        or dict(configured.source_scope) != capabilities.source_scope
        or not source_offering.provider_execution.supports_cancellation
        or (
            "official" in finality
            and (
                not isinstance(source_offering, AuthoritativeOfferingV1)
                or configured.official_offering_id != self.source_offering_id
                or (
                    self.reconciled
                    and source_offering.correction_policy != "immutable_correction"
                )
            )
        )
        or (
            "snapshot" in finality
            and (
                not isinstance(source_offering, ProvisionalSnapshotOfferingV1)
                or configured.snapshot_offering_id != self.source_offering_id
            )
        )
    ):
        raise failure("UNSUPPORTED_VERIFICATION")
    schedule = raw["schedule"]
    duration = iso_duration_milliseconds_v1(schedule["period_duration"])
    if not (
        iso_duration_milliseconds_v1(source_offering.windowing.minimum_window)
        <= duration
        <= iso_duration_milliseconds_v1(source_offering.windowing.maximum_window)
    ) or iso_duration_milliseconds_v1(schedule["delivery_sla"]) < iso_duration_milliseconds_v1(
        source_offering.worst_case_availability_lag
    ):
        raise failure("UNSUPPORTED_VERIFICATION")
    if isinstance(
        source_offering, ProvisionalSnapshotOfferingV1
    ) and duration < iso_duration_milliseconds_v1(source_offering.fastest_safe_cadence):
        raise failure("UNSUPPORTED_VERIFICATION")
    if (
        schedule["alignment"] == "source_timezone"
        and schedule.get("period_timezone") != source_offering.source_timezone
    ):
        raise failure("UNSUPPORTED_VERIFICATION")
    return capabilities
def configuration_schedule(self, configuration: ReportingConfiguration) ‑> dict[str, typing.Any]
Expand source code
def configuration_schedule(self, configuration: ReportingConfiguration) -> dict[str, Any]:
    """Project a proven legacy clock into the public schedule vocabulary.

    Existing captured clocks and their closed storage decoder stay intact.
    New production admission requires the normative public phase; an old
    explicit anchor is compatible only when it produces the same periods.
    Calendar-month/year clocks unsupported by that decoder are refused.
    """
    offered = self.wire()["schedule"]
    schedule = configuration.schedule
    alignment = offered["alignment"]
    expected_alignment = (
        alignment if alignment in {"utc", "account_timezone"} else "custom_timezone"
    )
    zone, duration, anchor = _schedule_clock(schedule, configuration.account_timezone)
    if (
        schedule.alignment != expected_alignment
        or schedule.period_duration != offered["period_duration"]
        or schedule.delivery_sla != offered["delivery_sla"]
        or (alignment != "billing_cycle" and (anchor - datetime(1970, 1, 1)) % duration)
        or (alignment == "billing_cycle" and schedule.period_anchor is None)
    ):
        raise failure("BINDING_MISMATCH")
    result: dict[str, Any] = {
        "period_duration": schedule.period_duration,
        "delivery_sla": schedule.delivery_sla,
        "alignment": alignment,
    }
    if alignment in {"source_timezone", "billing_cycle"}:
        result["period_timezone"] = zone.key
    if alignment == "billing_cycle":
        assert schedule.period_anchor is not None
        result["period_anchor"] = schedule.period_anchor.isoformat()
    if (
        offered.get("period_timezone_policy") == "fixed"
        or offered.get("period_anchor_policy") == "fixed"
    ) and result.get("period_timezone") != offered.get("period_timezone"):
        raise failure("BINDING_MISMATCH")
    if offered.get(
        "period_anchor_policy"
    ) == "fixed" and schedule.period_anchor != aware_timestamp(offered["period_anchor"]):
        raise failure("BINDING_MISMATCH")
    if (
        alignment == "source_timezone"
        and zone.key
        != self.check_source(effective=True).offering(self.source_offering_id).source_timezone
    ):
        raise failure("BINDING_MISMATCH")
    return result

Project a proven legacy clock into the public schedule vocabulary.

Existing captured clocks and their closed storage decoder stay intact. New production admission requires the normative public phase; an old explicit anchor is compatible only when it produces the same periods. Calendar-month/year clocks unsupported by that decoder are refused.

def source_binding(self, configuration: ReportingConfiguration) ‑> ReportingProductionSourceBinding
Expand source code
def source_binding(
    self, configuration: ReportingConfiguration
) -> ReportingProductionSourceBinding:
    require_account_work(configuration.account_id)
    try:
        capabilities = self.check_source(effective=True)
        source = self.producer._source
        assert isinstance(source, ReportingProductionSource)
        binding = source.configuration_binding(configuration)
        if binding is None:
            # Discovery can probe several generations before selecting a
            # dispatch. Refuse this candidate without revoking a turn that
            # may select another one; the producer's dispatch/publication
            # checks own the account-wide stop after a live denial.
            raise _SourceAuthorizationRevokedError()
        if type(binding) is not ReportingProductionSourceBinding:
            raise failure("BINDING_MISMATCH")
        binding.check(configuration, capabilities, self.source_offering_id)
        return binding
    except _SourceAuthorizationRevokedError:
        raise
    except Exception:
        raise failure("BINDING_MISMATCH") from None
def wire(self) ‑> dict[str, typing.Any]
Expand source code
def wire(self) -> dict[str, Any]:
    return dict(json.loads(self._wire))
class ReportingProductionSigning (resolver: ReportingSigningResolver,
algorithms: tuple[str, ...],
*,
brand_json_url: str)
Expand source code
@dataclass(frozen=True)
class ReportingProductionSigning:
    """A declared RFC 9421 algorithm contract checked on every resolved key.

    Key material remains in the trusted resolver. The public declaration and
    actual sender share this immutable contract; rotation cannot switch to an
    unadvertised algorithm or a legacy authentication path.
    """

    resolver: ReportingSigningResolver = field(repr=False)
    algorithms: tuple[str, ...]
    brand_json_url: str = field(kw_only=True)

    def __post_init__(self) -> None:
        object.__setattr__(self, "algorithms", tuple(self.algorithms))
        if (
            not self.algorithms
            or len(set(self.algorithms)) != len(self.algorithms)
            or not set(self.algorithms) <= ALLOWED_ALGS
            or not callable(getattr(self.resolver, "resolve", None))
        ):
            raise ReportingNotificationError("notification_signing_unready")
        try:
            if type(self.brand_json_url) is not str:
                raise ValueError
            identity = httpx.URL(self.brand_json_url)
            valid = (
                identity.scheme == "https"
                and bool(identity.host)
                and not identity.userinfo
                and not identity.query
                and not identity.fragment
                and identity.port in (None, 443)
            )
        except (TypeError, ValueError, httpx.InvalidURL):
            valid = False
        if not valid:
            raise ReportingNotificationError("notification_signing_unready")

    async def resolve(
        self, *, account_id: str, principal_id: str, signing_scope_id: str
    ) -> ReportingSigningMaterial:
        material = await self.resolver.resolve(
            account_id=account_id, principal_id=principal_id, signing_scope_id=signing_scope_id
        )
        if type(
            material
        ) is not ReportingSigningMaterial or material.advertised_algorithms != frozenset(
            self.algorithms
        ):
            raise ReportingNotificationError("notification_signing_unready")
        ReportingSigningMaterial.__post_init__(material)
        return material

    def wire(self) -> dict[str, Any]:
        return {
            "supported": True,
            "profile": "adcp/webhook-signing/v1",
            "algorithms": list(self.algorithms),
            "legacy_hmac_fallback": False,
            "delivery_retry_horizon_seconds": RETRY_HORIZON_SECONDS,
        }

A declared RFC 9421 algorithm contract checked on every resolved key.

Key material remains in the trusted resolver. The public declaration and actual sender share this immutable contract; rotation cannot switch to an unadvertised algorithm or a legacy authentication path.

Instance variables

var algorithms : tuple[str, ...]
var brand_json_url : str
var resolver : ReportingSigningResolver

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:
    material = await self.resolver.resolve(
        account_id=account_id, principal_id=principal_id, signing_scope_id=signing_scope_id
    )
    if type(
        material
    ) is not ReportingSigningMaterial or material.advertised_algorithms != frozenset(
        self.algorithms
    ):
        raise ReportingNotificationError("notification_signing_unready")
    ReportingSigningMaterial.__post_init__(material)
    return material
def wire(self) ‑> dict[str, typing.Any]
Expand source code
def wire(self) -> dict[str, Any]:
    return {
        "supported": True,
        "profile": "adcp/webhook-signing/v1",
        "algorithms": list(self.algorithms),
        "legacy_hmac_fallback": False,
        "delivery_retry_horizon_seconds": RETRY_HORIZON_SECONDS,
    }
class ReportingProductionSource (*args, **kwargs)
Expand source code
@runtime_checkable
class ReportingProductionSource(ReportingSourceExecutor, Protocol):
    """A source with an authenticated generation mapping available before I/O.

    Discovery may precede any account binding. Admission and every source turn
    require the applicable binding, obtained from the source's trusted account
    configuration. The SDK calls ``configuration_binding`` immediately before
    dispatch and again under the account lock before sealing or publishing the
    result, including replay after a restart. Returning ``None`` discards that
    result without a success checkpoint and stops this account's remaining
    work for the turn. Restoring the binding permits work on the next turn.

    An in-flight fetch is allowed to finish: revocation takes effect at the next
    dispatch or publish. Adopters own any caching or latency inside their
    callback; the SDK does not retain an authorization grant.
    """

    def configuration_binding(
        self, configuration: ReportingConfiguration
    ) -> ReportingProductionSourceBinding | None:
        """Return the current trusted binding, or ``None`` to deny source work.

        Adopters own callback caching and latency. In-flight fetches finish;
        revocation takes effect at the next dispatch or publish.
        """
        raise NotImplementedError

A source with an authenticated generation mapping available before I/O.

Discovery may precede any account binding. Admission and every source turn require the applicable binding, obtained from the source's trusted account configuration. The SDK calls configuration_binding immediately before dispatch and again under the account lock before sealing or publishing the result, including replay after a restart. Returning None discards that result without a success checkpoint and stops this account's remaining work for the turn. Restoring the binding permits work on the next turn.

An in-flight fetch is allowed to finish: revocation takes effect at the next dispatch or publish. Adopters own any caching or latency inside their callback; the SDK does not retain an authorization grant.

Ancestors

Methods

def configuration_binding(self, configuration: ReportingConfiguration) ‑> ReportingProductionSourceBinding | None
Expand source code
def configuration_binding(
    self, configuration: ReportingConfiguration
) -> ReportingProductionSourceBinding | None:
    """Return the current trusted binding, or ``None`` to deny source work.

    Adopters own callback caching and latency. In-flight fetches finish;
    revocation takes effect at the next dispatch or publish.
    """
    raise NotImplementedError

Return the current trusted binding, or None to deny source work.

Adopters own callback caching and latency. In-flight fetches finish; revocation takes effect at the next dispatch or publish.

Inherited members

class ReportingProductionSourceBinding (generation_key: ReportingConfigurationGenerationKey,
capabilities_sha256: str,
media_buy_products: tuple[tuple[str, str], ...],
*,
configuration_sha256: str,
service_context: ReportingProductionSourceContext | None = None)
Expand source code
@dataclass(frozen=True)
class ReportingProductionSourceBinding:
    """Trusted account-to-source mapping, fixed for a configuration generation.

    ``media_buy_products`` comes from the authenticated source/account mapping,
    never from buyer JSON or a report-definition identifier. The SDK persists
    it with admission, checks the effective capability digest on every source
    turn, and refuses a changed mapping. Reauthorization can withdraw a binding;
    it cannot silently change historical scope. No credentials belong here.
    """

    generation_key: ReportingConfigurationGenerationKey
    capabilities_sha256: str
    media_buy_products: tuple[tuple[str, str], ...]
    configuration_sha256: str = field(kw_only=True)
    service_context: ReportingProductionSourceContext | None = field(
        default=None, kw_only=True, repr=False
    )

    def __post_init__(self) -> None:
        from adcp.reporting.evidence import reporting_identifier, sha256_value

        if type(self.generation_key) is not ReportingConfigurationGenerationKey:
            raise ValueError("source binding requires an exact configuration generation")
        sha256_value(self.capabilities_sha256)
        sha256_value(self.configuration_sha256)
        if (
            self.service_context is not None
            and type(self.service_context) is not ReportingProductionSourceContext
        ):
            raise ValueError("source binding requires an exact service source context")
        pairs = tuple(tuple(pair) for pair in self.media_buy_products)
        if any(len(pair) != 2 for pair in pairs):
            raise ValueError("source binding requires media-buy/product pairs")
        for media_buy_id, product_id in pairs:
            reporting_identifier(media_buy_id, maximum=255)
            reporting_identifier(product_id, maximum=255)
        if len({pair[0] for pair in pairs}) != len(pairs):
            raise ValueError("source binding media buys must be unique")
        object.__setattr__(self, "media_buy_products", tuple(sorted(pairs)))

    @classmethod
    def for_configuration(
        cls,
        configuration: ReportingConfiguration,
        *,
        capabilities_sha256: str,
        media_buy_products: tuple[tuple[str, str], ...],
    ) -> ReportingProductionSourceBinding:
        """Freeze the exact generation semantics and explicitly resolved products.

        Lifecycle changes retain the same semantic generation, matching the
        ledger's immutable configuration contract. Captured projection inputs
        independently retain each historical activation/deactivation boundary.
        """
        return cls(
            configuration.generation_key,
            capabilities_sha256,
            media_buy_products,
            configuration_sha256=hashlib.sha256(
                canonical_json_utf8_v1(_config_payload(configuration))
            ).hexdigest(),
        )

    def document(self) -> dict[str, Any]:
        key = self.generation_key
        document = {
            "account_id": key.account_id,
            "consumer_id": key.consumer_id,
            "delivery_config_id": key.delivery_config_id,
            "delivery_config_version": key.delivery_config_version,
            "capabilities_sha256": self.capabilities_sha256,
            "configuration_sha256": self.configuration_sha256,
            "media_buy_products": [list(pair) for pair in self.media_buy_products],
        }
        if self.service_context is not None:
            document["service_context"] = self.service_context.document()
            document["service_context_sha256"] = hashlib.sha256(
                self.service_context._wire
            ).hexdigest()
        return document

    def check(
        self,
        configuration: ReportingConfiguration,
        capabilities: ReportingSourceCapabilitiesV1,
        offering_id: str,
    ) -> None:
        offering = capabilities.offering(offering_id)
        if (
            self.generation_key != configuration.generation_key
            or self.configuration_sha256
            != hashlib.sha256(canonical_json_utf8_v1(_config_payload(configuration))).hexdigest()
            or capabilities.scope != "effective_account"
            or self.capabilities_sha256 != capabilities.capabilities_sha256
            or {pair[0] for pair in self.media_buy_products} != set(configuration.media_buy_ids)
            or (self.media_buy_products and "media_buy" not in offering.constituent_kinds)
            or any(pair[1] not in offering.product_ids for pair in self.media_buy_products)
        ):
            raise failure("BINDING_MISMATCH")

    def constituents(self) -> tuple[ReportingConstituent, ...]:
        return tuple(
            MediaBuyConstituentV1(
                constituent_id=media_buy_id, media_buy_id=media_buy_id, product_id=product_id
            )
            for media_buy_id, product_id in self.media_buy_products
        )

Trusted account-to-source mapping, fixed for a configuration generation.

media_buy_products comes from the authenticated source/account mapping, never from buyer JSON or a report-definition identifier. The SDK persists it with admission, checks the effective capability digest on every source turn, and refuses a changed mapping. Reauthorization can withdraw a binding; it cannot silently change historical scope. No credentials belong here.

Static methods

def for_configuration(configuration: ReportingConfiguration,
*,
capabilities_sha256: str,
media_buy_products: tuple[tuple[str, str], ...]) ‑> ReportingProductionSourceBinding

Freeze the exact generation semantics and explicitly resolved products.

Lifecycle changes retain the same semantic generation, matching the ledger's immutable configuration contract. Captured projection inputs independently retain each historical activation/deactivation boundary.

Instance variables

var capabilities_sha256 : str
var configuration_sha256 : str
var generation_key : ReportingConfigurationGenerationKey
var media_buy_products : tuple[tuple[str, str], ...]
var service_context : ReportingProductionSourceContext | None

Methods

def check(self,
configuration: ReportingConfiguration,
capabilities: ReportingSourceCapabilitiesV1,
offering_id: str) ‑> None
Expand source code
def check(
    self,
    configuration: ReportingConfiguration,
    capabilities: ReportingSourceCapabilitiesV1,
    offering_id: str,
) -> None:
    offering = capabilities.offering(offering_id)
    if (
        self.generation_key != configuration.generation_key
        or self.configuration_sha256
        != hashlib.sha256(canonical_json_utf8_v1(_config_payload(configuration))).hexdigest()
        or capabilities.scope != "effective_account"
        or self.capabilities_sha256 != capabilities.capabilities_sha256
        or {pair[0] for pair in self.media_buy_products} != set(configuration.media_buy_ids)
        or (self.media_buy_products and "media_buy" not in offering.constituent_kinds)
        or any(pair[1] not in offering.product_ids for pair in self.media_buy_products)
    ):
        raise failure("BINDING_MISMATCH")
def constituents(self) ‑> tuple[ProductConstituentV1 | PackageItemConstituentV1 | MediaBuyConstituentV1, ...]
Expand source code
def constituents(self) -> tuple[ReportingConstituent, ...]:
    return tuple(
        MediaBuyConstituentV1(
            constituent_id=media_buy_id, media_buy_id=media_buy_id, product_id=product_id
        )
        for media_buy_id, product_id in self.media_buy_products
    )
def document(self) ‑> dict[str, typing.Any]
Expand source code
def document(self) -> dict[str, Any]:
    key = self.generation_key
    document = {
        "account_id": key.account_id,
        "consumer_id": key.consumer_id,
        "delivery_config_id": key.delivery_config_id,
        "delivery_config_version": key.delivery_config_version,
        "capabilities_sha256": self.capabilities_sha256,
        "configuration_sha256": self.configuration_sha256,
        "media_buy_products": [list(pair) for pair in self.media_buy_products],
    }
    if self.service_context is not None:
        document["service_context"] = self.service_context.document()
        document["service_context_sha256"] = hashlib.sha256(
            self.service_context._wire
        ).hexdigest()
    return document
class ReportingProductionSourceRegistry (*, account_context: ReportingContextResolver)
Expand source code
class ReportingProductionSourceRegistry:
    """Bind stable adapter names to the actual fixed production producers.

    ``account_context`` runs only at authenticated configuration admission,
    outside the storage transaction. Recovery validates the selected persisted
    generation against these profiles and the live production source binding;
    it never calls the resolver or rebuilds an in-memory account catalog.
    """

    def __init__(self, *, account_context: ReportingContextResolver) -> None:
        self._resolve = account_context
        self._entries: dict[str, _Registration] = {}
        self._frozen = False

    def register(self, adapter: str, producer: ReportingProducer) -> None:
        if self._frozen:
            raise ValueError("production source registry is frozen")
        if not re.fullmatch(r"[A-Za-z0-9_.:-]{1,128}", adapter):
            raise ValueError("adapter must be a stable identifier")
        if adapter in self._entries or any(v.producer is producer for v in self._entries.values()):
            raise ValueError("production source profile must have one unambiguous registration")
        self._entries[adapter] = _Registration(
            adapter, producer, canonical_json_utf8_v1(_profile(producer._offerings))
        )

    def freeze(self, offerings: tuple[ReportingProductionOffering, ...]) -> None:
        for offering in offerings:
            self._registration(offering)
        if any(
            not any(entry.producer is offering.producer for offering in offerings)
            for entry in self._entries.values()
        ):
            raise ValueError("registered production source profile has no offering")
        self._frozen = True

    def _registration(self, offering: ReportingProductionOffering) -> _Registration:
        matches = [v for v in self._entries.values() if v.producer is offering.producer]
        if len(matches) != 1:
            raise ValueError("production source profile is not registered")
        entry = matches[0]
        if entry.profile != canonical_json_utf8_v1(_profile(entry.producer._offerings)):
            raise ValueError("production source profile changed")
        return entry

    def _expected(
        self,
        configuration: ReportingConfiguration,
        offering: ReportingProductionOffering,
        capabilities: ReportingSourceCapabilitiesV1,
    ) -> dict[str, Any]:
        entry = self._registration(offering)
        return {
            "version": 1,
            "adapter": entry.adapter,
            "account_id": configuration.account_id,
            "account_timezone": configuration.account_timezone,
            "offering_id": offering.offering_id,
            "source_offering_id": offering.source_offering_id,
            "capabilities_sha256": capabilities.capabilities_sha256,
            **json.loads(entry.profile),
        }

    async def resolve(
        self,
        configuration: ReportingConfiguration,
        offering: ReportingProductionOffering,
        capabilities: ReportingSourceCapabilitiesV1,
    ) -> ReportingProductionSourceContext:
        from adcp.reporting.service import ReportingAccountContext, _thaw

        context = self._resolve(configuration)
        if inspect.isawaitable(context):
            context = await context
        expected = self._expected(configuration, offering, capabilities)
        if (
            type(context) is not ReportingAccountContext
            or context.account_id != configuration.account_id
            or context.account_timezone != configuration.account_timezone
            or context.adapter != expected["adapter"]
            or canonical_json_utf8_v1(_profile(context.producer_offerings()))
            != self._registration(offering).profile
            or (
                context.capability_offering
                and canonical_json_utf8_v1(_thaw(context.capability_offering))
                != canonical_json_utf8_v1(offering.wire())
            )
        ):
            raise ValueError("resolved account context does not match its fixed production profile")
        return ReportingProductionSourceContext(canonical_json_utf8_v1(expected))

    def recover(
        self,
        configuration: ReportingConfiguration,
        offering: ReportingProductionOffering,
        capabilities: ReportingSourceCapabilitiesV1,
        document: dict[str, Any] | None,
    ) -> ReportingProductionSourceContext:
        if document is None:
            raise ValueError("legacy service source context requires explicit migration")
        expected = canonical_json_utf8_v1(self._expected(configuration, offering, capabilities))
        if expected != canonical_json_utf8_v1(document):
            raise ValueError(
                "persisted account context does not match its fixed production profile"
            )
        return ReportingProductionSourceContext(expected)

Bind stable adapter names to the actual fixed production producers.

account_context runs only at authenticated configuration admission, outside the storage transaction. Recovery validates the selected persisted generation against these profiles and the live production source binding; it never calls the resolver or rebuilds an in-memory account catalog.

Methods

def freeze(self,
offerings: tuple[ReportingProductionOffering, ...]) ‑> None
Expand source code
def freeze(self, offerings: tuple[ReportingProductionOffering, ...]) -> None:
    for offering in offerings:
        self._registration(offering)
    if any(
        not any(entry.producer is offering.producer for offering in offerings)
        for entry in self._entries.values()
    ):
        raise ValueError("registered production source profile has no offering")
    self._frozen = True
def recover(self,
configuration: ReportingConfiguration,
offering: ReportingProductionOffering,
capabilities: ReportingSourceCapabilitiesV1,
document: dict[str, Any] | None) ‑> ReportingProductionSourceContext
Expand source code
def recover(
    self,
    configuration: ReportingConfiguration,
    offering: ReportingProductionOffering,
    capabilities: ReportingSourceCapabilitiesV1,
    document: dict[str, Any] | None,
) -> ReportingProductionSourceContext:
    if document is None:
        raise ValueError("legacy service source context requires explicit migration")
    expected = canonical_json_utf8_v1(self._expected(configuration, offering, capabilities))
    if expected != canonical_json_utf8_v1(document):
        raise ValueError(
            "persisted account context does not match its fixed production profile"
        )
    return ReportingProductionSourceContext(expected)
def register(self, adapter: str, producer: ReportingProducer) ‑> None
Expand source code
def register(self, adapter: str, producer: ReportingProducer) -> None:
    if self._frozen:
        raise ValueError("production source registry is frozen")
    if not re.fullmatch(r"[A-Za-z0-9_.:-]{1,128}", adapter):
        raise ValueError("adapter must be a stable identifier")
    if adapter in self._entries or any(v.producer is producer for v in self._entries.values()):
        raise ValueError("production source profile must have one unambiguous registration")
    self._entries[adapter] = _Registration(
        adapter, producer, canonical_json_utf8_v1(_profile(producer._offerings))
    )
async def resolve(self,
configuration: ReportingConfiguration,
offering: ReportingProductionOffering,
capabilities: ReportingSourceCapabilitiesV1) ‑> ReportingProductionSourceContext
Expand source code
async def resolve(
    self,
    configuration: ReportingConfiguration,
    offering: ReportingProductionOffering,
    capabilities: ReportingSourceCapabilitiesV1,
) -> ReportingProductionSourceContext:
    from adcp.reporting.service import ReportingAccountContext, _thaw

    context = self._resolve(configuration)
    if inspect.isawaitable(context):
        context = await context
    expected = self._expected(configuration, offering, capabilities)
    if (
        type(context) is not ReportingAccountContext
        or context.account_id != configuration.account_id
        or context.account_timezone != configuration.account_timezone
        or context.adapter != expected["adapter"]
        or canonical_json_utf8_v1(_profile(context.producer_offerings()))
        != self._registration(offering).profile
        or (
            context.capability_offering
            and canonical_json_utf8_v1(_thaw(context.capability_offering))
            != canonical_json_utf8_v1(offering.wire())
        )
    ):
        raise ValueError("resolved account context does not match its fixed production profile")
    return ReportingProductionSourceContext(canonical_json_utf8_v1(expected))
class ReportingProductionSupport (materializer: ReportingMaterializerService,
projection: InMemoryReportingStatusProjection | PgReportingStatusProjection,
*,
offerings: tuple[ReportingProductionOffering, ...],
configuration_task: ReportingProductionConfigurationTask,
resolve_account: ReceiptAccountResolver,
buyer_agents: BuyerAgentRegistry | None = None,
automated_recovery_window: timedelta = datetime.timedelta(seconds=21600),
status_retention_days: int = 400,
notification_workers: tuple[ReportingNotificationWorker, ...] = (),
poll_seconds: float = 0.25,
adcp_version: str | None = None,
source_registry: ReportingProductionSourceRegistry | None = None)
Expand source code
class ReportingProductionSupport:
    """Compose, mount, start, then activate accounts after draining old workers.

    The owned worker performs indexed producer/materializer/projector turns.
    Stopping it prevents fresh admission while registered discovery remains stable. An
    optional HTTP notification worker is independent of polling readiness;
    enabled logical enqueue is always part of the materializer transaction.

    Migrate with ``store.create_schema()`` before construction/start. For DDL,
    drain and close this support, migrate, and construct fresh support. Schema
    proofs do not detect arbitrary serving-time DDL or search-path changes.
    """

    def __init__(
        self,
        materializer: ReportingMaterializerService,
        projection: InMemoryReportingStatusProjection | PgReportingStatusProjection,
        *,
        offerings: tuple[ReportingProductionOffering, ...],
        configuration_task: ReportingProductionConfigurationTask,
        resolve_account: ReceiptAccountResolver,
        buyer_agents: BuyerAgentRegistry | None = None,
        automated_recovery_window: timedelta = timedelta(hours=6),
        status_retention_days: int = 400,
        notification_workers: tuple[ReportingNotificationWorker, ...] = (),
        poll_seconds: float = 0.25,
        adcp_version: str | None = None,
        source_registry: ReportingProductionSourceRegistry | None = None,
    ) -> None:
        from adcp.reporting.outbox.worker import ReportingNotificationWorker
        from adcp.reporting.production.handler import ReportingProductionHandler
        from adcp.reporting.production.pg import PgReportingProductionStore
        from adcp.reporting.projection.pg import PgReportingStatusProjection

        if (
            type(materializer) is not ReportingMaterializerService
            or type(materializer.store)
            not in {InMemoryReportingProductionStore, PgReportingProductionStore}
            or type(projection)
            not in {InMemoryReportingStatusProjection, PgReportingStatusProjection}
            or projection.ledger is not materializer.store
            or type(configuration_task) is not ReportingProductionConfigurationTask
            or type(offerings) is not tuple
            or not offerings
            or any(type(o) is not ReportingProductionOffering for o in offerings)
            or len({o.offering_id for o in offerings}) != len(offerings)
            or type(status_retention_days) is not int
            or not 1 <= status_retention_days <= 36500
            or automated_recovery_window <= timedelta(0)
            or not 0.01 <= poll_seconds <= 60
            or type(notification_workers) is not tuple
            or any(type(w) is not ReportingNotificationWorker for w in notification_workers)
        ):
            raise ValueError(
                "production support requires exact SDK components and bounded promises"
            )
        assert isinstance(
            materializer.store, (InMemoryReportingProductionStore, PgReportingProductionStore)
        )
        self.store = materializer.store
        previous = self.store._production_support
        if previous is not None and previous() is not None:
            raise ReportingNotificationError("reporting_production_owner_conflict")
        self.materializer, self.projection = materializer, projection
        self.offerings, self.configuration_task = offerings, configuration_task
        if source_registry is not None:
            if type(source_registry) is not ReportingProductionSourceRegistry:
                raise ValueError("production requires an exact source registry")
            source_registry.freeze(offerings)
        self.source_registry = source_registry
        self._service_lifecycle: _ServiceLifecycle | None = None
        self.automated_recovery_window, self.status_retention_days = (
            automated_recovery_window,
            status_retention_days,
        )
        self.notification_workers, self.poll_seconds = notification_workers, poll_seconds
        from adcp.reporting.production.notifications import worker_identity

        self._notification_identity = tuple(worker_identity(w) for w in notification_workers)
        self.handler = ReportingProductionHandler(
            self,
            resolve_account=resolve_account,
            buyer_agents=buyer_agents,
            adcp_version=adcp_version,
        )
        self._mounts: list[_MountedProduction] = []
        self._task: asyncio.Task[None] | None = None
        self._notification_task: asyncio.Task[None] | None = None
        self._stop = asyncio.Event()
        self._started = asyncio.Event()
        self._notification_started = asyncio.Event()
        self._closed = False
        self._failed = False
        self._producer_turn: ContextVar[ReportingProducer | None] = ContextVar(
            "reporting_production_source_turn", default=None
        )
        self._schema_task: asyncio.Task[bool] | None = None
        self._schema_epoch = 0
        self._schema_positive: tuple[int, int] | None = None
        self._protocol_version = self.handler.get_adcp_version()
        self._components = self._component_ids()
        self._methods = self._handler_methods()
        self.store._production_support = weakref.ref(self)
        self._assert_components(running=False, mounted=False)
        from adcp.reporting.production._declaration import declared_reporting_delivery
        from adcp.reporting.production.notifications import signing_capabilities, signing_identity

        writer = self.materializer.writer
        assert isinstance(writer, ReportingProductionDestination)
        declaration: dict[str, Any] = {
            "reporting_delivery": declared_reporting_delivery(
                offerings=[o.wire() for o in self.offerings],
                destination=writer,
                durable=not isinstance(self.store, InMemoryReportingProductionStore),
                escalation=self.projection.escalation or ReportingDeliveryEscalation(),
                consumer_status_enabled=self.projection.consumer_status_enabled,
                automated_recovery_window=self.automated_recovery_window,
                status_retention_days=self.status_retention_days,
                notifications=bool(self.notification_workers),
            )
        }
        if self.notification_workers and declaration["reporting_delivery"]:
            declaration["webhook_signing"] = signing_capabilities(self)
            declaration["identity"] = signing_identity(self)
        self._declaration = json.dumps(declaration)

    def _declared_capabilities(self) -> dict[str, Any]:
        """Detached promises captured when the concrete graph was registered."""
        # Worker state does not change registered promises. A changed protocol
        # pin invalidates the declaration itself, even with a warm schema proof.
        if (
            self.handler._adcp_version != self._protocol_version
            or self.handler.get_adcp_version() != self._protocol_version
        ):
            return {"reporting_delivery": {}}
        return dict(json.loads(self._declaration))

    @property
    def keys(self) -> tuple[ReportingVerificationKey, ...]:
        return tuple(dict.fromkeys(o.verification_key for o in self.offerings))

    @property
    def notifications_enabled(self) -> bool:
        return bool(self.projection.policy["notifications_enabled"])

    def _component_ids(self) -> tuple[int, ...]:
        return tuple(
            id(v)
            for v in (
                self.store,
                getattr(self.store, "_pool", None),
                self.materializer,
                self.materializer.io,
                self.materializer.io.registry,
                self.materializer.io.resolver,
                self.materializer.writer,
                self.projection,
                self.projection.outbox,
                self.configuration_task,
                self.source_registry,
                *self.offerings,
                *self.notification_workers,
            )
        )

    def _handler_methods(self) -> tuple[object, ...]:
        return tuple(
            getattr(getattr(self.handler, name), "__func__", None)
            for name in (
                "get_reporting_status",
                "get_media_buy_delivery",
                "sync_reporting_receipts",
                "sync_reporting_status",
                "sync_accounts",
                "get_adcp_capabilities",
                "get_adcp_version",
            )
        )

    def _mounted(self, *, receipts: bool = False) -> bool:
        required = {
            "get_adcp_capabilities",
            "get_reporting_status",
            "get_media_buy_delivery",
            "sync_accounts",
        }
        if receipts:
            required.add("sync_reporting_receipts")
        if self.projection.consumer_status_enabled:
            required.add("sync_reporting_status")
        live = [proof for proof in self._mounts if proof.mount() is not None]
        return bool(live) and all(
            required <= proof.entries.keys()
            and all(proof.current().get(name) == proof.entries[name] for name in required)
            for proof in live
        )

    def _assert_components(self, *, running: bool = True, mounted: bool = True) -> None:
        from adcp.reporting.production.notifications import check_workers

        check_workers(self)
        writer = self.materializer.writer
        if (
            self._closed
            or self._failed
            or (running and (self._task is None or self._task.done()))
            or (
                running
                and self.notification_workers
                and (self._notification_task is None or self._notification_task.done())
            )
            or self._components != self._component_ids()
            or self._methods != self._handler_methods()
            or type(self.handler._adcp_version) is not str
            or self.handler._adcp_version != self._protocol_version
            or (
                (running or mounted)
                and not isinstance(self.store, InMemoryReportingProductionStore)
                and self.store._pool.closed is not False
            )
            or self.materializer.store is not self.store
            or self.projection.ledger is not self.store
            or not self._projection_outbox_linked()
            or self.handler.receipt_store is not self.store
            or self.handler.reporting_feed_store is not self.store
            or self.handler._feed_consumer_status_enabled != self.projection.consumer_status_enabled
            or self.store._projection_read_policy != self.projection.policy
            or type(self.materializer.io) is not ReportingDestinationIO
            or isinstance(writer, ReferenceReportingDestinationWriter)
            or isinstance(self.materializer.io.resolver, ReferenceReportingResolver)
            or not isinstance(writer, ReportingProductionDestination)
            or writer is not self.materializer.io.resolver
            or writer.production_eligible is not True
            or type(writer.resource_retention_days) is not int
            or not 1 <= writer.resource_retention_days <= 36500
            or type(writer.authorization_revocation_seconds) is not int
            or not 0 <= writer.authorization_revocation_seconds <= 86400
            or (mounted and not self._mounted())
        ):
            raise ReportingNotificationError("reporting_production_component_unready")

    def _projection_outbox_linked(self) -> bool:
        # The projection queue participates even when optional HTTP delivery
        # workers are absent. Its connection ownership cannot ride a cached
        # catalog proof belonging to a different pool or in-memory ledger.
        if isinstance(self.store, InMemoryReportingProductionStore):
            from adcp.reporting.projection.memory import InMemoryReportingProjectionOutbox

            return (
                type(self.projection.outbox) is InMemoryReportingProjectionOutbox
                and self.projection.outbox._store is self.store
            )
        from adcp.reporting.projection.notifications import PgReportingProjectionOutbox

        return (
            type(self.projection.outbox) is PgReportingProjectionOutbox
            and self.projection.outbox._pool is self.store._pool
        )

    def _offering_ready(self, offering: ReportingProductionOffering) -> bool:
        try:
            if (
                offering.producer.store is not self.store
                or type(offering.producer._max_periods_per_turn) is not int
                or not 1 <= offering.producer._max_periods_per_turn <= 64
                or offering.producer.escalation != self.projection.escalation
                or not _method_ready(self.materializer.writer, offering)
                or (offering.reconciled and not self._mounted(receipts=True))
            ):
                return False
            verifier = self.materializer.io.registry.require(offering.verification_key)
            profile = offering.wire()["reporting_profile"]
            contract = json.loads(verifier.canonicalization_bytes)
            definition = json.loads(verifier.definition_bytes)
            if (
                offering.producer._revision_verifier is not verifier
                or profile["primary_keys"] != contract["primary_keys"]
                or profile["grain"] != definition.get("grain")
            ):
                return False
            offering.check_source()
            if self.source_registry is not None:
                self.source_registry._registration(offering)
            return True
        except Exception:
            return False

    def _admission_policy(self) -> dict[str, Any]:
        self._assert_components()
        # This fixes the installed contract, not its current authorization.
        # Each reservation separately proves its own offering; an unrelated
        # provider outage must not rewrite this epoch or veto healthy offerings.
        ready = self.offerings
        return {
            "version": 2,
            "verification_keys": sorted({verification_key_id(o.verification_key) for o in ready}),
            "offerings": sorted(hashlib.sha256(o._wire).hexdigest() for o in ready),
            "producers": sorted(o._producer_key for o in ready),
            "notifications_enabled": self.notifications_enabled,
            "projection": self.projection.policy,
        }

    def _configuration_offering(
        self,
        configuration: ReportingConfiguration,
        binding: ReportingDestinationBinding,
        *,
        offering_id: str | None = None,
        producer_key: str | None = None,
    ) -> ReportingProductionOffering:
        self._assert_components()
        writer = self.materializer.writer
        assert isinstance(writer, ReportingProductionDestination)
        matches = []
        for offering in self.offerings:
            if (offering_id is not None and offering.offering_id != offering_id) or (
                producer_key is not None and offering._producer_key != producer_key
            ):
                continue
            if not self._offering_ready(offering):
                continue
            try:
                offering.check_configuration(
                    configuration, binding, retention_days=writer.resource_retention_days
                )
                self._destination_binding(binding, offering)
            # A failed proof never qualifies this offering.
            except Exception:  # nosec B112
                continue
            matches.append(offering)
        if (
            len(matches) != 1
            or configuration.automated_recovery_window != self.automated_recovery_window
            or configuration.status_retention_days < self.status_retention_days
        ):
            raise failure("BINDING_MISMATCH")
        return matches[0]

    def _destination_binding(
        self, binding: ReportingDestinationBinding, offering: ReportingProductionOffering
    ) -> ReportingProductionDestinationBinding:
        try:
            writer = self.materializer.writer
            assert isinstance(writer, ReportingProductionDestination)
            resolved = writer.configuration_binding(binding)
            if (
                type(resolved) is not ReportingProductionDestinationBinding
                or resolved.binding != binding
                or resolved.method.capability != offering.verification_key.capability
                or resolved.method.wire() != offering.wire()["method"]
            ):
                raise failure("BINDING_MISMATCH")
            return resolved
        except Exception:
            raise failure("BINDING_MISMATCH") from None

    def _producer_keys(self) -> tuple[str, ...]:
        self._assert_components()
        producer = self._producer_turn.get()
        return tuple(
            o._producer_key
            for o in self.offerings
            if o.producer is producer and self._offering_ready(o)
        )

    def validate_configuration(self, value: ReportingConfigurationAdmission) -> None:
        value.check(self)

    def _source_binding(
        self,
        configuration: ReportingConfiguration,
        producer_key: str,
        *,
        document: dict[str, Any] | None = None,
        service_context: ReportingProductionSourceContext | None = None,
    ) -> ReportingProductionSourceBinding:
        self._assert_components()
        matches = [
            offering
            for offering in self.offerings
            if offering._producer_key == producer_key and self._offering_ready(offering)
        ]
        if len(matches) != 1:
            raise failure("BINDING_MISMATCH")
        offering = matches[0]
        binding = offering.source_binding(configuration)
        if self.source_registry is not None:
            context = self.source_registry.recover(
                configuration,
                offering,
                offering.check_source(effective=True),
                (
                    service_context.document()
                    if service_context is not None
                    else (document or {}).get("service_context")
                ),
            )
            binding = replace(binding, service_context=context)
        elif service_context is not None or (document or {}).get("service_context") is not None:
            raise failure("BINDING_MISMATCH")
        return binding

    async def _admit_configuration(self, value: ReportingConfigurationAdmission) -> None:
        context = None
        if self.source_registry is not None:
            offering = self._configuration_offering(
                value.configuration, value.binding, offering_id=value.offering_id
            )
            context = await self.source_registry.resolve(
                value.configuration, offering, offering.check_source(effective=True)
            )
        await self.store.admit_production_configuration(
            value.configuration,
            value.binding,
            offering_id=value.offering_id,
            service_context=context,
        )

    def _check_source_binding(
        self,
        configuration: ReportingConfiguration,
        producer_key: str,
        document: dict[str, Any] | None,
    ) -> ReportingProductionSourceBinding:
        try:
            binding = self._source_binding(configuration, producer_key, document=document)
        except (ValueError, TypeError):
            # Runtime incompatibility follows the existing scoped refusal path:
            # unavailable legacy/profile facts cannot stop unrelated work.
            raise failure("BINDING_MISMATCH") from None
        if binding.document() != document:
            raise failure("BINDING_MISMATCH")
        return binding

    def _check_context(
        self,
        context: MaterializerContext,
        key: ReportingVerificationKey,
        policy: dict[str, Any],
        producer_key: str | None,
        source_binding: dict[str, Any] | None,
        destination_binding: dict[str, Any] | None,
    ) -> None:
        if producer_key is None:
            raise failure("BINDING_MISMATCH")
        offering = self._configuration_offering(
            context.configuration, context.binding, producer_key=producer_key
        )
        self._check_source_binding(context.configuration, producer_key, source_binding)
        if self._destination_binding(context.binding, offering).wire() != destination_binding:
            raise failure("BINDING_MISMATCH")
        if (
            policy != self._admission_policy()
            or offering.verification_key != key
            or verification_key_id(key) not in policy["verification_keys"]
            or context.obligation.generation_key != context.configuration.generation_key
        ):
            raise failure("BINDING_MISMATCH")

    def invalidate_schema_validation(self) -> None:
        """Drain first; an older scan cannot publish a new epoch's proof."""
        self._schema_epoch += 1
        self._schema_positive = None
        self._schema_task = None

    async def _schema_ready(self) -> bool:
        from adcp.reporting.production.pg import PgReportingProductionStore
        from adcp.reporting.production.schema import validate_production_schema

        if not isinstance(self.store, PgReportingProductionStore):
            return False
        store = self.store
        pool, epoch = store._pool, self._schema_epoch
        identity = (id(pool), epoch)
        if self._schema_positive == identity:
            return True
        task = self._schema_task
        if task is None:

            async def scan() -> bool:
                try:
                    async with pool.connection() as connection:
                        await validate_production_schema(
                            connection,
                            notifications=self.notifications_enabled,
                            service_context=self.source_registry is not None,
                        )
                    if epoch == self._schema_epoch and pool is store._pool:
                        self._schema_positive = identity
                        return True
                    return False
                finally:
                    if self._schema_task is asyncio.current_task():
                        self._schema_task = None

            task = asyncio.create_task(scan())
            self._schema_task = task
            task.add_done_callback(lambda done: None if done.cancelled() else done.exception())
        if task.get_loop() is not asyncio.get_running_loop():
            return False
        return await asyncio.shield(task)

    async def start(self) -> None:
        if self._task is not None:
            self._assert_components()
            return
        self._assert_components(running=False)
        if (
            not isinstance(self.store, InMemoryReportingProductionStore)
            and not await self._schema_ready()
        ):
            raise ReportingNotificationError("reporting_production_schema_unready")
        self._assert_components(running=False)
        self._task = asyncio.create_task(self._run(), name="adcp-reporting-production")
        if self.notification_workers:
            self._notification_task = asyncio.create_task(
                self._run_notifications(), name="adcp-reporting-production-notifications"
            )
        await self._started.wait()
        if self.notification_workers:
            await self._notification_started.wait()
        self._assert_components()

    async def _run(self) -> None:
        boundary: _WorkerBoundary = "producer"
        try:
            while not self._stop.is_set():
                # Each producer leases a generation through the reviewed fair
                # indexed acquisition path; no account enumeration is required.
                with source_turn():
                    for producer in dict.fromkeys(o.producer for o in self.offerings):
                        boundary = "producer"
                        token = self._producer_turn.set(producer)
                        try:
                            await producer.run_worker()
                        finally:
                            self._producer_turn.reset(token)
                boundary = "materializer"
                await self.materializer.run_once()
                boundary = "projection"
                await self.projection.rebuild_one()
                boundary = "sweeper"
                await self.projection.sweep_one()
                self._started.set()
                try:
                    await asyncio.wait_for(self._stop.wait(), self.poll_seconds)
                except asyncio.TimeoutError:
                    pass
        except asyncio.CancelledError:
            self._stop.set()
            raise
        except Exception as error:
            already_stopping = self._stop.is_set()
            self._failed = True
            self._stop.set()
            if not already_stopping and not isinstance(error, ReportingNotificationError):
                _worker_stopped(boundary=boundary)
        finally:
            self._started.set()

    async def _run_notifications(self) -> None:
        from adcp.reporting.production.notifications import notification_turn

        try:
            while not self._stop.is_set():
                await notification_turn(self)
                self._notification_started.set()
                try:
                    await asyncio.wait_for(self._stop.wait(), self.poll_seconds)
                except asyncio.TimeoutError:
                    pass
        except asyncio.CancelledError:
            self._stop.set()
            raise
        except Exception as error:
            already_stopping = self._stop.is_set()
            self._failed = True
            self._stop.set()
            if not already_stopping and not isinstance(error, ReportingNotificationError):
                _worker_stopped(boundary="notifications")
        finally:
            self._notification_started.set()

    async def _check_notifications(self, account_id: str) -> None:
        from adcp.reporting.production.notifications import check_account_notifications

        self._assert_components()
        try:
            await check_account_notifications(self, account_id)
        except Exception:
            raise ReportingNotificationError("notification_chain_unready") from None
        self._assert_components()

    async def activate(self, *, account_id: str) -> bool:
        await self._check_notifications(account_id)
        await self.projection.activate(account_id=account_id)
        return await self.store._activate_production(account_id=account_id)

    async def aclose(self) -> None:
        self._stop.set()
        task = self._task
        if task is not None:
            # Drain short database turns instead of interrupting psycopg's
            # transaction entry/exit. Long I/O already has SDK-owned deadlines.
            await asyncio.gather(task, return_exceptions=True)
        if self._notification_task is not None:
            await asyncio.gather(self._notification_task, return_exceptions=True)
        self._notification_task = None
        self._task = None
        self._closed = True
        proof = self._schema_task
        if proof is not None:
            await asyncio.gather(proof, return_exceptions=True)
        if self.store._production_support is not None and self.store._production_support() is self:
            self.store._production_support = None

    async def reporting_delivery(self, *, extra: Mapping[str, Any] | None = None) -> dict[str, Any]:
        # Validate protected extra even when the concrete surface is unready.
        producer = self.offerings[0].producer
        payload = producer.advertised_reporting_delivery(
            consumer_status_task=self.projection.consumer_status_enabled,
            offerings=[],
            automated_recovery_window=self.automated_recovery_window,
            status_retention_days=self.status_retention_days,
            extra=extra,
        )
        try:
            if self._stop.is_set():
                return {}
            self._assert_components()
            if not await self._schema_ready():
                return {}
            # Catalog proof is cached; live component/mount ownership is not.
            # It must still hold after an awaited first scan or invalidation.
            self._assert_components()
            offerings = [o for o in self.offerings if self._offering_ready(o)]
            if not offerings:
                return {}
            writer = self.materializer.writer
            assert isinstance(writer, ReportingProductionDestination)
            payload.update(
                offerings=[o.wire() for o in offerings],
                managed_delivery=True,
                resource_retention_days=writer.resource_retention_days,
                authorization_revocation_seconds=writer.authorization_revocation_seconds,
            )
            if any(o.reconciled for o in offerings):
                payload.update(reconciled_billing=True, receipt_task="sync_reporting_receipts")
            if self.notification_workers:
                payload.update(
                    ledger_notification="reporting.ledger_changed",
                    readiness_notification="reporting.delivery_ready",
                    status_notification="reporting.status_changed",
                )
            return payload
        except Exception:
            return {}

Compose, mount, start, then activate accounts after draining old workers.

The owned worker performs indexed producer/materializer/projector turns. Stopping it prevents fresh admission while registered discovery remains stable. An optional HTTP notification worker is independent of polling readiness; enabled logical enqueue is always part of the materializer transaction.

Migrate with store.create_schema() before construction/start. For DDL, drain and close this support, migrate, and construct fresh support. Schema proofs do not detect arbitrary serving-time DDL or search-path changes.

Instance variables

prop keys : tuple[ReportingVerificationKey, ...]
Expand source code
@property
def keys(self) -> tuple[ReportingVerificationKey, ...]:
    return tuple(dict.fromkeys(o.verification_key for o in self.offerings))
prop notifications_enabled : bool
Expand source code
@property
def notifications_enabled(self) -> bool:
    return bool(self.projection.policy["notifications_enabled"])

Methods

async def aclose(self) ‑> None
Expand source code
async def aclose(self) -> None:
    self._stop.set()
    task = self._task
    if task is not None:
        # Drain short database turns instead of interrupting psycopg's
        # transaction entry/exit. Long I/O already has SDK-owned deadlines.
        await asyncio.gather(task, return_exceptions=True)
    if self._notification_task is not None:
        await asyncio.gather(self._notification_task, return_exceptions=True)
    self._notification_task = None
    self._task = None
    self._closed = True
    proof = self._schema_task
    if proof is not None:
        await asyncio.gather(proof, return_exceptions=True)
    if self.store._production_support is not None and self.store._production_support() is self:
        self.store._production_support = None
async def activate(self, *, account_id: str) ‑> bool
Expand source code
async def activate(self, *, account_id: str) -> bool:
    await self._check_notifications(account_id)
    await self.projection.activate(account_id=account_id)
    return await self.store._activate_production(account_id=account_id)
def invalidate_schema_validation(self) ‑> None
Expand source code
def invalidate_schema_validation(self) -> None:
    """Drain first; an older scan cannot publish a new epoch's proof."""
    self._schema_epoch += 1
    self._schema_positive = None
    self._schema_task = None

Drain first; an older scan cannot publish a new epoch's proof.

async def reporting_delivery(self, *, extra: Mapping[str, Any] | None = None) ‑> dict[str, typing.Any]
Expand source code
async def reporting_delivery(self, *, extra: Mapping[str, Any] | None = None) -> dict[str, Any]:
    # Validate protected extra even when the concrete surface is unready.
    producer = self.offerings[0].producer
    payload = producer.advertised_reporting_delivery(
        consumer_status_task=self.projection.consumer_status_enabled,
        offerings=[],
        automated_recovery_window=self.automated_recovery_window,
        status_retention_days=self.status_retention_days,
        extra=extra,
    )
    try:
        if self._stop.is_set():
            return {}
        self._assert_components()
        if not await self._schema_ready():
            return {}
        # Catalog proof is cached; live component/mount ownership is not.
        # It must still hold after an awaited first scan or invalidation.
        self._assert_components()
        offerings = [o for o in self.offerings if self._offering_ready(o)]
        if not offerings:
            return {}
        writer = self.materializer.writer
        assert isinstance(writer, ReportingProductionDestination)
        payload.update(
            offerings=[o.wire() for o in offerings],
            managed_delivery=True,
            resource_retention_days=writer.resource_retention_days,
            authorization_revocation_seconds=writer.authorization_revocation_seconds,
        )
        if any(o.reconciled for o in offerings):
            payload.update(reconciled_billing=True, receipt_task="sync_reporting_receipts")
        if self.notification_workers:
            payload.update(
                ledger_notification="reporting.ledger_changed",
                readiness_notification="reporting.delivery_ready",
                status_notification="reporting.status_changed",
            )
        return payload
    except Exception:
        return {}
async def start(self) ‑> None
Expand source code
async def start(self) -> None:
    if self._task is not None:
        self._assert_components()
        return
    self._assert_components(running=False)
    if (
        not isinstance(self.store, InMemoryReportingProductionStore)
        and not await self._schema_ready()
    ):
        raise ReportingNotificationError("reporting_production_schema_unready")
    self._assert_components(running=False)
    self._task = asyncio.create_task(self._run(), name="adcp-reporting-production")
    if self.notification_workers:
        self._notification_task = asyncio.create_task(
            self._run_notifications(), name="adcp-reporting-production-notifications"
        )
    await self._started.wait()
    if self.notification_workers:
        await self._notification_started.wait()
    self._assert_components()
def validate_configuration(self,
value: ReportingConfigurationAdmission) ‑> None
Expand source code
def validate_configuration(self, value: ReportingConfigurationAdmission) -> None:
    value.check(self)