Module adcp.decisioning.pg.lazy
Resolve-once wrappers for caller-owned PostgreSQL decisioning stores.
Factories open infrastructure on the serving loop, then construct the original concrete stores so constructor validation remains authoritative. A successful resolution is cached; failed initialization can be retried. These wrappers never open, close, replace or reopen pools themselves.
Classes
class LazyProposalStore (factory: LazyProposalStoreFactory)-
Expand source code
class LazyProposalStore(_LazyStore[PgProposalStore]): """Deferred PgProposalStore; durable before and after resolution.""" is_durable: ClassVar[bool] = True def __init__(self, factory: LazyProposalStoreFactory) -> None: super().__init__(factory, PgProposalStore) def _validate(self, store: PgProposalStore) -> None: if store.is_durable is not True: raise ValueError("LazyProposalStore requires a durable PostgreSQL store") @classmethod def migration_sql(cls, table_name: str = DEFAULT_TABLE_NAME) -> dict[str, str]: """Return the concrete store's migration SQL without opening infrastructure.""" return _migration_sql(table_name) async def create_schema(self) -> None: await (await self.resolve()).create_schema() async def put_draft( self, *, proposal_id: str, account_id: str, recipes: Mapping[str, Recipe], proposal_payload: Mapping[str, Any], ) -> None: await (await self.resolve()).put_draft( proposal_id=proposal_id, account_id=account_id, recipes=recipes, proposal_payload=proposal_payload, ) async def get(self, proposal_id: str, *, expected_account_id: str) -> ProposalRecord | None: return await (await self.resolve()).get( proposal_id, expected_account_id=expected_account_id ) async def commit( self, proposal_id: str, *, expires_at: datetime, proposal_payload: Mapping[str, Any], expected_account_id: str, ) -> None: await (await self.resolve()).commit( proposal_id, expires_at=expires_at, proposal_payload=proposal_payload, expected_account_id=expected_account_id, ) async def try_reserve_consumption( self, proposal_id: str, *, expected_account_id: str ) -> ProposalRecord: return await (await self.resolve()).try_reserve_consumption( proposal_id, expected_account_id=expected_account_id ) async def finalize_consumption( self, proposal_id: str, *, media_buy_id: str, expected_account_id: str ) -> None: await (await self.resolve()).finalize_consumption( proposal_id, media_buy_id=media_buy_id, expected_account_id=expected_account_id ) async def release_consumption(self, proposal_id: str, *, expected_account_id: str) -> None: await (await self.resolve()).release_consumption( proposal_id, expected_account_id=expected_account_id ) async def mark_consumed( self, proposal_id: str, *, media_buy_id: str, expected_account_id: str ) -> None: await (await self.resolve()).mark_consumed( proposal_id, media_buy_id=media_buy_id, expected_account_id=expected_account_id ) async def discard(self, proposal_id: str, *, expected_account_id: str) -> None: await (await self.resolve()).discard(proposal_id, expected_account_id=expected_account_id) async def get_by_media_buy_id( self, media_buy_id: str, *, expected_account_id: str ) -> ProposalRecord | None: return await (await self.resolve()).get_by_media_buy_id( media_buy_id, expected_account_id=expected_account_id )Deferred PgProposalStore; durable before and after resolution.
Ancestors
- adcp.decisioning.pg.lazy._LazyStore
- typing.Generic
Class variables
var is_durable : ClassVar[bool]
Static methods
def migration_sql(table_name: str = 'adcp_proposal_drafts') ‑> dict[str, str]-
Return the concrete store's migration SQL without opening infrastructure.
Methods
async def commit(self,
proposal_id: str,
*,
expires_at: datetime,
proposal_payload: Mapping[str, Any],
expected_account_id: str) ‑> None-
Expand source code
async def commit( self, proposal_id: str, *, expires_at: datetime, proposal_payload: Mapping[str, Any], expected_account_id: str, ) -> None: await (await self.resolve()).commit( proposal_id, expires_at=expires_at, proposal_payload=proposal_payload, expected_account_id=expected_account_id, ) async def create_schema(self) ‑> None-
Expand source code
async def create_schema(self) -> None: await (await self.resolve()).create_schema() async def discard(self, proposal_id: str, *, expected_account_id: str) ‑> None-
Expand source code
async def discard(self, proposal_id: str, *, expected_account_id: str) -> None: await (await self.resolve()).discard(proposal_id, expected_account_id=expected_account_id) async def finalize_consumption(self, proposal_id: str, *, media_buy_id: str, expected_account_id: str) ‑> None-
Expand source code
async def finalize_consumption( self, proposal_id: str, *, media_buy_id: str, expected_account_id: str ) -> None: await (await self.resolve()).finalize_consumption( proposal_id, media_buy_id=media_buy_id, expected_account_id=expected_account_id ) async def get(self, proposal_id: str, *, expected_account_id: str) ‑> ProposalRecord | None-
Expand source code
async def get(self, proposal_id: str, *, expected_account_id: str) -> ProposalRecord | None: return await (await self.resolve()).get( proposal_id, expected_account_id=expected_account_id ) async def get_by_media_buy_id(self, media_buy_id: str, *, expected_account_id: str) ‑> ProposalRecord | None-
Expand source code
async def get_by_media_buy_id( self, media_buy_id: str, *, expected_account_id: str ) -> ProposalRecord | None: return await (await self.resolve()).get_by_media_buy_id( media_buy_id, expected_account_id=expected_account_id ) async def mark_consumed(self, proposal_id: str, *, media_buy_id: str, expected_account_id: str) ‑> None-
Expand source code
async def mark_consumed( self, proposal_id: str, *, media_buy_id: str, expected_account_id: str ) -> None: await (await self.resolve()).mark_consumed( proposal_id, media_buy_id=media_buy_id, expected_account_id=expected_account_id ) async def put_draft(self,
*,
proposal_id: str,
account_id: str,
recipes: Mapping[str, Recipe],
proposal_payload: Mapping[str, Any]) ‑> None-
Expand source code
async def put_draft( self, *, proposal_id: str, account_id: str, recipes: Mapping[str, Recipe], proposal_payload: Mapping[str, Any], ) -> None: await (await self.resolve()).put_draft( proposal_id=proposal_id, account_id=account_id, recipes=recipes, proposal_payload=proposal_payload, ) async def release_consumption(self, proposal_id: str, *, expected_account_id: str) ‑> None-
Expand source code
async def release_consumption(self, proposal_id: str, *, expected_account_id: str) -> None: await (await self.resolve()).release_consumption( proposal_id, expected_account_id=expected_account_id ) async def try_reserve_consumption(self, proposal_id: str, *, expected_account_id: str) ‑> ProposalRecord-
Expand source code
async def try_reserve_consumption( self, proposal_id: str, *, expected_account_id: str ) -> ProposalRecord: return await (await self.resolve()).try_reserve_consumption( proposal_id, expected_account_id=expected_account_id )
class LazyTaskRegistry (factory: LazyTaskRegistryFactory)-
Expand source code
class LazyTaskRegistry(_LazyStore[PgTaskRegistry], _TaskLifecycleObservers): """Deferred PgTaskRegistry, including its listing and metrics contracts. A factory can return a registry, or an original (registry, outbox) pair. Build them together using one caller-owned pool. Pair and signing-scope validation are retained. Register observers before first use; they receive the first submitted transition too. No arbitrary delegate attributes are forwarded. For a server advertising SDK task-webhook signing, await resolve() before constructing the server: its synchronous boot validator needs the actual outbox/sender/horizon configuration. Polling-only servers can resolve on the first task operation. """ is_durable: ClassVar[bool] = True def __init__(self, factory: LazyTaskRegistryFactory) -> None: self._init_lifecycle_observers() async def resolve_registry() -> PgTaskRegistry: result = factory() stores = await result if isinstance(result, Awaitable) else result if isinstance(stores, tuple): registry, outbox = stores if not isinstance(registry, PgTaskRegistry) or not isinstance( outbox, PgTaskWebhookOutbox ): raise TypeError("Task factory must return a concrete registry/outbox pair") if registry.task_webhook_outbox is not outbox or registry._pool is not outbox._pool: raise ValueError( "Task registry/outbox pair must share one pool and registration" ) return registry return stores super().__init__(resolve_registry, PgTaskRegistry) def _validate(self, store: PgTaskRegistry) -> None: if store.is_durable is not True: raise ValueError("LazyTaskRegistry requires a durable PostgreSQL registry") if not callable(store.list): raise TypeError("LazyTaskRegistry requires PostgreSQL listing support") outbox = store.task_webhook_outbox if outbox is not None and outbox._pool is not store._pool: raise ValueError("Task registry and outbox must share one pool") store.add_lifecycle_observer(self._relay_transition) def _relay_transition( self, event: TaskTransition, *, task_id: str, account_id: str, task_type: str, created_at: float, updated_at: float, ) -> None: self._notify_lifecycle_observers( event, { "task_id": task_id, "account_id": account_id, "task_type": task_type, "created_at": created_at, "updated_at": updated_at, }, ) @property def task_webhook_outbox(self) -> PgTaskWebhookOutbox | None: """Original coupled outbox, available after resolution; no eager opening.""" store = self.resolved return store.task_webhook_outbox if store is not None else None @property def atomic_task_webhook_outbox(self) -> bool: store = self.resolved return store.atomic_task_webhook_outbox if store is not None else False async def create_schema(self) -> None: await (await self.resolve()).create_schema() async def issue( self, *, account_id: str, task_type: str, request_context: dict[str, Any] | None = None, webhook_url: str | None = None, webhook_operation_id: str | None = None, webhook_token: str | None = None, webhook_authentication: TaskWebhookAuthentication | None = None, webhook_signing_scope_id: str | None = None, **_extra: Any, ) -> str: return await (await self.resolve()).issue( account_id=account_id, task_type=task_type, request_context=request_context, webhook_url=webhook_url, webhook_operation_id=webhook_operation_id, webhook_token=webhook_token, webhook_authentication=webhook_authentication, webhook_signing_scope_id=webhook_signing_scope_id, **_extra, ) async def resolve_webhook_signing_scope(self, context: RequestContext[Any]) -> str | None: return await (await self.resolve()).resolve_webhook_signing_scope(context) async def update_progress(self, task_id: str, progress: dict[str, Any]) -> None: await (await self.resolve()).update_progress(task_id, progress) async def complete(self, task_id: str, result: dict[str, Any]) -> None: await (await self.resolve()).complete(task_id, result) async def fail(self, task_id: str, error: dict[str, Any]) -> None: await (await self.resolve()).fail(task_id, error) async def get( self, task_id: str, *, expected_account_id: str | None = None ) -> dict[str, Any] | None: return await (await self.resolve()).get(task_id, expected_account_id=expected_account_id) async def list( self, *, account_id: str, filters: dict[str, Any] | None = None, sort: dict[str, Any] | None = None, pagination: dict[str, Any] | None = None, ) -> dict[str, Any]: return await (await self.resolve()).list( account_id=account_id, filters=filters, sort=sort, pagination=pagination ) async def discard(self, task_id: str) -> None: await (await self.resolve()).discard(task_id)Deferred PgTaskRegistry, including its listing and metrics contracts.
A factory can return a registry, or an original (registry, outbox) pair. Build them together using one caller-owned pool. Pair and signing-scope validation are retained. Register observers before first use; they receive the first submitted transition too. No arbitrary delegate attributes are forwarded.
For a server advertising SDK task-webhook signing, await resolve() before constructing the server: its synchronous boot validator needs the actual outbox/sender/horizon configuration. Polling-only servers can resolve on the first task operation.
Ancestors
- adcp.decisioning.pg.lazy._LazyStore
- typing.Generic
- adcp.decisioning.task_registry._TaskLifecycleObservers
Class variables
var is_durable : ClassVar[bool]
Instance variables
prop atomic_task_webhook_outbox : bool-
Expand source code
@property def atomic_task_webhook_outbox(self) -> bool: store = self.resolved return store.atomic_task_webhook_outbox if store is not None else False prop task_webhook_outbox : PgTaskWebhookOutbox | None-
Expand source code
@property def task_webhook_outbox(self) -> PgTaskWebhookOutbox | None: """Original coupled outbox, available after resolution; no eager opening.""" store = self.resolved return store.task_webhook_outbox if store is not None else NoneOriginal coupled outbox, available after resolution; no eager opening.
Methods
async def complete(self, task_id: str, result: dict[str, Any]) ‑> None-
Expand source code
async def complete(self, task_id: str, result: dict[str, Any]) -> None: await (await self.resolve()).complete(task_id, result) async def create_schema(self) ‑> None-
Expand source code
async def create_schema(self) -> None: await (await self.resolve()).create_schema() async def discard(self, task_id: str) ‑> None-
Expand source code
async def discard(self, task_id: str) -> None: await (await self.resolve()).discard(task_id) async def fail(self, task_id: str, error: dict[str, Any]) ‑> None-
Expand source code
async def fail(self, task_id: str, error: dict[str, Any]) -> None: await (await self.resolve()).fail(task_id, error) async def get(self, task_id: str, *, expected_account_id: str | None = None) ‑> dict[str, typing.Any] | None-
Expand source code
async def get( self, task_id: str, *, expected_account_id: str | None = None ) -> dict[str, Any] | None: return await (await self.resolve()).get(task_id, expected_account_id=expected_account_id) async def issue(self,
*,
account_id: str,
task_type: str,
request_context: dict[str, Any] | None = None,
webhook_url: str | None = None,
webhook_operation_id: str | None = None,
webhook_token: str | None = None,
webhook_authentication: TaskWebhookAuthentication | None = None,
webhook_signing_scope_id: str | None = None,
**_extra: Any) ‑> str-
Expand source code
async def issue( self, *, account_id: str, task_type: str, request_context: dict[str, Any] | None = None, webhook_url: str | None = None, webhook_operation_id: str | None = None, webhook_token: str | None = None, webhook_authentication: TaskWebhookAuthentication | None = None, webhook_signing_scope_id: str | None = None, **_extra: Any, ) -> str: return await (await self.resolve()).issue( account_id=account_id, task_type=task_type, request_context=request_context, webhook_url=webhook_url, webhook_operation_id=webhook_operation_id, webhook_token=webhook_token, webhook_authentication=webhook_authentication, webhook_signing_scope_id=webhook_signing_scope_id, **_extra, ) async def list(self,
*,
account_id: str,
filters: dict[str, Any] | None = None,
sort: dict[str, Any] | None = None,
pagination: dict[str, Any] | None = None) ‑> dict[str, typing.Any]-
Expand source code
async def list( self, *, account_id: str, filters: dict[str, Any] | None = None, sort: dict[str, Any] | None = None, pagination: dict[str, Any] | None = None, ) -> dict[str, Any]: return await (await self.resolve()).list( account_id=account_id, filters=filters, sort=sort, pagination=pagination ) async def resolve_webhook_signing_scope(self, context: RequestContext[Any]) ‑> str | None-
Expand source code
async def resolve_webhook_signing_scope(self, context: RequestContext[Any]) -> str | None: return await (await self.resolve()).resolve_webhook_signing_scope(context) async def update_progress(self, task_id: str, progress: dict[str, Any]) ‑> None-
Expand source code
async def update_progress(self, task_id: str, progress: dict[str, Any]) -> None: await (await self.resolve()).update_progress(task_id, progress)
class LazyTaskWebhookOutbox (factory: LazyTaskWebhookOutboxFactory)-
Expand source code
class LazyTaskWebhookOutbox(_LazyStore[PgTaskWebhookOutbox]): """Deferred concrete outbox. Synchronous crypto helpers require resolve(). Use from_registry() to share a registry factory's original outbox. Pass the original objects, constructed inside that factory, to PgTaskRegistry; do not substitute this worker facade into the concrete constructor's pool checks. """ delivery_state_is_durable: ClassVar[bool] = True supports_atomic_task_outbox: ClassVar[bool] = True def __init__(self, factory: LazyTaskWebhookOutboxFactory) -> None: super().__init__(factory, PgTaskWebhookOutbox) def _validate(self, store: PgTaskWebhookOutbox) -> None: if ( store.delivery_state_is_durable is not True or store.supports_atomic_task_outbox is not True ): raise ValueError("LazyTaskWebhookOutbox requires a durable atomic PostgreSQL outbox") @classmethod def from_registry(cls, registry: LazyTaskRegistry) -> LazyTaskWebhookOutbox: """Create a worker facade resolving the same original registry/outbox pair.""" async def resolve_outbox() -> PgTaskWebhookOutbox: store = await registry.resolve() if store.task_webhook_outbox is None: raise ValueError("The registry factory did not configure a task webhook outbox") return store.task_webhook_outbox return cls(resolve_outbox) @property def delivery_retry_horizon_seconds(self) -> int: return self._require_resolved().delivery_retry_horizon_seconds @property def legacy_hmac_fallback(self) -> bool: return self._require_resolved().legacy_hmac_fallback async def create_schema(self) -> None: await (await self.resolve()).create_schema() async def enqueue_terminal( self, conn: Any, *, task_id: str, account_id: str, task_type: str, status: str, result: dict[str, Any], url: str, operation_id: str, token: str | None, authentication: TaskWebhookAuthentication | None = None, signing_scope_id: str | None = None, ) -> int: return await (await self.resolve()).enqueue_terminal( conn, task_id=task_id, account_id=account_id, task_type=task_type, status=status, result=result, url=url, operation_id=operation_id, token=token, authentication=authentication, signing_scope_id=signing_scope_id, ) def validate_registration( self, url: str, authentication: TaskWebhookAuthentication | None = None ) -> None: (self._require_resolved()).validate_registration(url, authentication) def protect_registration( self, *, account_id: str, task_id: str, task_type: str, url: str, operation_id: str, token: str | None, authentication: TaskWebhookAuthentication | None = None, signing_scope_id: str | None = None, ) -> tuple[bytes, bytes]: return (self._require_resolved()).protect_registration( account_id=account_id, task_id=task_id, task_type=task_type, url=url, operation_id=operation_id, token=token, authentication=authentication, signing_scope_id=signing_scope_id, ) def open_registration( self, *, account_id: str, task_id: str, task_type: str, encrypted_registration: bytes, nonce: bytes, ) -> tuple[str, str, str | None]: return (self._require_resolved()).open_registration( account_id=account_id, task_id=task_id, task_type=task_type, encrypted_registration=encrypted_registration, nonce=nonce, ) async def run_worker( self, *, poll_interval: float = 1.0, purge_interval: float = 300.0 ) -> None: await (await self.resolve()).run_worker( poll_interval=poll_interval, purge_interval=purge_interval ) async def process_one(self) -> bool: return await (await self.resolve()).process_one() async def purge_expired(self) -> None: await (await self.resolve()).purge_expired()Deferred concrete outbox. Synchronous crypto helpers require resolve().
Use from_registry() to share a registry factory's original outbox. Pass the original objects, constructed inside that factory, to PgTaskRegistry; do not substitute this worker facade into the concrete constructor's pool checks.
Ancestors
- adcp.decisioning.pg.lazy._LazyStore
- typing.Generic
Class variables
var delivery_state_is_durable : ClassVar[bool]var supports_atomic_task_outbox : ClassVar[bool]
Static methods
def from_registry(registry: LazyTaskRegistry) ‑> LazyTaskWebhookOutbox-
Create a worker facade resolving the same original registry/outbox pair.
Instance variables
prop delivery_retry_horizon_seconds : int-
Expand source code
@property def delivery_retry_horizon_seconds(self) -> int: return self._require_resolved().delivery_retry_horizon_seconds prop legacy_hmac_fallback : bool-
Expand source code
@property def legacy_hmac_fallback(self) -> bool: return self._require_resolved().legacy_hmac_fallback
Methods
async def create_schema(self) ‑> None-
Expand source code
async def create_schema(self) -> None: await (await self.resolve()).create_schema() async def enqueue_terminal(self,
conn: Any,
*,
task_id: str,
account_id: str,
task_type: str,
status: str,
result: dict[str, Any],
url: str,
operation_id: str,
token: str | None,
authentication: TaskWebhookAuthentication | None = None,
signing_scope_id: str | None = None) ‑> int-
Expand source code
async def enqueue_terminal( self, conn: Any, *, task_id: str, account_id: str, task_type: str, status: str, result: dict[str, Any], url: str, operation_id: str, token: str | None, authentication: TaskWebhookAuthentication | None = None, signing_scope_id: str | None = None, ) -> int: return await (await self.resolve()).enqueue_terminal( conn, task_id=task_id, account_id=account_id, task_type=task_type, status=status, result=result, url=url, operation_id=operation_id, token=token, authentication=authentication, signing_scope_id=signing_scope_id, ) def open_registration(self,
*,
account_id: str,
task_id: str,
task_type: str,
encrypted_registration: bytes,
nonce: bytes) ‑> tuple[str, str, str | None]-
Expand source code
def open_registration( self, *, account_id: str, task_id: str, task_type: str, encrypted_registration: bytes, nonce: bytes, ) -> tuple[str, str, str | None]: return (self._require_resolved()).open_registration( account_id=account_id, task_id=task_id, task_type=task_type, encrypted_registration=encrypted_registration, nonce=nonce, ) async def process_one(self) ‑> bool-
Expand source code
async def process_one(self) -> bool: return await (await self.resolve()).process_one() def protect_registration(self,
*,
account_id: str,
task_id: str,
task_type: str,
url: str,
operation_id: str,
token: str | None,
authentication: TaskWebhookAuthentication | None = None,
signing_scope_id: str | None = None) ‑> tuple[bytes, bytes]-
Expand source code
def protect_registration( self, *, account_id: str, task_id: str, task_type: str, url: str, operation_id: str, token: str | None, authentication: TaskWebhookAuthentication | None = None, signing_scope_id: str | None = None, ) -> tuple[bytes, bytes]: return (self._require_resolved()).protect_registration( account_id=account_id, task_id=task_id, task_type=task_type, url=url, operation_id=operation_id, token=token, authentication=authentication, signing_scope_id=signing_scope_id, ) async def purge_expired(self) ‑> None-
Expand source code
async def purge_expired(self) -> None: await (await self.resolve()).purge_expired() async def run_worker(self, *, poll_interval: float = 1.0, purge_interval: float = 300.0) ‑> None-
Expand source code
async def run_worker( self, *, poll_interval: float = 1.0, purge_interval: float = 300.0 ) -> None: await (await self.resolve()).run_worker( poll_interval=poll_interval, purge_interval=purge_interval ) def validate_registration(self, url: str, authentication: TaskWebhookAuthentication | None = None) ‑> None-
Expand source code
def validate_registration( self, url: str, authentication: TaskWebhookAuthentication | None = None ) -> None: (self._require_resolved()).validate_registration(url, authentication)