Module adcp.decisioning.pg.task_registry
PostgreSQL-backed :class:~adcp.decisioning.TaskRegistry implementation.
Durable counterpart to :class:~adcp.decisioning.InMemoryTaskRegistry:
task state survives process restarts and is safe for multi-worker deployments
sharing a single Postgres database.
The caller supplies an :class:psycopg_pool.AsyncConnectionPool. We don't
open, own, or close the pool — adopters typically share an existing pool with
their main application database.
Quickstart
::
import asyncio
from psycopg_pool import AsyncConnectionPool
from adcp.decisioning import PgTaskRegistry, serve
from myapp import MyPlatform
async def main():
async with AsyncConnectionPool(
"postgresql://user:pass@localhost/mydb",
min_size=2,
max_size=10,
) as pool:
registry = PgTaskRegistry(pool=pool)
await registry.create_schema() # idempotent; safe on every boot
serve(MyPlatform(), registry=registry)
asyncio.run(main())
Schema Bootstrap
Call :meth:create_schema once per deployment (or every boot — it is
idempotent via CREATE TABLE IF NOT EXISTS). The equivalent raw DDL
ships at :file:src/adcp/decisioning/pg/decisioning_tasks.sql for adopters
using a migration tool (Alembic, Flyway, psql).
Cross-tenant safety
:meth:get enforces account isolation at the SQL level —
WHERE account_id = %s is part of the query predicate, not a Python-level
filter. A mis-matched expected_account_id returns None without
materializing the row.
Multi-worker concurrency
Terminal-state transitions (:meth:complete, :meth:fail) use an atomic
UPDATE ... WHERE state NOT IN ('completed', 'failed') RETURNING task_id
pattern. If the UPDATE lands zero rows, a follow-up SELECT determines whether
the task is unknown or already terminal, enabling correct idempotency
behavior across workers without optimistic-lock retries.
:meth:update_progress similarly uses a conditional UPDATE that silently
no-ops on terminal rows, so a straggler progress write can never resurrect a
completed task.
Classes
class PgTaskRegistry (*,
pool: AsyncConnectionPool,
task_webhook_outbox: PgTaskWebhookOutbox | None = None,
webhook_signing_scope_resolver: WebhookSigningScopeResolver | None = None)-
Expand source code
class PgTaskRegistry(_TaskLifecycleObservers): """PostgreSQL-backed :class:`~adcp.decisioning.TaskRegistry` — v6.1. Durable counterpart to :class:`~adcp.decisioning.InMemoryTaskRegistry`. Set ``is_durable = True`` so the production-mode gate in :func:`adcp.decisioning.serve.create_adcp_server_from_platform` accepts it without requiring ``ADCP_DECISIONING_ALLOW_INMEMORY_TASKS=1``. Parameters ---------- pool: An :class:`psycopg_pool.AsyncConnectionPool` owned by the caller. Each registry operation acquires a short-lived connection from the pool and returns it immediately after the query. No long-lived transactions, no cross-operation state. task_webhook_outbox: Optional :class:`PgTaskWebhookOutbox` sharing this pool. When a task carries callback registration, ``complete`` / ``fail`` enqueue its immutable terminal webhook on the same transaction connection. Notes ----- Unlike :class:`~adcp.signing.PgReplayStore`, this class uses a fixed ``decisioning_tasks`` table name. Multi-tenant table-name isolation is not supported in this release — callers requiring strict schema separation should use separate databases or schemas. """ is_durable: ClassVar[bool] = True def __init__( self, *, pool: AsyncConnectionPool, task_webhook_outbox: PgTaskWebhookOutbox | None = None, webhook_signing_scope_resolver: WebhookSigningScopeResolver | None = None, _table: str = _DEFAULT_TABLE, ) -> None: if not PG_AVAILABLE: raise ImportError(_INSTALL_HINT) if not _SAFE_IDENTIFIER_RE.fullmatch(_table): raise ValueError(f"_table must match [a-z_][a-z0-9_]* (ASCII only), got {_table!r}") if task_webhook_outbox is not None and task_webhook_outbox._pool is not pool: raise ValueError( "PgTaskRegistry and PgTaskWebhookOutbox must use the same connection pool" ) uses_sender_resolver = ( task_webhook_outbox is not None and task_webhook_outbox._sender_resolver is not None ) if uses_sender_resolver != (webhook_signing_scope_resolver is not None): raise ValueError( "webhook_signing_scope_resolver is required exactly when the " "PgTaskWebhookOutbox uses sender_resolver" ) self._init_lifecycle_observers() self._pool = pool self._table = _table self.task_webhook_outbox = task_webhook_outbox self._webhook_signing_scope_resolver = webhook_signing_scope_resolver self.atomic_task_webhook_outbox = task_webhook_outbox is not None # Pre-format queries at construction so the hot path avoids f-strings per call. # _table is whitelisted by _SAFE_IDENTIFIER_RE above. self._sql_insert = ( # noqa: S608 — table name is whitelisted f"INSERT INTO {self._table}" # nosec B608 — table identifier is constructor-validated f" (task_id, account_id, state, task_type, request_context," f" webhook_registration, webhook_registration_nonce, created_at, updated_at)" f" VALUES (%s, %s, 'submitted', %s, %s::jsonb, %s, %s, %s, %s)" ) self._sql_update_progress = ( # noqa: S608 f"WITH previous AS (SELECT task_id, state FROM {self._table}" # nosec B608 — table identifier is constructor-validated f" WHERE task_id = %s AND state NOT IN ('completed', 'failed') FOR UPDATE)" f" UPDATE {self._table} AS task" f" SET state = CASE task.state WHEN 'submitted' THEN 'working' ELSE task.state END," f" progress = %s::jsonb, updated_at = %s" f" FROM previous WHERE task.task_id = previous.task_id" f" RETURNING previous.state, task.task_id, task.account_id, task.task_type," f" task.created_at, task.updated_at" ) self._sql_complete = ( # noqa: S608 f"UPDATE {self._table}" # nosec B608 — table identifier is constructor-validated f" SET state = 'completed', result = %s::jsonb, updated_at = %s" f" WHERE task_id = %s AND state NOT IN ('completed', 'failed')" f" RETURNING task_id, account_id, task_type, webhook_registration," f" webhook_registration_nonce, created_at, updated_at" ) self._sql_fail = ( # noqa: S608 f"UPDATE {self._table}" # nosec B608 — table identifier is constructor-validated f" SET state = 'failed', error = %s::jsonb, updated_at = %s" f" WHERE task_id = %s AND state NOT IN ('completed', 'failed')" f" RETURNING task_id, account_id, task_type, webhook_registration," f" webhook_registration_nonce, created_at, updated_at" ) self._sql_clear_webhook_registration = ( # noqa: S608 f"UPDATE {self._table} SET webhook_registration = NULL," # nosec B608 — table identifier is constructor-validated f" webhook_registration_nonce = NULL WHERE task_id = %s" ) # Explicit ``::text`` cast on the optional account-filter # parameter so psycopg's bind-param type inference doesn't # fail with ``IndeterminateDatatype: could not determine # data type of parameter $2``. Without the cast, the # ``%s IS NULL`` predicate gives psycopg no type context # for the parameter and the query fails at prepare time. self._sql_get = ( # noqa: S608 f"SELECT task_id, account_id, state, task_type," # nosec B608 — table identifier is constructor-validated f" progress, result, error, request_context, created_at, updated_at" f" FROM {self._table}" f" WHERE task_id = %s AND (%s::text IS NULL OR account_id = %s)" ) self._sql_get_state_result = ( # noqa: S608 f"SELECT state, result, task_type FROM {self._table}" # nosec B608 — validated identifier " WHERE task_id = %s" ) self._sql_get_task_type = f"SELECT task_type FROM {self._table} WHERE task_id = %s" # noqa: S608 # nosec B608 — table identifier is constructor-validated self._sql_get_state_error = f"SELECT state, error FROM {self._table} WHERE task_id = %s" # noqa: S608 # nosec B608 — table identifier is constructor-validated self._sql_discard = ( # noqa: S608 f"DELETE FROM {self._table} WHERE task_id = %s" # nosec B608 — table identifier is constructor-validated f" RETURNING task_id, account_id, task_type, created_at, updated_at" ) self._sql_ddl = ( # noqa: S608 f"CREATE TABLE IF NOT EXISTS {self._table} (" f' task_id TEXT COLLATE "C" NOT NULL PRIMARY KEY,' f' account_id TEXT COLLATE "C" NOT NULL,' f" state TEXT NOT NULL DEFAULT 'submitted'," f" task_type TEXT NOT NULL," f" progress JSONB," f" result JSONB," f" error JSONB," f" request_context JSONB," f" webhook_registration BYTEA," f" webhook_registration_nonce BYTEA," f" created_at DOUBLE PRECISION NOT NULL," f" updated_at DOUBLE PRECISION NOT NULL" f");" f"ALTER TABLE {self._table} ADD COLUMN IF NOT EXISTS request_context JSONB;" f"ALTER TABLE {self._table} ADD COLUMN IF NOT EXISTS webhook_registration BYTEA;" f"ALTER TABLE {self._table} ADD COLUMN IF NOT EXISTS webhook_registration_nonce BYTEA;" f"CREATE INDEX IF NOT EXISTS {self._table}_account_idx" # noqa: S608 f" ON {self._table} (account_id);" ) # -- schema bootstrap ----------------------------------------------- async def create_schema(self) -> None: """Create the task registry table and supporting index. Honors the ``_table`` kwarg the store was constructed with. Idempotent via ``CREATE TABLE IF NOT EXISTS`` — safe to call on every application boot. The equivalent raw DDL ships at ``adcp/decisioning/pg/decisioning_tasks.sql`` in the installed package for adopters using a migration tool (Alembic, Flyway, psql). """ async with self._pool.connection() as conn: await conn.execute(self._sql_ddl) # -- TaskRegistry Protocol ------------------------------------------ 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: """Allocate a task_id, persist a ``submitted`` row, return the id. Mirrors :meth:`~adcp.decisioning.InMemoryTaskRegistry.issue` including the account_id validation guard — empty or sentinel account_ids would allow cross-tenant task-id probing via the ``WHERE account_id = %s`` predicate collapsing multiple tenants into one slot. """ if not account_id or not account_id.strip() or account_id == "<unset>": raise ValueError( f"account_id must be a non-empty, non-default string; " f"got {account_id!r}. AccountStore.resolve must always " "return Account(id=<non-empty>) so cross-tenant cache " "scoping works correctly." ) if (webhook_url is None) != (webhook_operation_id is None): raise ValueError("webhook_url and webhook_operation_id must be supplied together") if webhook_url is not None and not webhook_url: raise ValueError("webhook_url must be non-empty when supplied") if webhook_operation_id is not None and not webhook_operation_id: raise ValueError("webhook_operation_id must be non-empty when supplied") if webhook_url is None and webhook_signing_scope_id is not None: raise ValueError("webhook_signing_scope_id requires webhook_url") if webhook_url is None and webhook_authentication is not None: raise ValueError("webhook_authentication requires webhook_url") if webhook_authentication is not None and webhook_signing_scope_id is not None: raise ValueError( "legacy webhook_authentication must not carry an RFC 9421 signing scope" ) outbox = self.task_webhook_outbox if webhook_url is not None: if outbox is None: raise ValueError("webhook registration requires the registry's task_webhook_outbox") outbox.validate_registration(webhook_url, webhook_authentication) task_id = f"task_{uuid.uuid4().hex[:16]}" encrypted_registration: bytes | None = None registration_nonce: bytes | None = None if webhook_url is not None: if webhook_operation_id is None or outbox is None: raise RuntimeError("validated webhook registration became incomplete") encrypted_registration, registration_nonce = outbox.protect_registration( account_id=account_id, task_id=task_id, task_type=task_type, url=webhook_url, operation_id=webhook_operation_id, token=webhook_token, authentication=webhook_authentication, signing_scope_id=webhook_signing_scope_id, ) now = time.time() async with self._pool.connection() as conn: await conn.execute( self._sql_insert, ( task_id, account_id, task_type, json.dumps(request_context) if request_context is not None else None, encrypted_registration, registration_nonce, now, now, ), ) self._notify_lifecycle_observers( "submitted", { "task_id": task_id, "account_id": account_id, "task_type": task_type, "created_at": now, "updated_at": now, }, ) return task_id async def resolve_webhook_signing_scope( self, context: RequestContext[Any], ) -> str | None: """Derive an opaque signing scope from trusted framework context. The callback is operator wiring and receives the hydrated :class:`RequestContext`. It must use internal tenant/platform metadata, never buyer request ``context``, ``push_notification_config``, or an unqualified buyer account id. """ resolver = self._webhook_signing_scope_resolver if resolver is None: return None from adcp.decisioning.types import AdcpError try: value: object = resolver(context) if inspect.isawaitable(value): value = await value except Exception: raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolution failed", recovery="terminal", ) from None if not isinstance(value, str): raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolver returned an invalid value", recovery="terminal", ) # Reuse the outbox's bounded opaque-ID validation before any task row # is issued. This is trusted server state, but it still crosses a DB # and authenticated-envelope boundary. if self.task_webhook_outbox is None: raise RuntimeError("signing scope resolver requires a task webhook outbox") self.task_webhook_outbox._validate_signing_scope_id(value) return value async def update_progress( self, task_id: str, progress: dict[str, Any], ) -> None: """Write a progress payload; transition ``submitted`` → ``working``. Silently no-ops when the task is already in a terminal state or unknown — the dispatch wrapper expects this method never to raise on transient conditions (see :class:`~adcp.decisioning.TaskRegistry` docstring). The ``state NOT IN ('completed', 'failed')`` predicate is evaluated server-side so a concurrent terminal write cannot be overwritten by a straggler progress event. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_update_progress, (task_id, json.dumps(progress), time.time()), ) row = await cur.fetchone() if row is not None and row[0] == "submitted": self._notify_lifecycle_observers("working", self._lifecycle_record(row[1:])) @staticmethod def _lifecycle_record(row: tuple[Any, ...]) -> dict[str, Any]: return dict(zip(("task_id", "account_id", "task_type", "created_at", "updated_at"), row)) async def complete(self, task_id: str, result: dict[str, Any]) -> None: """Complete once; notify metrics only after the transaction commits.""" row = await self._complete(task_id, result) if row is not None: self._notify_lifecycle_observers("completed", self._lifecycle_record(row[:3] + row[5:])) async def _complete( self, task_id: str, result: dict[str, Any], ) -> tuple[Any, ...] | None: """Mark the task ``completed`` with ``result`` as the terminal artifact. Idempotent on repeated calls with an equal ``result``; raises :class:`ValueError` on conflicting re-completion. Uses an atomic ``UPDATE ... RETURNING`` so concurrent workers cannot race each other into double-completion without detection. """ async with self._pool.connection() as conn: # Do not rely on the pool connection's implicit transaction. # Explicitly bind the terminal transition, outbox insert, and # callback-registration clear even when adopters configure the # pool with autocommit=True. async with conn.transaction(): type_cursor = await conn.execute(self._sql_get_task_type, (task_id,)) type_row = await type_cursor.fetchone() safe_result = ( strip_credentials_from_wire_result(type_row[0], result) if type_row is not None else result ) cur = await conn.execute( self._sql_complete, (json.dumps(safe_result), time.time(), task_id), ) row: tuple[Any, ...] | None = await cur.fetchone() if row is not None: await self._enqueue_terminal_if_registered( conn, row=row, status="completed", payload=safe_result, ) return row # connection and transaction exit before the caller notifies # Zero rows in RETURNING — task is unknown or already terminal. cur2 = await conn.execute(self._sql_get_state_result, (task_id,)) row = await cur2.fetchone() if row is None: raise ValueError(f"Task {task_id!r} not found") state, existing_result, task_type = row safe_result = strip_credentials_from_wire_result(task_type, result) if state == "completed": if existing_result == safe_result: return None # idempotent raise ValueError(f"Task {task_id!r} already completed with a different result") raise ValueError(f"Task {task_id!r} already in terminal state {state!r}") async def fail(self, task_id: str, error: dict[str, Any]) -> None: """Fail once; notify metrics only after the transaction commits.""" row = await self._fail(task_id, error) if row is not None: self._notify_lifecycle_observers("failed", self._lifecycle_record(row[:3] + row[5:])) async def _fail( self, task_id: str, error: dict[str, Any], ) -> tuple[Any, ...] | None: """Mark the task ``failed`` with ``error`` as the terminal payload. Idempotent on repeated calls with an equal ``error``; raises :class:`ValueError` on conflicting re-failure. """ async with self._pool.connection() as conn: async with conn.transaction(): cur = await conn.execute( self._sql_fail, (json.dumps(error), time.time(), task_id), ) row: tuple[Any, ...] | None = await cur.fetchone() if row is not None: await self._enqueue_terminal_if_registered( conn, row=row, status="failed", payload=error, ) return row # notify only after connection/transaction exit # Zero rows in RETURNING — task is unknown or already terminal. cur2 = await conn.execute(self._sql_get_state_error, (task_id,)) row = await cur2.fetchone() if row is None: raise ValueError(f"Task {task_id!r} not found") state, existing_error = row if state == "failed": if existing_error == error: return None # idempotent raise ValueError(f"Task {task_id!r} already failed with a different error") raise ValueError(f"Task {task_id!r} already in terminal state {state!r}") async def get( self, task_id: str, *, expected_account_id: str | None = None, ) -> dict[str, Any] | None: """Look up a task record; cross-tenant probes return ``None``. The ``expected_account_id`` predicate is enforced at the SQL level (``WHERE account_id = %s``), not as a Python-level filter after fetch. This guarantees the row is never materialized for a mismatched probe, eliminating the fetch-then-filter anti-pattern. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_get, (task_id, expected_account_id, expected_account_id) ) row = await cur.fetchone() if row is None: return None return { "task_id": row[0], "account_id": row[1], "state": row[2], "task_type": row[3], "progress": row[4], "result": row[5], "error": row[6], "created_at": row[8], "updated_at": row[9], **({"context": row[7]} if row[7] is not None else {}), } 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]: from adcp.decisioning.task_queries import list_task_records # Account isolation is a SQL predicate, before rows are materialized. async with self._pool.connection() as conn: cur = await conn.execute( f"SELECT task_id, account_id, state, task_type, progress, result, error," # nosec B608 — table identifier is constructor-validated f" request_context, created_at, updated_at," f" (webhook_registration IS NOT NULL) FROM {self._table} WHERE account_id = %s", (account_id,), ) rows = await cur.fetchall() records = [ dict( zip( ( "task_id", "account_id", "state", "task_type", "progress", "result", "error", "context", "created_at", "updated_at", "has_webhook", ), row, ) ) for row in rows ] return list_task_records( records, account_id=account_id, filters=filters, sort=sort, pagination=pagination ) async def _enqueue_terminal_if_registered( self, conn: Any, *, row: tuple[Any, ...], status: str, payload: dict[str, Any], ) -> None: """Enqueue on ``conn`` so task state and webhook commit atomically.""" task_id, account_id, task_type, encrypted_registration, registration_nonce = row[:5] if encrypted_registration is None: return if self.task_webhook_outbox is None: raise RuntimeError( "Task carries push_notification_config but PgTaskRegistry has no " "task_webhook_outbox; refusing a non-atomic terminal transition" ) if registration_nonce is None: raise RuntimeError(f"Task {task_id!r} has incomplete webhook registration") url, operation_id, token, authentication, signing_scope_id = ( self.task_webhook_outbox._open_registration_with_scope( account_id=account_id, task_id=task_id, task_type=task_type, encrypted_registration=bytes(encrypted_registration), nonce=bytes(registration_nonce), ) ) await self.task_webhook_outbox.enqueue_terminal( conn, task_id=task_id, account_id=account_id, task_type=task_type, status=status, result=payload, url=url, operation_id=operation_id, token=token, authentication=authentication, signing_scope_id=signing_scope_id, ) # The encrypted outbox envelope now owns the callback registration. # Clear the task-row copy in this same transaction. await conn.execute(self._sql_clear_webhook_registration, (task_id,)) async def discard(self, task_id: str) -> None: """Remove a task_id from the registry — rollback path. Idempotent: discarding an unknown task_id is a no-op (no raise), matching the :class:`~adcp.decisioning.InMemoryTaskRegistry` contract. """ async with self._pool.connection() as conn: cur = await conn.execute(self._sql_discard, (task_id,)) row = await cur.fetchone() if row is not None: event_record = self._lifecycle_record(row) event_record["updated_at"] = time.time() self._notify_lifecycle_observers("discarded", event_record)PostgreSQL-backed :class:
~adcp.decisioning.TaskRegistry— v6.1.Durable counterpart to :class:
~adcp.decisioning.InMemoryTaskRegistry. Setis_durable = Trueso the production-mode gate in :func:adcp.decisioning.serve.create_adcp_server_from_platformaccepts it without requiringADCP_DECISIONING_ALLOW_INMEMORY_TASKS=1.Parameters
pool: An :class:
psycopg_pool.AsyncConnectionPoolowned by the caller. Each registry operation acquires a short-lived connection from the pool and returns it immediately after the query. No long-lived transactions, no cross-operation state. task_webhook_outbox: Optional :class:PgTaskWebhookOutboxsharing this pool. When a task carries callback registration,complete/failenqueue its immutable terminal webhook on the same transaction connection.Notes
Unlike :class:
~adcp.signing.PgReplayStore, this class uses a fixeddecisioning_taskstable name. Multi-tenant table-name isolation is not supported in this release — callers requiring strict schema separation should use separate databases or schemas.Ancestors
- adcp.decisioning.task_registry._TaskLifecycleObservers
Class variables
var is_durable : ClassVar[bool]
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: """Complete once; notify metrics only after the transaction commits.""" row = await self._complete(task_id, result) if row is not None: self._notify_lifecycle_observers("completed", self._lifecycle_record(row[:3] + row[5:]))Complete once; notify metrics only after the transaction commits.
async def create_schema(self) ‑> None-
Expand source code
async def create_schema(self) -> None: """Create the task registry table and supporting index. Honors the ``_table`` kwarg the store was constructed with. Idempotent via ``CREATE TABLE IF NOT EXISTS`` — safe to call on every application boot. The equivalent raw DDL ships at ``adcp/decisioning/pg/decisioning_tasks.sql`` in the installed package for adopters using a migration tool (Alembic, Flyway, psql). """ async with self._pool.connection() as conn: await conn.execute(self._sql_ddl)Create the task registry table and supporting index.
Honors the
_tablekwarg the store was constructed with. Idempotent viaCREATE TABLE IF NOT EXISTS— safe to call on every application boot. The equivalent raw DDL ships atadcp/decisioning/pg/decisioning_tasks.sqlin the installed package for adopters using a migration tool (Alembic, Flyway, psql). async def discard(self, task_id: str) ‑> None-
Expand source code
async def discard(self, task_id: str) -> None: """Remove a task_id from the registry — rollback path. Idempotent: discarding an unknown task_id is a no-op (no raise), matching the :class:`~adcp.decisioning.InMemoryTaskRegistry` contract. """ async with self._pool.connection() as conn: cur = await conn.execute(self._sql_discard, (task_id,)) row = await cur.fetchone() if row is not None: event_record = self._lifecycle_record(row) event_record["updated_at"] = time.time() self._notify_lifecycle_observers("discarded", event_record)Remove a task_id from the registry — rollback path.
Idempotent: discarding an unknown task_id is a no-op (no raise), matching the :class:
~adcp.decisioning.InMemoryTaskRegistrycontract. 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: """Fail once; notify metrics only after the transaction commits.""" row = await self._fail(task_id, error) if row is not None: self._notify_lifecycle_observers("failed", self._lifecycle_record(row[:3] + row[5:]))Fail once; notify metrics only after the transaction commits.
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: """Look up a task record; cross-tenant probes return ``None``. The ``expected_account_id`` predicate is enforced at the SQL level (``WHERE account_id = %s``), not as a Python-level filter after fetch. This guarantees the row is never materialized for a mismatched probe, eliminating the fetch-then-filter anti-pattern. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_get, (task_id, expected_account_id, expected_account_id) ) row = await cur.fetchone() if row is None: return None return { "task_id": row[0], "account_id": row[1], "state": row[2], "task_type": row[3], "progress": row[4], "result": row[5], "error": row[6], "created_at": row[8], "updated_at": row[9], **({"context": row[7]} if row[7] is not None else {}), }Look up a task record; cross-tenant probes return
None.The
expected_account_idpredicate is enforced at the SQL level (WHERE account_id = %s), not as a Python-level filter after fetch. This guarantees the row is never materialized for a mismatched probe, eliminating the fetch-then-filter anti-pattern. 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: """Allocate a task_id, persist a ``submitted`` row, return the id. Mirrors :meth:`~adcp.decisioning.InMemoryTaskRegistry.issue` including the account_id validation guard — empty or sentinel account_ids would allow cross-tenant task-id probing via the ``WHERE account_id = %s`` predicate collapsing multiple tenants into one slot. """ if not account_id or not account_id.strip() or account_id == "<unset>": raise ValueError( f"account_id must be a non-empty, non-default string; " f"got {account_id!r}. AccountStore.resolve must always " "return Account(id=<non-empty>) so cross-tenant cache " "scoping works correctly." ) if (webhook_url is None) != (webhook_operation_id is None): raise ValueError("webhook_url and webhook_operation_id must be supplied together") if webhook_url is not None and not webhook_url: raise ValueError("webhook_url must be non-empty when supplied") if webhook_operation_id is not None and not webhook_operation_id: raise ValueError("webhook_operation_id must be non-empty when supplied") if webhook_url is None and webhook_signing_scope_id is not None: raise ValueError("webhook_signing_scope_id requires webhook_url") if webhook_url is None and webhook_authentication is not None: raise ValueError("webhook_authentication requires webhook_url") if webhook_authentication is not None and webhook_signing_scope_id is not None: raise ValueError( "legacy webhook_authentication must not carry an RFC 9421 signing scope" ) outbox = self.task_webhook_outbox if webhook_url is not None: if outbox is None: raise ValueError("webhook registration requires the registry's task_webhook_outbox") outbox.validate_registration(webhook_url, webhook_authentication) task_id = f"task_{uuid.uuid4().hex[:16]}" encrypted_registration: bytes | None = None registration_nonce: bytes | None = None if webhook_url is not None: if webhook_operation_id is None or outbox is None: raise RuntimeError("validated webhook registration became incomplete") encrypted_registration, registration_nonce = outbox.protect_registration( account_id=account_id, task_id=task_id, task_type=task_type, url=webhook_url, operation_id=webhook_operation_id, token=webhook_token, authentication=webhook_authentication, signing_scope_id=webhook_signing_scope_id, ) now = time.time() async with self._pool.connection() as conn: await conn.execute( self._sql_insert, ( task_id, account_id, task_type, json.dumps(request_context) if request_context is not None else None, encrypted_registration, registration_nonce, now, now, ), ) self._notify_lifecycle_observers( "submitted", { "task_id": task_id, "account_id": account_id, "task_type": task_type, "created_at": now, "updated_at": now, }, ) return task_idAllocate a task_id, persist a
submittedrow, return the id.Mirrors :meth:
~adcp.decisioning.InMemoryTaskRegistry.issueincluding the account_id validation guard — empty or sentinel account_ids would allow cross-tenant task-id probing via theWHERE account_id = %spredicate collapsing multiple tenants into one slot. 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]: from adcp.decisioning.task_queries import list_task_records # Account isolation is a SQL predicate, before rows are materialized. async with self._pool.connection() as conn: cur = await conn.execute( f"SELECT task_id, account_id, state, task_type, progress, result, error," # nosec B608 — table identifier is constructor-validated f" request_context, created_at, updated_at," f" (webhook_registration IS NOT NULL) FROM {self._table} WHERE account_id = %s", (account_id,), ) rows = await cur.fetchall() records = [ dict( zip( ( "task_id", "account_id", "state", "task_type", "progress", "result", "error", "context", "created_at", "updated_at", "has_webhook", ), row, ) ) for row in rows ] return list_task_records( records, 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: """Derive an opaque signing scope from trusted framework context. The callback is operator wiring and receives the hydrated :class:`RequestContext`. It must use internal tenant/platform metadata, never buyer request ``context``, ``push_notification_config``, or an unqualified buyer account id. """ resolver = self._webhook_signing_scope_resolver if resolver is None: return None from adcp.decisioning.types import AdcpError try: value: object = resolver(context) if inspect.isawaitable(value): value = await value except Exception: raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolution failed", recovery="terminal", ) from None if not isinstance(value, str): raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolver returned an invalid value", recovery="terminal", ) # Reuse the outbox's bounded opaque-ID validation before any task row # is issued. This is trusted server state, but it still crosses a DB # and authenticated-envelope boundary. if self.task_webhook_outbox is None: raise RuntimeError("signing scope resolver requires a task webhook outbox") self.task_webhook_outbox._validate_signing_scope_id(value) return valueDerive an opaque signing scope from trusted framework context.
The callback is operator wiring and receives the hydrated :class:
RequestContext. It must use internal tenant/platform metadata, never buyer requestcontext,push_notification_config, or an unqualified buyer account id. 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: """Write a progress payload; transition ``submitted`` → ``working``. Silently no-ops when the task is already in a terminal state or unknown — the dispatch wrapper expects this method never to raise on transient conditions (see :class:`~adcp.decisioning.TaskRegistry` docstring). The ``state NOT IN ('completed', 'failed')`` predicate is evaluated server-side so a concurrent terminal write cannot be overwritten by a straggler progress event. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_update_progress, (task_id, json.dumps(progress), time.time()), ) row = await cur.fetchone() if row is not None and row[0] == "submitted": self._notify_lifecycle_observers("working", self._lifecycle_record(row[1:]))Write a progress payload; transition
submitted→working.Silently no-ops when the task is already in a terminal state or unknown — the dispatch wrapper expects this method never to raise on transient conditions (see :class:
~adcp.decisioning.TaskRegistrydocstring).The
state NOT IN ('completed', 'failed')predicate is evaluated server-side so a concurrent terminal write cannot be overwritten by a straggler progress event.
class PostgresTaskRegistry (*,
pool: AsyncConnectionPool,
task_webhook_outbox: PgTaskWebhookOutbox | None = None,
webhook_signing_scope_resolver: WebhookSigningScopeResolver | None = None)-
Expand source code
class PgTaskRegistry(_TaskLifecycleObservers): """PostgreSQL-backed :class:`~adcp.decisioning.TaskRegistry` — v6.1. Durable counterpart to :class:`~adcp.decisioning.InMemoryTaskRegistry`. Set ``is_durable = True`` so the production-mode gate in :func:`adcp.decisioning.serve.create_adcp_server_from_platform` accepts it without requiring ``ADCP_DECISIONING_ALLOW_INMEMORY_TASKS=1``. Parameters ---------- pool: An :class:`psycopg_pool.AsyncConnectionPool` owned by the caller. Each registry operation acquires a short-lived connection from the pool and returns it immediately after the query. No long-lived transactions, no cross-operation state. task_webhook_outbox: Optional :class:`PgTaskWebhookOutbox` sharing this pool. When a task carries callback registration, ``complete`` / ``fail`` enqueue its immutable terminal webhook on the same transaction connection. Notes ----- Unlike :class:`~adcp.signing.PgReplayStore`, this class uses a fixed ``decisioning_tasks`` table name. Multi-tenant table-name isolation is not supported in this release — callers requiring strict schema separation should use separate databases or schemas. """ is_durable: ClassVar[bool] = True def __init__( self, *, pool: AsyncConnectionPool, task_webhook_outbox: PgTaskWebhookOutbox | None = None, webhook_signing_scope_resolver: WebhookSigningScopeResolver | None = None, _table: str = _DEFAULT_TABLE, ) -> None: if not PG_AVAILABLE: raise ImportError(_INSTALL_HINT) if not _SAFE_IDENTIFIER_RE.fullmatch(_table): raise ValueError(f"_table must match [a-z_][a-z0-9_]* (ASCII only), got {_table!r}") if task_webhook_outbox is not None and task_webhook_outbox._pool is not pool: raise ValueError( "PgTaskRegistry and PgTaskWebhookOutbox must use the same connection pool" ) uses_sender_resolver = ( task_webhook_outbox is not None and task_webhook_outbox._sender_resolver is not None ) if uses_sender_resolver != (webhook_signing_scope_resolver is not None): raise ValueError( "webhook_signing_scope_resolver is required exactly when the " "PgTaskWebhookOutbox uses sender_resolver" ) self._init_lifecycle_observers() self._pool = pool self._table = _table self.task_webhook_outbox = task_webhook_outbox self._webhook_signing_scope_resolver = webhook_signing_scope_resolver self.atomic_task_webhook_outbox = task_webhook_outbox is not None # Pre-format queries at construction so the hot path avoids f-strings per call. # _table is whitelisted by _SAFE_IDENTIFIER_RE above. self._sql_insert = ( # noqa: S608 — table name is whitelisted f"INSERT INTO {self._table}" # nosec B608 — table identifier is constructor-validated f" (task_id, account_id, state, task_type, request_context," f" webhook_registration, webhook_registration_nonce, created_at, updated_at)" f" VALUES (%s, %s, 'submitted', %s, %s::jsonb, %s, %s, %s, %s)" ) self._sql_update_progress = ( # noqa: S608 f"WITH previous AS (SELECT task_id, state FROM {self._table}" # nosec B608 — table identifier is constructor-validated f" WHERE task_id = %s AND state NOT IN ('completed', 'failed') FOR UPDATE)" f" UPDATE {self._table} AS task" f" SET state = CASE task.state WHEN 'submitted' THEN 'working' ELSE task.state END," f" progress = %s::jsonb, updated_at = %s" f" FROM previous WHERE task.task_id = previous.task_id" f" RETURNING previous.state, task.task_id, task.account_id, task.task_type," f" task.created_at, task.updated_at" ) self._sql_complete = ( # noqa: S608 f"UPDATE {self._table}" # nosec B608 — table identifier is constructor-validated f" SET state = 'completed', result = %s::jsonb, updated_at = %s" f" WHERE task_id = %s AND state NOT IN ('completed', 'failed')" f" RETURNING task_id, account_id, task_type, webhook_registration," f" webhook_registration_nonce, created_at, updated_at" ) self._sql_fail = ( # noqa: S608 f"UPDATE {self._table}" # nosec B608 — table identifier is constructor-validated f" SET state = 'failed', error = %s::jsonb, updated_at = %s" f" WHERE task_id = %s AND state NOT IN ('completed', 'failed')" f" RETURNING task_id, account_id, task_type, webhook_registration," f" webhook_registration_nonce, created_at, updated_at" ) self._sql_clear_webhook_registration = ( # noqa: S608 f"UPDATE {self._table} SET webhook_registration = NULL," # nosec B608 — table identifier is constructor-validated f" webhook_registration_nonce = NULL WHERE task_id = %s" ) # Explicit ``::text`` cast on the optional account-filter # parameter so psycopg's bind-param type inference doesn't # fail with ``IndeterminateDatatype: could not determine # data type of parameter $2``. Without the cast, the # ``%s IS NULL`` predicate gives psycopg no type context # for the parameter and the query fails at prepare time. self._sql_get = ( # noqa: S608 f"SELECT task_id, account_id, state, task_type," # nosec B608 — table identifier is constructor-validated f" progress, result, error, request_context, created_at, updated_at" f" FROM {self._table}" f" WHERE task_id = %s AND (%s::text IS NULL OR account_id = %s)" ) self._sql_get_state_result = ( # noqa: S608 f"SELECT state, result, task_type FROM {self._table}" # nosec B608 — validated identifier " WHERE task_id = %s" ) self._sql_get_task_type = f"SELECT task_type FROM {self._table} WHERE task_id = %s" # noqa: S608 # nosec B608 — table identifier is constructor-validated self._sql_get_state_error = f"SELECT state, error FROM {self._table} WHERE task_id = %s" # noqa: S608 # nosec B608 — table identifier is constructor-validated self._sql_discard = ( # noqa: S608 f"DELETE FROM {self._table} WHERE task_id = %s" # nosec B608 — table identifier is constructor-validated f" RETURNING task_id, account_id, task_type, created_at, updated_at" ) self._sql_ddl = ( # noqa: S608 f"CREATE TABLE IF NOT EXISTS {self._table} (" f' task_id TEXT COLLATE "C" NOT NULL PRIMARY KEY,' f' account_id TEXT COLLATE "C" NOT NULL,' f" state TEXT NOT NULL DEFAULT 'submitted'," f" task_type TEXT NOT NULL," f" progress JSONB," f" result JSONB," f" error JSONB," f" request_context JSONB," f" webhook_registration BYTEA," f" webhook_registration_nonce BYTEA," f" created_at DOUBLE PRECISION NOT NULL," f" updated_at DOUBLE PRECISION NOT NULL" f");" f"ALTER TABLE {self._table} ADD COLUMN IF NOT EXISTS request_context JSONB;" f"ALTER TABLE {self._table} ADD COLUMN IF NOT EXISTS webhook_registration BYTEA;" f"ALTER TABLE {self._table} ADD COLUMN IF NOT EXISTS webhook_registration_nonce BYTEA;" f"CREATE INDEX IF NOT EXISTS {self._table}_account_idx" # noqa: S608 f" ON {self._table} (account_id);" ) # -- schema bootstrap ----------------------------------------------- async def create_schema(self) -> None: """Create the task registry table and supporting index. Honors the ``_table`` kwarg the store was constructed with. Idempotent via ``CREATE TABLE IF NOT EXISTS`` — safe to call on every application boot. The equivalent raw DDL ships at ``adcp/decisioning/pg/decisioning_tasks.sql`` in the installed package for adopters using a migration tool (Alembic, Flyway, psql). """ async with self._pool.connection() as conn: await conn.execute(self._sql_ddl) # -- TaskRegistry Protocol ------------------------------------------ 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: """Allocate a task_id, persist a ``submitted`` row, return the id. Mirrors :meth:`~adcp.decisioning.InMemoryTaskRegistry.issue` including the account_id validation guard — empty or sentinel account_ids would allow cross-tenant task-id probing via the ``WHERE account_id = %s`` predicate collapsing multiple tenants into one slot. """ if not account_id or not account_id.strip() or account_id == "<unset>": raise ValueError( f"account_id must be a non-empty, non-default string; " f"got {account_id!r}. AccountStore.resolve must always " "return Account(id=<non-empty>) so cross-tenant cache " "scoping works correctly." ) if (webhook_url is None) != (webhook_operation_id is None): raise ValueError("webhook_url and webhook_operation_id must be supplied together") if webhook_url is not None and not webhook_url: raise ValueError("webhook_url must be non-empty when supplied") if webhook_operation_id is not None and not webhook_operation_id: raise ValueError("webhook_operation_id must be non-empty when supplied") if webhook_url is None and webhook_signing_scope_id is not None: raise ValueError("webhook_signing_scope_id requires webhook_url") if webhook_url is None and webhook_authentication is not None: raise ValueError("webhook_authentication requires webhook_url") if webhook_authentication is not None and webhook_signing_scope_id is not None: raise ValueError( "legacy webhook_authentication must not carry an RFC 9421 signing scope" ) outbox = self.task_webhook_outbox if webhook_url is not None: if outbox is None: raise ValueError("webhook registration requires the registry's task_webhook_outbox") outbox.validate_registration(webhook_url, webhook_authentication) task_id = f"task_{uuid.uuid4().hex[:16]}" encrypted_registration: bytes | None = None registration_nonce: bytes | None = None if webhook_url is not None: if webhook_operation_id is None or outbox is None: raise RuntimeError("validated webhook registration became incomplete") encrypted_registration, registration_nonce = outbox.protect_registration( account_id=account_id, task_id=task_id, task_type=task_type, url=webhook_url, operation_id=webhook_operation_id, token=webhook_token, authentication=webhook_authentication, signing_scope_id=webhook_signing_scope_id, ) now = time.time() async with self._pool.connection() as conn: await conn.execute( self._sql_insert, ( task_id, account_id, task_type, json.dumps(request_context) if request_context is not None else None, encrypted_registration, registration_nonce, now, now, ), ) self._notify_lifecycle_observers( "submitted", { "task_id": task_id, "account_id": account_id, "task_type": task_type, "created_at": now, "updated_at": now, }, ) return task_id async def resolve_webhook_signing_scope( self, context: RequestContext[Any], ) -> str | None: """Derive an opaque signing scope from trusted framework context. The callback is operator wiring and receives the hydrated :class:`RequestContext`. It must use internal tenant/platform metadata, never buyer request ``context``, ``push_notification_config``, or an unqualified buyer account id. """ resolver = self._webhook_signing_scope_resolver if resolver is None: return None from adcp.decisioning.types import AdcpError try: value: object = resolver(context) if inspect.isawaitable(value): value = await value except Exception: raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolution failed", recovery="terminal", ) from None if not isinstance(value, str): raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolver returned an invalid value", recovery="terminal", ) # Reuse the outbox's bounded opaque-ID validation before any task row # is issued. This is trusted server state, but it still crosses a DB # and authenticated-envelope boundary. if self.task_webhook_outbox is None: raise RuntimeError("signing scope resolver requires a task webhook outbox") self.task_webhook_outbox._validate_signing_scope_id(value) return value async def update_progress( self, task_id: str, progress: dict[str, Any], ) -> None: """Write a progress payload; transition ``submitted`` → ``working``. Silently no-ops when the task is already in a terminal state or unknown — the dispatch wrapper expects this method never to raise on transient conditions (see :class:`~adcp.decisioning.TaskRegistry` docstring). The ``state NOT IN ('completed', 'failed')`` predicate is evaluated server-side so a concurrent terminal write cannot be overwritten by a straggler progress event. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_update_progress, (task_id, json.dumps(progress), time.time()), ) row = await cur.fetchone() if row is not None and row[0] == "submitted": self._notify_lifecycle_observers("working", self._lifecycle_record(row[1:])) @staticmethod def _lifecycle_record(row: tuple[Any, ...]) -> dict[str, Any]: return dict(zip(("task_id", "account_id", "task_type", "created_at", "updated_at"), row)) async def complete(self, task_id: str, result: dict[str, Any]) -> None: """Complete once; notify metrics only after the transaction commits.""" row = await self._complete(task_id, result) if row is not None: self._notify_lifecycle_observers("completed", self._lifecycle_record(row[:3] + row[5:])) async def _complete( self, task_id: str, result: dict[str, Any], ) -> tuple[Any, ...] | None: """Mark the task ``completed`` with ``result`` as the terminal artifact. Idempotent on repeated calls with an equal ``result``; raises :class:`ValueError` on conflicting re-completion. Uses an atomic ``UPDATE ... RETURNING`` so concurrent workers cannot race each other into double-completion without detection. """ async with self._pool.connection() as conn: # Do not rely on the pool connection's implicit transaction. # Explicitly bind the terminal transition, outbox insert, and # callback-registration clear even when adopters configure the # pool with autocommit=True. async with conn.transaction(): type_cursor = await conn.execute(self._sql_get_task_type, (task_id,)) type_row = await type_cursor.fetchone() safe_result = ( strip_credentials_from_wire_result(type_row[0], result) if type_row is not None else result ) cur = await conn.execute( self._sql_complete, (json.dumps(safe_result), time.time(), task_id), ) row: tuple[Any, ...] | None = await cur.fetchone() if row is not None: await self._enqueue_terminal_if_registered( conn, row=row, status="completed", payload=safe_result, ) return row # connection and transaction exit before the caller notifies # Zero rows in RETURNING — task is unknown or already terminal. cur2 = await conn.execute(self._sql_get_state_result, (task_id,)) row = await cur2.fetchone() if row is None: raise ValueError(f"Task {task_id!r} not found") state, existing_result, task_type = row safe_result = strip_credentials_from_wire_result(task_type, result) if state == "completed": if existing_result == safe_result: return None # idempotent raise ValueError(f"Task {task_id!r} already completed with a different result") raise ValueError(f"Task {task_id!r} already in terminal state {state!r}") async def fail(self, task_id: str, error: dict[str, Any]) -> None: """Fail once; notify metrics only after the transaction commits.""" row = await self._fail(task_id, error) if row is not None: self._notify_lifecycle_observers("failed", self._lifecycle_record(row[:3] + row[5:])) async def _fail( self, task_id: str, error: dict[str, Any], ) -> tuple[Any, ...] | None: """Mark the task ``failed`` with ``error`` as the terminal payload. Idempotent on repeated calls with an equal ``error``; raises :class:`ValueError` on conflicting re-failure. """ async with self._pool.connection() as conn: async with conn.transaction(): cur = await conn.execute( self._sql_fail, (json.dumps(error), time.time(), task_id), ) row: tuple[Any, ...] | None = await cur.fetchone() if row is not None: await self._enqueue_terminal_if_registered( conn, row=row, status="failed", payload=error, ) return row # notify only after connection/transaction exit # Zero rows in RETURNING — task is unknown or already terminal. cur2 = await conn.execute(self._sql_get_state_error, (task_id,)) row = await cur2.fetchone() if row is None: raise ValueError(f"Task {task_id!r} not found") state, existing_error = row if state == "failed": if existing_error == error: return None # idempotent raise ValueError(f"Task {task_id!r} already failed with a different error") raise ValueError(f"Task {task_id!r} already in terminal state {state!r}") async def get( self, task_id: str, *, expected_account_id: str | None = None, ) -> dict[str, Any] | None: """Look up a task record; cross-tenant probes return ``None``. The ``expected_account_id`` predicate is enforced at the SQL level (``WHERE account_id = %s``), not as a Python-level filter after fetch. This guarantees the row is never materialized for a mismatched probe, eliminating the fetch-then-filter anti-pattern. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_get, (task_id, expected_account_id, expected_account_id) ) row = await cur.fetchone() if row is None: return None return { "task_id": row[0], "account_id": row[1], "state": row[2], "task_type": row[3], "progress": row[4], "result": row[5], "error": row[6], "created_at": row[8], "updated_at": row[9], **({"context": row[7]} if row[7] is not None else {}), } 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]: from adcp.decisioning.task_queries import list_task_records # Account isolation is a SQL predicate, before rows are materialized. async with self._pool.connection() as conn: cur = await conn.execute( f"SELECT task_id, account_id, state, task_type, progress, result, error," # nosec B608 — table identifier is constructor-validated f" request_context, created_at, updated_at," f" (webhook_registration IS NOT NULL) FROM {self._table} WHERE account_id = %s", (account_id,), ) rows = await cur.fetchall() records = [ dict( zip( ( "task_id", "account_id", "state", "task_type", "progress", "result", "error", "context", "created_at", "updated_at", "has_webhook", ), row, ) ) for row in rows ] return list_task_records( records, account_id=account_id, filters=filters, sort=sort, pagination=pagination ) async def _enqueue_terminal_if_registered( self, conn: Any, *, row: tuple[Any, ...], status: str, payload: dict[str, Any], ) -> None: """Enqueue on ``conn`` so task state and webhook commit atomically.""" task_id, account_id, task_type, encrypted_registration, registration_nonce = row[:5] if encrypted_registration is None: return if self.task_webhook_outbox is None: raise RuntimeError( "Task carries push_notification_config but PgTaskRegistry has no " "task_webhook_outbox; refusing a non-atomic terminal transition" ) if registration_nonce is None: raise RuntimeError(f"Task {task_id!r} has incomplete webhook registration") url, operation_id, token, authentication, signing_scope_id = ( self.task_webhook_outbox._open_registration_with_scope( account_id=account_id, task_id=task_id, task_type=task_type, encrypted_registration=bytes(encrypted_registration), nonce=bytes(registration_nonce), ) ) await self.task_webhook_outbox.enqueue_terminal( conn, task_id=task_id, account_id=account_id, task_type=task_type, status=status, result=payload, url=url, operation_id=operation_id, token=token, authentication=authentication, signing_scope_id=signing_scope_id, ) # The encrypted outbox envelope now owns the callback registration. # Clear the task-row copy in this same transaction. await conn.execute(self._sql_clear_webhook_registration, (task_id,)) async def discard(self, task_id: str) -> None: """Remove a task_id from the registry — rollback path. Idempotent: discarding an unknown task_id is a no-op (no raise), matching the :class:`~adcp.decisioning.InMemoryTaskRegistry` contract. """ async with self._pool.connection() as conn: cur = await conn.execute(self._sql_discard, (task_id,)) row = await cur.fetchone() if row is not None: event_record = self._lifecycle_record(row) event_record["updated_at"] = time.time() self._notify_lifecycle_observers("discarded", event_record)PostgreSQL-backed :class:
~adcp.decisioning.TaskRegistry— v6.1.Durable counterpart to :class:
~adcp.decisioning.InMemoryTaskRegistry. Setis_durable = Trueso the production-mode gate in :func:adcp.decisioning.serve.create_adcp_server_from_platformaccepts it without requiringADCP_DECISIONING_ALLOW_INMEMORY_TASKS=1.Parameters
pool: An :class:
psycopg_pool.AsyncConnectionPoolowned by the caller. Each registry operation acquires a short-lived connection from the pool and returns it immediately after the query. No long-lived transactions, no cross-operation state. task_webhook_outbox: Optional :class:PgTaskWebhookOutboxsharing this pool. When a task carries callback registration,complete/failenqueue its immutable terminal webhook on the same transaction connection.Notes
Unlike :class:
~adcp.signing.PgReplayStore, this class uses a fixeddecisioning_taskstable name. Multi-tenant table-name isolation is not supported in this release — callers requiring strict schema separation should use separate databases or schemas.Ancestors
- adcp.decisioning.task_registry._TaskLifecycleObservers
Class variables
var is_durable : ClassVar[bool]
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: """Complete once; notify metrics only after the transaction commits.""" row = await self._complete(task_id, result) if row is not None: self._notify_lifecycle_observers("completed", self._lifecycle_record(row[:3] + row[5:]))Complete once; notify metrics only after the transaction commits.
async def create_schema(self) ‑> None-
Expand source code
async def create_schema(self) -> None: """Create the task registry table and supporting index. Honors the ``_table`` kwarg the store was constructed with. Idempotent via ``CREATE TABLE IF NOT EXISTS`` — safe to call on every application boot. The equivalent raw DDL ships at ``adcp/decisioning/pg/decisioning_tasks.sql`` in the installed package for adopters using a migration tool (Alembic, Flyway, psql). """ async with self._pool.connection() as conn: await conn.execute(self._sql_ddl)Create the task registry table and supporting index.
Honors the
_tablekwarg the store was constructed with. Idempotent viaCREATE TABLE IF NOT EXISTS— safe to call on every application boot. The equivalent raw DDL ships atadcp/decisioning/pg/decisioning_tasks.sqlin the installed package for adopters using a migration tool (Alembic, Flyway, psql). async def discard(self, task_id: str) ‑> None-
Expand source code
async def discard(self, task_id: str) -> None: """Remove a task_id from the registry — rollback path. Idempotent: discarding an unknown task_id is a no-op (no raise), matching the :class:`~adcp.decisioning.InMemoryTaskRegistry` contract. """ async with self._pool.connection() as conn: cur = await conn.execute(self._sql_discard, (task_id,)) row = await cur.fetchone() if row is not None: event_record = self._lifecycle_record(row) event_record["updated_at"] = time.time() self._notify_lifecycle_observers("discarded", event_record)Remove a task_id from the registry — rollback path.
Idempotent: discarding an unknown task_id is a no-op (no raise), matching the :class:
~adcp.decisioning.InMemoryTaskRegistrycontract. 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: """Fail once; notify metrics only after the transaction commits.""" row = await self._fail(task_id, error) if row is not None: self._notify_lifecycle_observers("failed", self._lifecycle_record(row[:3] + row[5:]))Fail once; notify metrics only after the transaction commits.
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: """Look up a task record; cross-tenant probes return ``None``. The ``expected_account_id`` predicate is enforced at the SQL level (``WHERE account_id = %s``), not as a Python-level filter after fetch. This guarantees the row is never materialized for a mismatched probe, eliminating the fetch-then-filter anti-pattern. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_get, (task_id, expected_account_id, expected_account_id) ) row = await cur.fetchone() if row is None: return None return { "task_id": row[0], "account_id": row[1], "state": row[2], "task_type": row[3], "progress": row[4], "result": row[5], "error": row[6], "created_at": row[8], "updated_at": row[9], **({"context": row[7]} if row[7] is not None else {}), }Look up a task record; cross-tenant probes return
None.The
expected_account_idpredicate is enforced at the SQL level (WHERE account_id = %s), not as a Python-level filter after fetch. This guarantees the row is never materialized for a mismatched probe, eliminating the fetch-then-filter anti-pattern. 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: """Allocate a task_id, persist a ``submitted`` row, return the id. Mirrors :meth:`~adcp.decisioning.InMemoryTaskRegistry.issue` including the account_id validation guard — empty or sentinel account_ids would allow cross-tenant task-id probing via the ``WHERE account_id = %s`` predicate collapsing multiple tenants into one slot. """ if not account_id or not account_id.strip() or account_id == "<unset>": raise ValueError( f"account_id must be a non-empty, non-default string; " f"got {account_id!r}. AccountStore.resolve must always " "return Account(id=<non-empty>) so cross-tenant cache " "scoping works correctly." ) if (webhook_url is None) != (webhook_operation_id is None): raise ValueError("webhook_url and webhook_operation_id must be supplied together") if webhook_url is not None and not webhook_url: raise ValueError("webhook_url must be non-empty when supplied") if webhook_operation_id is not None and not webhook_operation_id: raise ValueError("webhook_operation_id must be non-empty when supplied") if webhook_url is None and webhook_signing_scope_id is not None: raise ValueError("webhook_signing_scope_id requires webhook_url") if webhook_url is None and webhook_authentication is not None: raise ValueError("webhook_authentication requires webhook_url") if webhook_authentication is not None and webhook_signing_scope_id is not None: raise ValueError( "legacy webhook_authentication must not carry an RFC 9421 signing scope" ) outbox = self.task_webhook_outbox if webhook_url is not None: if outbox is None: raise ValueError("webhook registration requires the registry's task_webhook_outbox") outbox.validate_registration(webhook_url, webhook_authentication) task_id = f"task_{uuid.uuid4().hex[:16]}" encrypted_registration: bytes | None = None registration_nonce: bytes | None = None if webhook_url is not None: if webhook_operation_id is None or outbox is None: raise RuntimeError("validated webhook registration became incomplete") encrypted_registration, registration_nonce = outbox.protect_registration( account_id=account_id, task_id=task_id, task_type=task_type, url=webhook_url, operation_id=webhook_operation_id, token=webhook_token, authentication=webhook_authentication, signing_scope_id=webhook_signing_scope_id, ) now = time.time() async with self._pool.connection() as conn: await conn.execute( self._sql_insert, ( task_id, account_id, task_type, json.dumps(request_context) if request_context is not None else None, encrypted_registration, registration_nonce, now, now, ), ) self._notify_lifecycle_observers( "submitted", { "task_id": task_id, "account_id": account_id, "task_type": task_type, "created_at": now, "updated_at": now, }, ) return task_idAllocate a task_id, persist a
submittedrow, return the id.Mirrors :meth:
~adcp.decisioning.InMemoryTaskRegistry.issueincluding the account_id validation guard — empty or sentinel account_ids would allow cross-tenant task-id probing via theWHERE account_id = %spredicate collapsing multiple tenants into one slot. 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]: from adcp.decisioning.task_queries import list_task_records # Account isolation is a SQL predicate, before rows are materialized. async with self._pool.connection() as conn: cur = await conn.execute( f"SELECT task_id, account_id, state, task_type, progress, result, error," # nosec B608 — table identifier is constructor-validated f" request_context, created_at, updated_at," f" (webhook_registration IS NOT NULL) FROM {self._table} WHERE account_id = %s", (account_id,), ) rows = await cur.fetchall() records = [ dict( zip( ( "task_id", "account_id", "state", "task_type", "progress", "result", "error", "context", "created_at", "updated_at", "has_webhook", ), row, ) ) for row in rows ] return list_task_records( records, 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: """Derive an opaque signing scope from trusted framework context. The callback is operator wiring and receives the hydrated :class:`RequestContext`. It must use internal tenant/platform metadata, never buyer request ``context``, ``push_notification_config``, or an unqualified buyer account id. """ resolver = self._webhook_signing_scope_resolver if resolver is None: return None from adcp.decisioning.types import AdcpError try: value: object = resolver(context) if inspect.isawaitable(value): value = await value except Exception: raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolution failed", recovery="terminal", ) from None if not isinstance(value, str): raise AdcpError( "INTERNAL_ERROR", message="Webhook signing scope resolver returned an invalid value", recovery="terminal", ) # Reuse the outbox's bounded opaque-ID validation before any task row # is issued. This is trusted server state, but it still crosses a DB # and authenticated-envelope boundary. if self.task_webhook_outbox is None: raise RuntimeError("signing scope resolver requires a task webhook outbox") self.task_webhook_outbox._validate_signing_scope_id(value) return valueDerive an opaque signing scope from trusted framework context.
The callback is operator wiring and receives the hydrated :class:
RequestContext. It must use internal tenant/platform metadata, never buyer requestcontext,push_notification_config, or an unqualified buyer account id. 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: """Write a progress payload; transition ``submitted`` → ``working``. Silently no-ops when the task is already in a terminal state or unknown — the dispatch wrapper expects this method never to raise on transient conditions (see :class:`~adcp.decisioning.TaskRegistry` docstring). The ``state NOT IN ('completed', 'failed')`` predicate is evaluated server-side so a concurrent terminal write cannot be overwritten by a straggler progress event. """ async with self._pool.connection() as conn: cur = await conn.execute( self._sql_update_progress, (task_id, json.dumps(progress), time.time()), ) row = await cur.fetchone() if row is not None and row[0] == "submitted": self._notify_lifecycle_observers("working", self._lifecycle_record(row[1:]))Write a progress payload; transition
submitted→working.Silently no-ops when the task is already in a terminal state or unknown — the dispatch wrapper expects this method never to raise on transient conditions (see :class:
~adcp.decisioning.TaskRegistrydocstring).The
state NOT IN ('completed', 'failed')predicate is evaluated server-side so a concurrent terminal write cannot be overwritten by a straggler progress event.