Module adcp.accounts
Buyer-side provisioning registry, scoped to a seller and buyer identity.
Global variables
var AccountPolicy-
Natural-key account policy for client task calls.
offsends account references unchanged, as 8.0 did, and is the 8.x default.autoandstrictopt in to registry-aware preflight;autobecomes the default in the next major release.
Functions
def has_natural_reference(request: Any) ‑> bool-
Expand source code
def has_natural_reference(request: Any) -> bool: """Return True when ``request`` carries a buyer-declared natural account key.""" if "account" not in getattr(type(request), "model_fields", {}): return False ref = getattr(request, "account", None) return ref is not None and not _payload(ref).get("account_id")Return True when
requestcarries a buyer-declared natural account key.
Classes
class AccountRecord (account_id: str | None, status: str | None)-
Expand source code
@dataclass(frozen=True) class AccountRecord: """A successful provisioning result; pending setup still means provisioned.""" account_id: str | None status: str | NoneA successful provisioning result; pending setup still means provisioned.
Instance variables
var account_id : str | Nonevar status : str | None
class AccountRegistry (client: ADCPClient,
storage: AccountStorage | None = None)-
Expand source code
class AccountRegistry: """Track provisioned keys for one seller. Exposed as ``client.accounts``.""" def __init__(self, client: ADCPClient, storage: AccountStorage | None = None) -> None: self._client = client self._storage = storage if storage is not None else InMemoryAccountStorage() self._seller = client.agent_config.agent_uri self._ensure_lock = asyncio.Lock() async def get(self, key: AccountReference | Mapping[str, Any]) -> AccountRecord | None: return await self._storage.load(self._seller, account_key(key)) async def forget(self, key: AccountReference | Mapping[str, Any]) -> None: """Invalidate a stale record, e.g. after a seller reports a missing account.""" await self._storage.delete(self._seller, account_key(key)) async def ensure( self, key: AccountReference | Mapping[str, Any], *, billing: str, payment_terms: str | None = None, billing_entity: Any = None, ) -> AccountRecord: """Provision once through sync_accounts, retaining the ID and status. Billing and payment terms are supplied only on first provisioning; use sync_accounts explicitly for subsequent settings changes. """ ref = account_key_payload(TypeAdapter(AccountReference).validate_python(key)) if "account_id" in ref: raise ValueError("ensure requires a buyer-declared natural key") async with self._ensure_lock: existing = await self.get(ref) if existing is not None: return existing caps = await self._client.fetch_capabilities() if getattr(caps.account, "require_operator_auth", False) is True: self._raise_setup("sync_accounts", "Discover an account_id through list_accounts") entry = {**ref, "billing": billing} if payment_terms is not None: entry["payment_terms"] = payment_terms if billing_entity is not None: entry["billing_entity"] = _payload(billing_entity) result = await self._client.sync_accounts( SyncAccountsRequest.model_validate( { "idempotency_key": str(uuid.uuid4()), "accounts": [entry], } ) ) body = _payload(result.data) if not result.success or body.get("errors"): raise classify_task_error( "sync_accounts", body.get("errors") or [ result.adcp_error or { "code": "ACCOUNT_NOT_FOUND", "message": result.error or "Provisioning failed", } ], agent_id=self._client.agent_config.id, ) rows = body.get("accounts", []) if not rows or rows[0].get("action") == "failed": raise classify_task_error( "sync_accounts", (rows[0].get("errors") if rows else None) or [{"code": "ACCOUNT_NOT_FOUND", "message": "Provisioning failed"}], agent_id=self._client.agent_config.id, ) record = await self.get(ref) if record is None: self._raise_setup("sync_accounts", "Seller did not confirm account provisioning") assert record is not None return record async def observe( self, operation: str, request: Any, result: Any, *, fetch_capabilities: bool = True, ) -> None: """Update from successful sync/list results; never learn from a dry run. With ``fetch_capabilities=False`` only cached capabilities are used, so observation never adds a network request; rows whose key depends on uncached capabilities are skipped rather than guessed. """ if operation not in {"sync_accounts", "list_accounts"} or not result.success: return req = _payload(request) body = _payload(result.data) if req.get("dry_run") or body.get("dry_run"): return inputs = req.get("accounts", []) for row in body.get("accounts", []): if operation == "sync_accounts": caps = self._client.capabilities if ( caps is None and fetch_capabilities and row.get("timezone") and row.get("brand") and row.get("operator") ): caps = await self._client.fetch_capabilities() timezone_caps = getattr(getattr(caps, "account", None), "timezone", None) key = _sync_key( row, inputs, buyer_selected_timezone=getattr(timezone_caps, "account_selection", None) == "buyer_selected", ) if key is None: continue else: try: key = account_key_payload(_row_reference(row)) except ValueError: continue if row.get("action") == "failed": if any( error.get("code") == "ACCOUNT_NOT_FOUND" for error in row.get("errors") or [] ): await self.forget(key) continue if operation == "list_accounts" and "account_id" not in key: caps = ( await self._client.fetch_capabilities() if fetch_capabilities else self._client.capabilities ) if caps is None: continue timezone_caps = getattr(getattr(caps, "account", None), "timezone", None) if getattr(timezone_caps, "account_selection", None) != "buyer_selected": key.pop("timezone", None) if row.get("action") is None and operation == "sync_accounts": continue await self._storage.save( self._seller, account_key(key), AccountRecord(account_id=row.get("account_id"), status=row.get("status")), ) async def prepare(self, operation: str, request: BaseModel, policy: AccountPolicy) -> BaseModel: """Apply account policy to a copy, preserving the caller's request. ``off`` returns the request unchanged with no capability fetch. """ if policy not in ACCOUNT_POLICIES: raise ValueError("account_policy must be 'off', 'auto', or 'strict'") if policy == "off": return request if operation == "list_accounts" or "account" not in type(request).model_fields: return request ref = getattr(request, "account", None) wire = _payload(ref) # Explicit IDs are already seller-assigned handles and may be supplied # out of band. Never infer that a returned handle replaces a natural key. if wire.get("account_id"): return request caps = self._client.capabilities if ref is not None: caps = await self._client.fetch_capabilities() account_caps = getattr(caps, "account", None) required = type(request).model_fields["account"].is_required() or ( operation in {"get_products", "list_products"} and getattr(account_caps, "required_for_products", False) is True ) if ref is None: if required: self._raise_setup( operation, "An account is required; provision it before this call" ) return request requires_operator_auth = getattr(account_caps, "require_operator_auth", False) is True if requires_operator_auth and ( required or policy == "strict" or operation not in _PUBLIC_DISCOVERY_TASKS ): self._raise_setup(operation, "This seller requires an account_id from list_accounts") record = None if requires_operator_auth else await self.get(wire) if ( record is not None and operation in {"get_products", "list_products"} and record.status not in {"active", "payment_required"} ): if required or policy == "strict": self._raise_setup(operation, "Account setup must allow product discovery") # Inactive accounts can still browse the public feed where allowed. record = None if record is None: if required or policy == "strict" or operation not in _PUBLIC_DISCOVERY_TASKS: raise AccountNotFoundError( operation, [ { "code": "ACCOUNT_NOT_FOUND", "message": "Provision this key with client.accounts.ensure first", } ], agent_id=self._client.agent_config.id, ) updates: dict[str, Any] = {"account": None} if "brand" in type(request).model_fields and getattr(request, "brand", None) is None: brand = wire.get("brand") if operation == "list_products": brand = account_key_payload(wire)["brand"] updates["brand"] = TypeAdapter( type(request).model_fields["brand"].annotation ).validate_python(brand) # Account-scoped version tokens cannot be reused for a public read. for name in ("if_wholesale_feed_version", "if_feed_version", "if_pricing_version"): if name in type(request).model_fields: updates[name] = None return request.model_copy(update=updates) if operation in {"create_media_buy", "buy_products", "accept_proposal", "activate_signal"}: if record.status == "payment_required": raise AccountPaymentRequiredError( operation, [{"code": "ACCOUNT_PAYMENT_REQUIRED", "message": "Account needs payment"}], agent_id=self._client.agent_config.id, ) if record.status != "active": self._raise_setup(operation, "Account must be active before committing spend") return request def _raise_setup(self, operation: str, message: str) -> None: raise AccountSetupRequiredError( operation, [{"code": "ACCOUNT_SETUP_REQUIRED", "message": message}], agent_id=self._client.agent_config.id, )Track provisioned keys for one seller. Exposed as
client.accounts.Methods
async def ensure(self,
key: AccountReference | Mapping[str, Any],
*,
billing: str,
payment_terms: str | None = None,
billing_entity: Any = None) ‑> AccountRecord-
Expand source code
async def ensure( self, key: AccountReference | Mapping[str, Any], *, billing: str, payment_terms: str | None = None, billing_entity: Any = None, ) -> AccountRecord: """Provision once through sync_accounts, retaining the ID and status. Billing and payment terms are supplied only on first provisioning; use sync_accounts explicitly for subsequent settings changes. """ ref = account_key_payload(TypeAdapter(AccountReference).validate_python(key)) if "account_id" in ref: raise ValueError("ensure requires a buyer-declared natural key") async with self._ensure_lock: existing = await self.get(ref) if existing is not None: return existing caps = await self._client.fetch_capabilities() if getattr(caps.account, "require_operator_auth", False) is True: self._raise_setup("sync_accounts", "Discover an account_id through list_accounts") entry = {**ref, "billing": billing} if payment_terms is not None: entry["payment_terms"] = payment_terms if billing_entity is not None: entry["billing_entity"] = _payload(billing_entity) result = await self._client.sync_accounts( SyncAccountsRequest.model_validate( { "idempotency_key": str(uuid.uuid4()), "accounts": [entry], } ) ) body = _payload(result.data) if not result.success or body.get("errors"): raise classify_task_error( "sync_accounts", body.get("errors") or [ result.adcp_error or { "code": "ACCOUNT_NOT_FOUND", "message": result.error or "Provisioning failed", } ], agent_id=self._client.agent_config.id, ) rows = body.get("accounts", []) if not rows or rows[0].get("action") == "failed": raise classify_task_error( "sync_accounts", (rows[0].get("errors") if rows else None) or [{"code": "ACCOUNT_NOT_FOUND", "message": "Provisioning failed"}], agent_id=self._client.agent_config.id, ) record = await self.get(ref) if record is None: self._raise_setup("sync_accounts", "Seller did not confirm account provisioning") assert record is not None return recordProvision once through sync_accounts, retaining the ID and status.
Billing and payment terms are supplied only on first provisioning; use sync_accounts explicitly for subsequent settings changes.
async def forget(self, key: AccountReference | Mapping[str, Any]) ‑> None-
Expand source code
async def forget(self, key: AccountReference | Mapping[str, Any]) -> None: """Invalidate a stale record, e.g. after a seller reports a missing account.""" await self._storage.delete(self._seller, account_key(key))Invalidate a stale record, e.g. after a seller reports a missing account.
async def get(self, key: AccountReference | Mapping[str, Any]) ‑> AccountRecord | None-
Expand source code
async def get(self, key: AccountReference | Mapping[str, Any]) -> AccountRecord | None: return await self._storage.load(self._seller, account_key(key)) async def observe(self, operation: str, request: Any, result: Any, *, fetch_capabilities: bool = True) ‑> None-
Expand source code
async def observe( self, operation: str, request: Any, result: Any, *, fetch_capabilities: bool = True, ) -> None: """Update from successful sync/list results; never learn from a dry run. With ``fetch_capabilities=False`` only cached capabilities are used, so observation never adds a network request; rows whose key depends on uncached capabilities are skipped rather than guessed. """ if operation not in {"sync_accounts", "list_accounts"} or not result.success: return req = _payload(request) body = _payload(result.data) if req.get("dry_run") or body.get("dry_run"): return inputs = req.get("accounts", []) for row in body.get("accounts", []): if operation == "sync_accounts": caps = self._client.capabilities if ( caps is None and fetch_capabilities and row.get("timezone") and row.get("brand") and row.get("operator") ): caps = await self._client.fetch_capabilities() timezone_caps = getattr(getattr(caps, "account", None), "timezone", None) key = _sync_key( row, inputs, buyer_selected_timezone=getattr(timezone_caps, "account_selection", None) == "buyer_selected", ) if key is None: continue else: try: key = account_key_payload(_row_reference(row)) except ValueError: continue if row.get("action") == "failed": if any( error.get("code") == "ACCOUNT_NOT_FOUND" for error in row.get("errors") or [] ): await self.forget(key) continue if operation == "list_accounts" and "account_id" not in key: caps = ( await self._client.fetch_capabilities() if fetch_capabilities else self._client.capabilities ) if caps is None: continue timezone_caps = getattr(getattr(caps, "account", None), "timezone", None) if getattr(timezone_caps, "account_selection", None) != "buyer_selected": key.pop("timezone", None) if row.get("action") is None and operation == "sync_accounts": continue await self._storage.save( self._seller, account_key(key), AccountRecord(account_id=row.get("account_id"), status=row.get("status")), )Update from successful sync/list results; never learn from a dry run.
With
fetch_capabilities=Falseonly cached capabilities are used, so observation never adds a network request; rows whose key depends on uncached capabilities are skipped rather than guessed. async def prepare(self,
operation: str,
request: BaseModel,
policy: AccountPolicy) ‑> BaseModel-
Expand source code
async def prepare(self, operation: str, request: BaseModel, policy: AccountPolicy) -> BaseModel: """Apply account policy to a copy, preserving the caller's request. ``off`` returns the request unchanged with no capability fetch. """ if policy not in ACCOUNT_POLICIES: raise ValueError("account_policy must be 'off', 'auto', or 'strict'") if policy == "off": return request if operation == "list_accounts" or "account" not in type(request).model_fields: return request ref = getattr(request, "account", None) wire = _payload(ref) # Explicit IDs are already seller-assigned handles and may be supplied # out of band. Never infer that a returned handle replaces a natural key. if wire.get("account_id"): return request caps = self._client.capabilities if ref is not None: caps = await self._client.fetch_capabilities() account_caps = getattr(caps, "account", None) required = type(request).model_fields["account"].is_required() or ( operation in {"get_products", "list_products"} and getattr(account_caps, "required_for_products", False) is True ) if ref is None: if required: self._raise_setup( operation, "An account is required; provision it before this call" ) return request requires_operator_auth = getattr(account_caps, "require_operator_auth", False) is True if requires_operator_auth and ( required or policy == "strict" or operation not in _PUBLIC_DISCOVERY_TASKS ): self._raise_setup(operation, "This seller requires an account_id from list_accounts") record = None if requires_operator_auth else await self.get(wire) if ( record is not None and operation in {"get_products", "list_products"} and record.status not in {"active", "payment_required"} ): if required or policy == "strict": self._raise_setup(operation, "Account setup must allow product discovery") # Inactive accounts can still browse the public feed where allowed. record = None if record is None: if required or policy == "strict" or operation not in _PUBLIC_DISCOVERY_TASKS: raise AccountNotFoundError( operation, [ { "code": "ACCOUNT_NOT_FOUND", "message": "Provision this key with client.accounts.ensure first", } ], agent_id=self._client.agent_config.id, ) updates: dict[str, Any] = {"account": None} if "brand" in type(request).model_fields and getattr(request, "brand", None) is None: brand = wire.get("brand") if operation == "list_products": brand = account_key_payload(wire)["brand"] updates["brand"] = TypeAdapter( type(request).model_fields["brand"].annotation ).validate_python(brand) # Account-scoped version tokens cannot be reused for a public read. for name in ("if_wholesale_feed_version", "if_feed_version", "if_pricing_version"): if name in type(request).model_fields: updates[name] = None return request.model_copy(update=updates) if operation in {"create_media_buy", "buy_products", "accept_proposal", "activate_signal"}: if record.status == "payment_required": raise AccountPaymentRequiredError( operation, [{"code": "ACCOUNT_PAYMENT_REQUIRED", "message": "Account needs payment"}], agent_id=self._client.agent_config.id, ) if record.status != "active": self._raise_setup(operation, "Account must be active before committing spend") return requestApply account policy to a copy, preserving the caller's request.
offreturns the request unchanged with no capability fetch.
class AccountStorage (*args, **kwargs)-
Expand source code
class AccountStorage(Protocol): """Persist records under seller URI and complete natural key. An adapter MUST be private to one authenticated buyer identity. Sharing it across credentials requires a separate namespace for each buyer. """ async def load(self, seller: str, key: str) -> AccountRecord | None: ... async def save(self, seller: str, key: str, record: AccountRecord) -> None: ... async def delete(self, seller: str, key: str) -> None: ...Persist records under seller URI and complete natural key.
An adapter MUST be private to one authenticated buyer identity. Sharing it across credentials requires a separate namespace for each buyer.
Ancestors
- typing.Protocol
- typing.Generic
Methods
async def delete(self, seller: str, key: str) ‑> None-
Expand source code
async def delete(self, seller: str, key: str) -> None: ... async def load(self, seller: str, key: str) ‑> AccountRecord | None-
Expand source code
async def load(self, seller: str, key: str) -> AccountRecord | None: ... async def save(self,
seller: str,
key: str,
record: AccountRecord) ‑> None-
Expand source code
async def save(self, seller: str, key: str, record: AccountRecord) -> None: ...
class InMemoryAccountStorage-
Expand source code
class InMemoryAccountStorage: """Default storage; pass a durable adapter to survive process restarts.""" def __init__(self) -> None: self._records: dict[tuple[str, str], AccountRecord] = {} async def load(self, seller: str, key: str) -> AccountRecord | None: return self._records.get((seller, key)) async def save(self, seller: str, key: str, record: AccountRecord) -> None: self._records[seller, key] = record async def delete(self, seller: str, key: str) -> None: self._records.pop((seller, key), None)Default storage; pass a durable adapter to survive process restarts.
Methods
async def delete(self, seller: str, key: str) ‑> None-
Expand source code
async def delete(self, seller: str, key: str) -> None: self._records.pop((seller, key), None) async def load(self, seller: str, key: str) ‑> AccountRecord | None-
Expand source code
async def load(self, seller: str, key: str) -> AccountRecord | None: return self._records.get((seller, key)) async def save(self,
seller: str,
key: str,
record: AccountRecord) ‑> None-
Expand source code
async def save(self, seller: str, key: str, record: AccountRecord) -> None: self._records[seller, key] = record