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_outboxThe 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
- InMemoryReportingProjectionStore
- InMemoryReportingFeedStore
- InMemoryReportingReceiptStore
- InMemoryReportingMaterializerStore
- InMemoryReportingReconciliationStore
- InMemoryReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
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
- PgReportingStatusOutbox
- PgReportingOutbox
- adcp.reporting.outbox._activity_pg._PgReportingActivity
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
- PgReportingProjectionStore
- PgReportingFeedStore
- PgReportingReceiptStore
- PgReportingMaterializerStore
- PgReportingReconciliationStore
- PgReportingLedgerStore
- adcp.reporting.ledger.delivery._ReconciliationOperations
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) -
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 : ReportingDestinationBindingvar configuration : ReportingConfigurationvar 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.
handleremains the application's authenticated, caller-owned desired state account task. It resolves provider grants and opaque bindings, calls the suppliedadmitfor 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 : Accountvar 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
- ReportingDestinationWriter
- ReportingDestinationResolver
- typing.Protocol
- typing.Generic
Instance variables
-
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.configurationis 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 : ReportingDestinationBindingvar 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
- ReportingReceiptHandler
- ADCPHandler
- abc.ABC
- typing.Generic
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
ReportingReceiptHandler:accept_proposalacquire_rightsactivate_signalbuild_creative_legacybuy_productscalibrate_contentcheck_governancecomply_test_controllercontext_matchcontrol_media_buycreate_collection_listcreate_content_standardscreate_media_buycreate_property_listdecline_proposalsdelete_collection_listdelete_property_listget_account_financialsget_adcp_capabilitiesget_adcp_versionget_brand_identityget_collection_listget_content_standardsget_creative_deliveryget_creative_featuresget_media_buy_artifactsget_media_buy_deliveryget_media_buysget_plan_audit_logsget_principalget_productsget_property_listget_reporting_statusget_rightsget_signalsget_task_statusidentity_matchlist_account_changeslist_accountslist_collection_listslist_content_standardslist_creative_formats_legacylist_creativeslist_productslist_property_listslist_taskslist_transformerslog_eventpreview_creative_legacyprovide_performance_feedbackrefine_proposalsreport_plan_adjustmentreport_plan_outcomereport_usagerequest_proposalssi_get_offeringsi_initiate_sessionsi_send_messagesi_terminate_sessionsync_accountssync_agent_notification_configssync_audiencessync_catalogssync_creativessync_event_sourcessync_governancesync_planssync_principalsync_reporting_receiptssync_reporting_statusupdate_collection_listupdate_content_standardsupdate_media_buyupdate_property_listupdate_rightsvalidate_content_deliveryvalidate_inputverify_brand_claimverify_brand_claims
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 : ReportingWriterCapabilityvar 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 : ReportingDeliveryOfferingprop offering_id : str-
Expand source code
@property def offering_id(self) -> str: return str(self.wire()["offering_id"]) var producer : ReportingProducerprop reconciled : bool-
Expand source code
@property def reconciled(self) -> bool: return bool(self.wire()["reconciliation_mode"] == "consumer_receipt") var source_offering_id : strvar 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 resultProject 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 : strvar 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 NotImplementedErrorA 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_bindingimmediately before dispatch and again under the account lock before sealing or publishing the result, including replay after a restart. ReturningNonediscards 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
- ReportingSourceExecutor
- typing.Protocol
- typing.Generic
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 NotImplementedErrorReturn the current trusted binding, or
Noneto 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_productscomes 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 : strvar configuration_sha256 : strvar generation_key : ReportingConfigurationGenerationKeyvar 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_contextruns 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 = NoneDrain 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)