Module adcp.server.idempotency
Server-side idempotency middleware for AdCP mutating tool handlers.
Implements the seller side of AdCP #2315: extract idempotency_key, look up
cached responses scoped by (authenticated_principal, idempotency_key), and
replay the cached response verbatim when a subsequent request carries the same
key + canonicalized-equivalent payload. Reject key reuse with a different
payload as IDEMPOTENCY_CONFLICT.
The spec contract lives at
adcontextprotocol/adcp/docs/building/implementation/security.mdx#idempotency.
Typical usage::
from adcp.server import ADCPHandler, IdempotencyStore, MemoryBackend, ToolContext
from adcp.server.responses import capabilities_response
idempotency = IdempotencyStore(
backend=MemoryBackend(),
ttl_seconds=86400, # 24 hours, matches spec minimum
)
class MySeller(ADCPHandler):
@idempotency.wrap
async def create_media_buy(self, params, context=None):
# <code>params</code> carries <code>idempotency\_key</code>; <code>context.caller\_identity</code>
# identifies the principal. Without either, the middleware falls
# through to this handler with no dedup (schema validation above
# us rejects missing keys per AdCP #2315).
return my_create_logic(params)
async def get_adcp_capabilities(self, params, context=None):
return capabilities_response(
["media_buy"],
idempotency=idempotency.capability(),
)
Callers who invoke the handler directly (tests, non-HTTP code paths) must
pass a :class:ToolContext with caller_identity set —
it's how the middleware scopes the cache namespace per-principal::
ctx = ToolContext(caller_identity="buyer-acme")
result = await seller.create_media_buy(
{"idempotency_key": key, ...}, ctx
)
Backends:
- :class:
MemoryBackend— in-process dict with TTL; use for tests and single-process reference implementations. - :class:
PgBackend— Postgres-backed store for multi-worker durable replay. Requires theadcp[pg]extra.await backend.create_schema()once at boot; commits go through a fresh pool connection. UseIdempotencyStore.reserve()for an explicit business/replay transaction with commit-before-unlock ordering.
Sub-modules
adcp.server.idempotency.backends-
Storage backends for :class:
~adcp.server.idempotency.IdempotencyStore… adcp.server.idempotency.canonicalize-
Canonical payload hashing for AdCP idempotency replay detection …
adcp.server.idempotency.lazy-
Deferred-construction wrapper for :class:
IdempotencyBackend… adcp.server.idempotency.reservation-
Explicit PostgreSQL reservations with business/replay transaction ownership.
adcp.server.idempotency.store-
The :class:
IdempotencyStorecoordinator: canonical hashing + backend + decorator … adcp.server.idempotency.webhook_dedup-
Webhook receiver-side claim and completion store …
Functions
def canonical_json_sha256(payload: dict[str, Any]) ‑> str-
Expand source code
def canonical_json_sha256(payload: dict[str, Any]) -> str: """Compute the spec's canonical payload fingerprint. 1. Strip the spec's exclusion list (see :func:`strip_excluded_fields`). 2. RFC 8785 JCS canonicalize (stable key order, compact, UTF-8, spec- compliant number serialization). 3. SHA-256 over the canonical bytes; return hex digest. The result is stable across all conforming JCS implementations. Two payloads whose hashes match are equivalent under AdCP replay semantics; two with different hashes are distinct and MUST be treated as a conflict when the caller supplies the same ``idempotency_key``. """ stripped = strip_excluded_fields(payload) canonical = rfc8785.dumps(stripped) return hashlib.sha256(canonical).hexdigest()Compute the spec's canonical payload fingerprint.
- Strip the spec's exclusion list (see :func:
strip_excluded_fields()). - RFC 8785 JCS canonicalize (stable key order, compact, UTF-8, spec- compliant number serialization).
- SHA-256 over the canonical bytes; return hex digest.
The result is stable across all conforming JCS implementations. Two payloads whose hashes match are equivalent under AdCP replay semantics; two with different hashes are distinct and MUST be treated as a conflict when the caller supplies the same
idempotency_key. - Strip the spec's exclusion list (see :func:
def create_lazy_backend(factory: LazyBackendFactory, *, allow_clear_all: bool = False) ‑> LazyBackend-
Expand source code
def create_lazy_backend( factory: LazyBackendFactory, *, allow_clear_all: bool = False, ) -> LazyBackend: """Construct a :class:`LazyBackend` — functional alias mirroring the JS ``createLazyBackend`` factory shape. See :class:`LazyBackend` for semantics and parameter documentation. """ return LazyBackend(factory, allow_clear_all=allow_clear_all)Construct a :class:
LazyBackend— functional alias mirroring the JScreateLazyBackendfactory shape.See :class:
LazyBackendfor semantics and parameter documentation. def is_wrapped(fn: Any) ‑> bool-
Expand source code
def is_wrapped(fn: Any) -> bool: """Return whether a callable's chain contains an SDK-registered wrapper. Accepts bound methods and plain callables. Outer decorators must use ``functools.wraps`` (or set ``__wrapped__``) so the registered wrapper remains reachable. Copied attributes alone do not establish membership. Used by :mod:`adcp.decisioning.validate_idempotency`. """ seen: set[int] = set() # Bound the traversal to guard against pathological __wrapped__ cycles # and chains that manufacture a new callable on every attribute access. for _ in range(16): target = getattr(fn, "__func__", fn) if target is None or id(target) in seen: return False seen.add(id(target)) try: if target in _WRAPPED_FUNCTIONS: return True except TypeError: # Unhashable/non-weak-referenceable objects cannot be registered. pass fn = getattr(target, "__wrapped__", None) return FalseReturn whether a callable's chain contains an SDK-registered wrapper.
Accepts bound methods and plain callables. Outer decorators must use
functools.wraps(or set__wrapped__) so the registered wrapper remains reachable. Copied attributes alone do not establish membership. Used by :mod:adcp.decisioning.validate_idempotency. def strip_excluded_fields(payload: dict[str, Any]) ‑> dict[str, typing.Any]-
Expand source code
def strip_excluded_fields(payload: dict[str, Any]) -> dict[str, Any]: """Return a deep copy of ``payload`` with the spec's exclusion list removed. Top-level keys in :data:`EXCLUDED_FIELDS` are dropped. Nested paths in ``_NESTED_EXCLUSIONS`` are traversed; the final leaf key is removed if the traversal reaches it. Missing intermediate keys are a no-op — the caller's payload is free to omit the push_notification_config entirely. The input dict is never mutated. """ out: dict[str, Any] = copy.deepcopy(payload) for key in EXCLUDED_FIELDS: out.pop(key, None) for path in _NESTED_EXCLUSIONS: _drop_nested(out, path) return outReturn a deep copy of
payloadwith the spec's exclusion list removed.Top-level keys in :data:
EXCLUDED_FIELDSare dropped. Nested paths in_NESTED_EXCLUSIONSare traversed; the final leaf key is removed if the traversal reaches it. Missing intermediate keys are a no-op — the caller's payload is free to omit the push_notification_config entirely.The input dict is never mutated.
Classes
class CachedResponse (payload_hash: str, response: dict[str, Any], expires_at_epoch: float)-
Expand source code
@dataclass(frozen=True) class CachedResponse: """A single cached handler response keyed by ``(principal_id, key)``. :param payload_hash: Canonical JSON SHA-256 of the *original* request. On replay we compare the new request's hash to this value; mismatch is ``IDEMPOTENCY_CONFLICT``. :param response: The response dict the handler returned. On replay, :meth:`IdempotencyStore.wrap` injects ``replayed: true`` at the envelope level per AdCP L1/security idempotency rule 4 before returning to the caller — the cached value here stays clean so the same entry can serve multiple replays without compounding. :param expires_at_epoch: Unix timestamp (seconds) when this entry becomes eligible for eviction. Reads after this time return None. """ payload_hash: str response: dict[str, Any] expires_at_epoch: floatA single cached handler response keyed by
(principal_id, key).:param payload_hash: Canonical JSON SHA-256 of the original request. On replay we compare the new request's hash to this value; mismatch is
IDEMPOTENCY_CONFLICT. :param response: The response dict the handler returned. On replay, :meth:IdempotencyStore.wrap()injectsreplayed: trueat the envelope level per AdCP L1/security idempotency rule 4 before returning to the caller — the cached value here stays clean so the same entry can serve multiple replays without compounding. :param expires_at_epoch: Unix timestamp (seconds) when this entry becomes eligible for eviction. Reads after this time return None.Instance variables
var expires_at_epoch : floatvar payload_hash : strvar response : dict[str, typing.Any]
class IdempotencyBackend-
Expand source code
class IdempotencyBackend(ABC): """Abstract storage backend contract. All methods are async. Implementations MUST be safe to call concurrently from multiple asyncio tasks. ``hold`` must coordinate all processes that share the backend namespace. """ @abstractmethod async def get(self, scope_key: str, key: str) -> CachedResponse | None: """Return the cached entry, or None if missing or expired. ``scope_key`` is the caller-composed identity scope — typically ``tenant_id + caller_identity``. Backends treat it as an opaque string; the composition is owned by :class:`~adcp.server.idempotency.IdempotencyStore`. """ @abstractmethod async def put( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Store ``entry`` under ``(scope_key, key)``. Overwrites any prior entry — the store only calls ``put`` after verifying the slot is empty or expired, so an overwrite in that window is a legitimate retry of the write itself.""" async def replace( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Replace a live entry while the caller holds this key's lock. Request idempotency uses first-writer-wins :meth:`put` semantics. Stateful protocols such as webhook delivery claims also need an owner-fenced live-state transition. Callers MUST invoke ``replace`` only inside :meth:`hold` after validating the current entry. The default delegates to ``put`` for legacy/custom backends whose put is a true overwrite; backends with first-writer-wins put semantics must override it. """ await self.put(scope_key, key, entry) async def current_time(self) -> float: """Return the clock authoritative for this backend's expiry checks.""" return time.time() def hold(self, scope_key: str, key: str) -> Any: """Return an async context manager holding the key's execution lock. Custom backends must implement this operation to be usable by :class:`IdempotencyStore`. It is concrete only to keep existing backend subclasses importable while adopters migrate. """ raise NotImplementedError( f"{type(self).__name__} must implement atomic hold(scope_key, key)" ) async def supports_atomic_hold(self) -> bool: """Return whether ``hold`` provides the backend's atomic guarantee. This explicit capability is delegable by wrappers such as :class:`LazyBackend`; method-identity inspection is not. """ return type(self).hold is not IdempotencyBackend.hold async def supports_atomic_put_if_absent(self) -> bool: """Return whether ``put_if_absent`` is implemented atomically.""" return ( type(self).put_if_absent is not IdempotencyBackend.put_if_absent or await self.supports_atomic_hold() ) async def put_if_absent(self, scope_key: str, key: str, entry: CachedResponse) -> bool: """Atomically insert a fresh/expired slot; return whether it won. The default composes ``hold`` with the legacy ``get``/``put`` API, allowing custom backends to implement one locking primitive. Backends with a native conditional insert should override this method. """ async with self.hold(scope_key, key): if await self.get(scope_key, key) is not None: return False await self.put(scope_key, key, entry) return True @abstractmethod async def delete_expired(self, now_epoch: float | None = None) -> int: """Best-effort sweep of expired entries. Returns the count removed. Sweeping is optional — :meth:`get` MUST self-filter expired entries. Backends that have natural TTL primitives (Redis ``EXPIRE``, Postgres partial indexes) may implement this as a no-op."""Abstract storage backend contract.
All methods are async. Implementations MUST be safe to call concurrently from multiple asyncio tasks.
holdmust coordinate all processes that share the backend namespace.Ancestors
- abc.ABC
Subclasses
Methods
async def current_time(self) ‑> float-
Expand source code
async def current_time(self) -> float: """Return the clock authoritative for this backend's expiry checks.""" return time.time()Return the clock authoritative for this backend's expiry checks.
async def delete_expired(self, now_epoch: float | None = None) ‑> int-
Expand source code
@abstractmethod async def delete_expired(self, now_epoch: float | None = None) -> int: """Best-effort sweep of expired entries. Returns the count removed. Sweeping is optional — :meth:`get` MUST self-filter expired entries. Backends that have natural TTL primitives (Redis ``EXPIRE``, Postgres partial indexes) may implement this as a no-op."""Best-effort sweep of expired entries. Returns the count removed.
Sweeping is optional — :meth:
getMUST self-filter expired entries. Backends that have natural TTL primitives (RedisEXPIRE, Postgres partial indexes) may implement this as a no-op. async def get(self, scope_key: str, key: str) ‑> CachedResponse | None-
Expand source code
@abstractmethod async def get(self, scope_key: str, key: str) -> CachedResponse | None: """Return the cached entry, or None if missing or expired. ``scope_key`` is the caller-composed identity scope — typically ``tenant_id + caller_identity``. Backends treat it as an opaque string; the composition is owned by :class:`~adcp.server.idempotency.IdempotencyStore`. """Return the cached entry, or None if missing or expired.
scope_keyis the caller-composed identity scope — typicallytenant_id + caller_identity. Backends treat it as an opaque string; the composition is owned by :class:~adcp.server.idempotency.IdempotencyStore. def hold(self, scope_key: str, key: str) ‑> Any-
Expand source code
def hold(self, scope_key: str, key: str) -> Any: """Return an async context manager holding the key's execution lock. Custom backends must implement this operation to be usable by :class:`IdempotencyStore`. It is concrete only to keep existing backend subclasses importable while adopters migrate. """ raise NotImplementedError( f"{type(self).__name__} must implement atomic hold(scope_key, key)" )Return an async context manager holding the key's execution lock.
Custom backends must implement this operation to be usable by :class:
IdempotencyStore. It is concrete only to keep existing backend subclasses importable while adopters migrate. async def put(self,
scope_key: str,
key: str,
entry: CachedResponse) ‑> None-
Expand source code
@abstractmethod async def put( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Store ``entry`` under ``(scope_key, key)``. Overwrites any prior entry — the store only calls ``put`` after verifying the slot is empty or expired, so an overwrite in that window is a legitimate retry of the write itself."""Store
entryunder(scope_key, key). Overwrites any prior entry — the store only callsputafter verifying the slot is empty or expired, so an overwrite in that window is a legitimate retry of the write itself. async def put_if_absent(self,
scope_key: str,
key: str,
entry: CachedResponse) ‑> bool-
Expand source code
async def put_if_absent(self, scope_key: str, key: str, entry: CachedResponse) -> bool: """Atomically insert a fresh/expired slot; return whether it won. The default composes ``hold`` with the legacy ``get``/``put`` API, allowing custom backends to implement one locking primitive. Backends with a native conditional insert should override this method. """ async with self.hold(scope_key, key): if await self.get(scope_key, key) is not None: return False await self.put(scope_key, key, entry) return TrueAtomically insert a fresh/expired slot; return whether it won.
The default composes
holdwith the legacyget/putAPI, allowing custom backends to implement one locking primitive. Backends with a native conditional insert should override this method. async def replace(self,
scope_key: str,
key: str,
entry: CachedResponse) ‑> None-
Expand source code
async def replace( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Replace a live entry while the caller holds this key's lock. Request idempotency uses first-writer-wins :meth:`put` semantics. Stateful protocols such as webhook delivery claims also need an owner-fenced live-state transition. Callers MUST invoke ``replace`` only inside :meth:`hold` after validating the current entry. The default delegates to ``put`` for legacy/custom backends whose put is a true overwrite; backends with first-writer-wins put semantics must override it. """ await self.put(scope_key, key, entry)Replace a live entry while the caller holds this key's lock.
Request idempotency uses first-writer-wins :meth:
putsemantics. Stateful protocols such as webhook delivery claims also need an owner-fenced live-state transition. Callers MUST invokereplaceonly inside :meth:holdafter validating the current entry. The default delegates toputfor legacy/custom backends whose put is a true overwrite; backends with first-writer-wins put semantics must override it. async def supports_atomic_hold(self) ‑> bool-
Expand source code
async def supports_atomic_hold(self) -> bool: """Return whether ``hold`` provides the backend's atomic guarantee. This explicit capability is delegable by wrappers such as :class:`LazyBackend`; method-identity inspection is not. """ return type(self).hold is not IdempotencyBackend.holdReturn whether
holdprovides the backend's atomic guarantee.This explicit capability is delegable by wrappers such as :class:
LazyBackend; method-identity inspection is not. async def supports_atomic_put_if_absent(self) ‑> bool-
Expand source code
async def supports_atomic_put_if_absent(self) -> bool: """Return whether ``put_if_absent`` is implemented atomically.""" return ( type(self).put_if_absent is not IdempotencyBackend.put_if_absent or await self.supports_atomic_hold() )Return whether
put_if_absentis implemented atomically.
class IdempotencyReservationError (*args, **kwargs)-
Expand source code
class IdempotencyReservationError(RuntimeError): """The reservation lifecycle cannot safely commit a replayable operation."""The reservation lifecycle cannot safely commit a replayable operation.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class IdempotencyStore (backend: IdempotencyBackend,
ttl_seconds: int = 86400,
hash_fn: Callable[[dict[str, Any]], str] = <function canonical_json_sha256>,
*,
clock: Callable[[], float] = <built-in function time>,
raise_on_persist_error: bool = False)-
Expand source code
class IdempotencyStore: """Coordinator that binds canonical hashing to a storage backend. :param backend: A concrete :class:`IdempotencyBackend`. :param ttl_seconds: How long cached responses remain replayable. Must be within the spec's ``[3600, 604800]`` range (1h to 7d). 86400 (24h) is the recommended floor and matches the compliance storyboard. :param hash_fn: Optional override for the canonical hash function. Defaults to :func:`canonical_json_sha256`. Exposed for tests and for anyone who wants to experiment with alternative equivalence rules — though note the spec mandates RFC 8785 JCS for interop. :param raise_on_persist_error: When true, translate a backend write failure after handler completion to a retryable ``SERVICE_UNAVAILABLE`` error instead of returning an unrecorded success. Production handlers using this mode must independently deduplicate their business side effects, because the handler may already have crossed that boundary. """ def __init__( self, backend: IdempotencyBackend, ttl_seconds: int = 86400, hash_fn: Callable[[dict[str, Any]], str] = canonical_json_sha256, *, clock: Callable[[], float] = time.time, raise_on_persist_error: bool = False, ) -> None: if not _MIN_TTL_SECONDS <= ttl_seconds <= _MAX_TTL_SECONDS: raise ValueError( f"ttl_seconds must be in [{_MIN_TTL_SECONDS}, {_MAX_TTL_SECONDS}] " f"per AdCP spec (capabilities.idempotency.replay_ttl_seconds), " f"got {ttl_seconds}" ) self.backend = backend self.ttl_seconds = ttl_seconds self._hash_fn = hash_fn self._clock = clock self.raise_on_persist_error = raise_on_persist_error self._warned_fallback_hold = False @asynccontextmanager async def _hold(self, scope_key: str, key: str) -> AsyncIterator[None]: """Use native backend locking or a deprecated process-local fallback.""" if await self.backend.supports_atomic_hold(): async with self.backend.hold(scope_key, key): yield return if not self._warned_fallback_hold: warnings.warn( f"{type(self.backend).__name__} does not implement hold(); using " "process-local idempotency locking. Implement hold() for " "cross-process atomicity; this compatibility fallback is deprecated.", DeprecationWarning, stacklevel=3, ) self._warned_fallback_hold = True slot = (scope_key, key) state = _legacy_backend_lock_state(self.backend) async with state.guard: lock = state.locks.get(slot) if lock is None: lock = asyncio.Lock() state.locks[slot] = lock async with lock: yield def capability(self) -> dict[str, Any]: """Return the capabilities fragment declaring this store's replay window. Embed under ``capabilities.adcp.idempotency`` on the seller's ``get_adcp_capabilities`` response. Buyers read this to reason about retry-safe windows (AdCP #2315):: caps.adcp.idempotency = idempotency.capability() # → {"supported": True, "replay_ttl_seconds": 86400} ``supported`` became REQUIRED in AdCP 3.0 GA — agents emitting only ``replay_ttl_seconds`` fail strict schema validation on the new capabilities response. """ return {"supported": True, "replay_ttl_seconds": self.ttl_seconds} def reserve( self, params: Any, context: Any, *, connection: AsyncConnection[Any] | None = None, operation: str = "handler", ) -> AbstractAsyncContextManager[PgReservation]: """Explicitly reserve a scoped request with PostgreSQL transaction ownership. Requires a direct ``PgBackend`` and an authenticated request containing an idempotency key. Use instead of ``@wrap`` when business and replay writes must commit atomically. Business SQL must use the yielded slot's connection; call ``slot.record(response)`` before context exit. """ if not isinstance(self.backend, PgBackend): raise TypeError("Transactional reservations require a direct PgBackend") scope, key, payload = self._prepare(params, context, operation) if scope is None or key is None: raise ValueError("Transactional reservations require an idempotency_key") return self.backend.reserve( scope, key, self._hash_fn(payload), ttl_seconds=self.ttl_seconds, connection=connection, operation=operation, ) def wrap(self, handler: HandlerFn) -> HandlerFn: """Decorator that adds idempotency semantics to an AdCP handler method. Supports three calling conventions the framework dispatches with: 1. **Positional** ``handler(self, params, context)`` — the default for non-projected tools (``get_products``, ``create_media_buy``, etc.). 2. **Keyword** ``handler(self, params=..., context=...)`` — same shape, just kwargs. 3. **Arg-projected** ``handler(self, **arg_projector_kwargs, ctx=...)`` where ``params`` is split into per-field kwargs by the framework dispatcher (e.g. ``update_media_buy`` is called as ``handler(self, media_buy_id=..., patch=..., ctx=...)``). In this mode the wrap searches the kwargs for a Pydantic model (``patch`` for update_media_buy) to extract the idempotency key and hash payload from. Adopters whose projection contains no Pydantic model (e.g. a method projecting only a list of ids) get fall-through behavior: no key found → handler runs without dedup. ``params`` is normalized to a dict before hashing; the return value is coerced to a dict for caching (via ``model_dump`` if Pydantic). The decorator always returns the handler's original object on a cache miss and a best-effort Pydantic re-validation on a hit (when the handler's declared return type exposes ``model_validate``). Callers that return raw dicts get dicts back. """ @wraps(handler) async def _wrapped(*args: Any, **kwargs: Any) -> Any: if _FROZEN_FEED_DISPATCH.get() or ( getattr(handler, "__name__", None) == "get_reporting_status" and getattr(getattr(handler, "__self__", None), "reporting_feed_store", None) is not None ): # A frozen page is replayable data, never cached authorization. return await handler(*args, **kwargs) if ( _RECEIPT_BATCH_DISPATCH.get() or getattr(handler, "__name__", None) == "sync_reporting_receipts" ): # Receipt ingress owns durable whole-batch replay after exact # account/consumer authorization. This generic key omits the # resolved account and cached hits bypass that authorization. # It also adds replayed=True, changing the immutable response. return await handler(*args, **kwargs) handler_self, hash_source, context = _resolve_call_args(args, kwargs) operation = getattr(handler, "__name__", "handler") scope_key, idempotency_key, params_dict = self._prepare(hash_source, context, operation) if scope_key is None or idempotency_key is None: # No key → spec says the server MUST reject with INVALID_REQUEST. # We let the handler run so validation layers above us (Pydantic, # FastAPI, etc.) can reject with a typed error; the middleware's # job is only to dedup when a key IS present. # # Forward the call exactly as received so all three calling # conventions (positional / keyword / arg-projected) reach # the inner handler unchanged. The wrap is signature- # transparent on the no-key path. return await handler(*args, **kwargs) payload_hash = self._hash_fn(params_dict) def _replay_cached(cached: CachedResponse) -> Any: if cached.payload_hash == payload_hash: logger.debug( "idempotency replay: scope=%s key_prefix=%s", _scope_log_id(scope_key), idempotency_key[:8], ) replay = _clone_response(cached.response) replay["replayed"] = True return replay raise IdempotencyConflictError( operation=operation, errors=[ { "code": "IDEMPOTENCY_CONFLICT", "message": ( "idempotency_key reused with a different payload " "(canonical hash mismatch)" ), } ], ) # A pure replay does not need an execution lock. Read once before # acquiring the backend hold, then re-check inside the hold on a # miss to preserve first-writer-wins under concurrent requests. cached = await self.backend.get(scope_key, idempotency_key) if cached is not None: return _replay_cached(cached) async def _execute_locked() -> Any: # The backend lock spans lookup, handler execution, and commit. # This is the critical invariant: a concurrent request for the # same scoped key waits, then observes the winner's cached result. async with self._hold(scope_key, idempotency_key): cached = await self.backend.get(scope_key, idempotency_key) if cached is not None: return _replay_cached(cached) response = await handler(*args, **kwargs) # Deep-copy when caching so post-return mutation of the caller's # copy can't poison future replays. `_clone_response` also deep- # copies on the hit path, giving independent objects per replay. response_dict = copy.deepcopy(_to_dict(response)) entry = CachedResponse( payload_hash=payload_hash, response=response_dict, expires_at_epoch=self._clock() + self.ttl_seconds, ) # Commit while the execution lock is still held. This does not # make unrelated business writes transactional with the cache, # but it prevents a concurrent duplicate handler execution. try: await self.backend.put(scope_key, idempotency_key, entry) except Exception as exc: logger.warning( "Idempotency cache put failed for scope=%s key_prefix=%s — " "handler completed but a subsequent retry with this key will " "re-execute rather than replay. This indicates an operational " "issue with the idempotency backend.", _scope_log_id(scope_key), idempotency_key[:8], exc_info=True, ) if self.raise_on_persist_error: # Local import avoids making the server middleware # module depend on the decisioning package at import # time. The exception is translated consistently by # both MCP and A2A server boundaries. from adcp.decisioning.types import AdcpError raise AdcpError( "SERVICE_UNAVAILABLE", message=( "Idempotency persistence failed after handler completion; " "the outcome may be unknown. Retry with the same key only " "when downstream effects are independently deduplicated." ), recovery="transient", ) from exc return response execution_task = asyncio.create_task(_execute_locked()) _SUPERVISED_OPERATIONS.add(execution_task) execution_task.add_done_callback(_finish_supervised_operation) try: return await asyncio.shield(execution_task) except asyncio.CancelledError: def _log_late_failure(task: asyncio.Task[Any]) -> None: if task.cancelled(): return exc = task.exception() if exc is not None: logger.error( "Idempotency operation failed after its request was cancelled", exc_info=(type(exc), exc, exc.__traceback__), ) execution_task.add_done_callback(_log_late_failure) raise # Register the wrapper for the boot-time validator at # adcp.decisioning.validate_idempotency. WeakSet membership — # not a public attribute — so adopters can't spoof "wrapped" # by stamping an attr on a plain function. The wrapper is # registered, not the original handler: re-decorating a forked # copy of `handler` would otherwise falsely flag both. # # ``is_wrapped`` checks registry membership at each __wrapped__ link. # Full inspect.unwrap()-then-check would reach the original handler, # which is not registered, and miss this closure under outer decorators. _WRAPPED_FUNCTIONS.add(_wrapped) return _wrapped def _prepare( self, params: Any, context: Any, operation: str ) -> tuple[str | None, str | None, dict[str, Any]]: """Normalize inputs and extract the (scope_key, key, params_dict) tuple. ``scope_key`` composes ``tenant_id`` (when present) with ``caller_identity`` so cache entries are isolated across tenants even if the seller's principal IDs are only unique within each tenant. Returns ``(None, None, params_dict)`` when idempotency doesn't apply because no key was supplied. If a key is supplied without caller identity, raises :class:`IdempotencyScopeError` so the request fails closed before side effects execute. """ params_dict = _to_dict(params) idempotency_key = params_dict.get("idempotency_key") if not isinstance(idempotency_key, str) or not idempotency_key: return None, None, params_dict scope_key = _extract_scope_key(context) if scope_key is None: # No caller identity: we can't safely scope the key. Spec requires # per-principal scope; anything else is a cross-principal replay # attack surface. Fail closed instead of executing an unscoped # idempotent operation. self._warn_missing_principal_once() raise IdempotencyScopeError(operation=operation) return scope_key, idempotency_key, params_dict _missing_principal_warned: bool = False def _warn_missing_principal_once(self) -> None: """Emit a one-time warning when the middleware sees a key but no principal. Silent fall-through is the worst DX: the seller drops in ``@idempotency.wrap``, ships, and doesn't discover until incident review that no dedup ever happened. Fire once per store instance so operators see the signal without filling logs on every request. """ if self._missing_principal_warned: return self._missing_principal_warned = True warnings.warn( "IdempotencyStore received a request with idempotency_key but no " "caller_identity on ToolContext — request is rejected. This usually " "means your transport isn't populating the authenticated principal. " "A2A: wire an a2a-sdk auth middleware that sets ServerCallContext.user; " "MCP: populate ToolContext.caller_identity from your FastMCP auth " "middleware (see adcp.server.idempotency README). " "This warning fires once per IdempotencyStore instance.", UserWarning, stacklevel=3, )Coordinator that binds canonical hashing to a storage backend.
:param backend: A concrete :class:
IdempotencyBackend. :param ttl_seconds: How long cached responses remain replayable. Must be within the spec's[3600, 604800]range (1h to 7d). 86400 (24h) is the recommended floor and matches the compliance storyboard. :param hash_fn: Optional override for the canonical hash function. Defaults to :func:canonical_json_sha256(). Exposed for tests and for anyone who wants to experiment with alternative equivalence rules — though note the spec mandates RFC 8785 JCS for interop. :param raise_on_persist_error: When true, translate a backend write failure after handler completion to a retryableSERVICE_UNAVAILABLEerror instead of returning an unrecorded success. Production handlers using this mode must independently deduplicate their business side effects, because the handler may already have crossed that boundary.Methods
def capability(self) ‑> dict[str, typing.Any]-
Expand source code
def capability(self) -> dict[str, Any]: """Return the capabilities fragment declaring this store's replay window. Embed under ``capabilities.adcp.idempotency`` on the seller's ``get_adcp_capabilities`` response. Buyers read this to reason about retry-safe windows (AdCP #2315):: caps.adcp.idempotency = idempotency.capability() # → {"supported": True, "replay_ttl_seconds": 86400} ``supported`` became REQUIRED in AdCP 3.0 GA — agents emitting only ``replay_ttl_seconds`` fail strict schema validation on the new capabilities response. """ return {"supported": True, "replay_ttl_seconds": self.ttl_seconds}Return the capabilities fragment declaring this store's replay window.
Embed under
capabilities.adcp.idempotencyon the seller'sget_adcp_capabilitiesresponse. Buyers read this to reason about retry-safe windows (AdCP #2315)::caps.adcp.idempotency = idempotency.capability() # → {"supported": True, "replay_ttl_seconds": 86400}supportedbecame REQUIRED in AdCP 3.0 GA — agents emitting onlyreplay_ttl_secondsfail strict schema validation on the new capabilities response. def reserve(self,
params: Any,
context: Any,
*,
connection: AsyncConnection[Any] | None = None,
operation: str = 'handler') ‑> AbstractAsyncContextManager[PgReservation]-
Expand source code
def reserve( self, params: Any, context: Any, *, connection: AsyncConnection[Any] | None = None, operation: str = "handler", ) -> AbstractAsyncContextManager[PgReservation]: """Explicitly reserve a scoped request with PostgreSQL transaction ownership. Requires a direct ``PgBackend`` and an authenticated request containing an idempotency key. Use instead of ``@wrap`` when business and replay writes must commit atomically. Business SQL must use the yielded slot's connection; call ``slot.record(response)`` before context exit. """ if not isinstance(self.backend, PgBackend): raise TypeError("Transactional reservations require a direct PgBackend") scope, key, payload = self._prepare(params, context, operation) if scope is None or key is None: raise ValueError("Transactional reservations require an idempotency_key") return self.backend.reserve( scope, key, self._hash_fn(payload), ttl_seconds=self.ttl_seconds, connection=connection, operation=operation, )Explicitly reserve a scoped request with PostgreSQL transaction ownership.
Requires a direct
PgBackendand an authenticated request containing an idempotency key. Use instead of@wrapwhen business and replay writes must commit atomically. Business SQL must use the yielded slot's connection; callslot.record(response)before context exit. def wrap(self, handler: HandlerFn) ‑> Callable[..., Awaitable[typing.Any]]-
Expand source code
def wrap(self, handler: HandlerFn) -> HandlerFn: """Decorator that adds idempotency semantics to an AdCP handler method. Supports three calling conventions the framework dispatches with: 1. **Positional** ``handler(self, params, context)`` — the default for non-projected tools (``get_products``, ``create_media_buy``, etc.). 2. **Keyword** ``handler(self, params=..., context=...)`` — same shape, just kwargs. 3. **Arg-projected** ``handler(self, **arg_projector_kwargs, ctx=...)`` where ``params`` is split into per-field kwargs by the framework dispatcher (e.g. ``update_media_buy`` is called as ``handler(self, media_buy_id=..., patch=..., ctx=...)``). In this mode the wrap searches the kwargs for a Pydantic model (``patch`` for update_media_buy) to extract the idempotency key and hash payload from. Adopters whose projection contains no Pydantic model (e.g. a method projecting only a list of ids) get fall-through behavior: no key found → handler runs without dedup. ``params`` is normalized to a dict before hashing; the return value is coerced to a dict for caching (via ``model_dump`` if Pydantic). The decorator always returns the handler's original object on a cache miss and a best-effort Pydantic re-validation on a hit (when the handler's declared return type exposes ``model_validate``). Callers that return raw dicts get dicts back. """ @wraps(handler) async def _wrapped(*args: Any, **kwargs: Any) -> Any: if _FROZEN_FEED_DISPATCH.get() or ( getattr(handler, "__name__", None) == "get_reporting_status" and getattr(getattr(handler, "__self__", None), "reporting_feed_store", None) is not None ): # A frozen page is replayable data, never cached authorization. return await handler(*args, **kwargs) if ( _RECEIPT_BATCH_DISPATCH.get() or getattr(handler, "__name__", None) == "sync_reporting_receipts" ): # Receipt ingress owns durable whole-batch replay after exact # account/consumer authorization. This generic key omits the # resolved account and cached hits bypass that authorization. # It also adds replayed=True, changing the immutable response. return await handler(*args, **kwargs) handler_self, hash_source, context = _resolve_call_args(args, kwargs) operation = getattr(handler, "__name__", "handler") scope_key, idempotency_key, params_dict = self._prepare(hash_source, context, operation) if scope_key is None or idempotency_key is None: # No key → spec says the server MUST reject with INVALID_REQUEST. # We let the handler run so validation layers above us (Pydantic, # FastAPI, etc.) can reject with a typed error; the middleware's # job is only to dedup when a key IS present. # # Forward the call exactly as received so all three calling # conventions (positional / keyword / arg-projected) reach # the inner handler unchanged. The wrap is signature- # transparent on the no-key path. return await handler(*args, **kwargs) payload_hash = self._hash_fn(params_dict) def _replay_cached(cached: CachedResponse) -> Any: if cached.payload_hash == payload_hash: logger.debug( "idempotency replay: scope=%s key_prefix=%s", _scope_log_id(scope_key), idempotency_key[:8], ) replay = _clone_response(cached.response) replay["replayed"] = True return replay raise IdempotencyConflictError( operation=operation, errors=[ { "code": "IDEMPOTENCY_CONFLICT", "message": ( "idempotency_key reused with a different payload " "(canonical hash mismatch)" ), } ], ) # A pure replay does not need an execution lock. Read once before # acquiring the backend hold, then re-check inside the hold on a # miss to preserve first-writer-wins under concurrent requests. cached = await self.backend.get(scope_key, idempotency_key) if cached is not None: return _replay_cached(cached) async def _execute_locked() -> Any: # The backend lock spans lookup, handler execution, and commit. # This is the critical invariant: a concurrent request for the # same scoped key waits, then observes the winner's cached result. async with self._hold(scope_key, idempotency_key): cached = await self.backend.get(scope_key, idempotency_key) if cached is not None: return _replay_cached(cached) response = await handler(*args, **kwargs) # Deep-copy when caching so post-return mutation of the caller's # copy can't poison future replays. `_clone_response` also deep- # copies on the hit path, giving independent objects per replay. response_dict = copy.deepcopy(_to_dict(response)) entry = CachedResponse( payload_hash=payload_hash, response=response_dict, expires_at_epoch=self._clock() + self.ttl_seconds, ) # Commit while the execution lock is still held. This does not # make unrelated business writes transactional with the cache, # but it prevents a concurrent duplicate handler execution. try: await self.backend.put(scope_key, idempotency_key, entry) except Exception as exc: logger.warning( "Idempotency cache put failed for scope=%s key_prefix=%s — " "handler completed but a subsequent retry with this key will " "re-execute rather than replay. This indicates an operational " "issue with the idempotency backend.", _scope_log_id(scope_key), idempotency_key[:8], exc_info=True, ) if self.raise_on_persist_error: # Local import avoids making the server middleware # module depend on the decisioning package at import # time. The exception is translated consistently by # both MCP and A2A server boundaries. from adcp.decisioning.types import AdcpError raise AdcpError( "SERVICE_UNAVAILABLE", message=( "Idempotency persistence failed after handler completion; " "the outcome may be unknown. Retry with the same key only " "when downstream effects are independently deduplicated." ), recovery="transient", ) from exc return response execution_task = asyncio.create_task(_execute_locked()) _SUPERVISED_OPERATIONS.add(execution_task) execution_task.add_done_callback(_finish_supervised_operation) try: return await asyncio.shield(execution_task) except asyncio.CancelledError: def _log_late_failure(task: asyncio.Task[Any]) -> None: if task.cancelled(): return exc = task.exception() if exc is not None: logger.error( "Idempotency operation failed after its request was cancelled", exc_info=(type(exc), exc, exc.__traceback__), ) execution_task.add_done_callback(_log_late_failure) raise # Register the wrapper for the boot-time validator at # adcp.decisioning.validate_idempotency. WeakSet membership — # not a public attribute — so adopters can't spoof "wrapped" # by stamping an attr on a plain function. The wrapper is # registered, not the original handler: re-decorating a forked # copy of `handler` would otherwise falsely flag both. # # ``is_wrapped`` checks registry membership at each __wrapped__ link. # Full inspect.unwrap()-then-check would reach the original handler, # which is not registered, and miss this closure under outer decorators. _WRAPPED_FUNCTIONS.add(_wrapped) return _wrappedDecorator that adds idempotency semantics to an AdCP handler method.
Supports three calling conventions the framework dispatches with:
- Positional
handler(self, params, context)— the default for non-projected tools (get_products,create_media_buy, etc.). - Keyword
handler(self, params=..., context=...)— same shape, just kwargs. - Arg-projected
handler(self, **arg_projector_kwargs, ctx=...)whereparamsis split into per-field kwargs by the framework dispatcher (e.g.update_media_buyis called ashandler(self, media_buy_id=..., patch=..., ctx=...)). In this mode the wrap searches the kwargs for a Pydantic model (patchfor update_media_buy) to extract the idempotency key and hash payload from. Adopters whose projection contains no Pydantic model (e.g. a method projecting only a list of ids) get fall-through behavior: no key found → handler runs without dedup.
paramsis normalized to a dict before hashing; the return value is coerced to a dict for caching (viamodel_dumpif Pydantic). The decorator always returns the handler's original object on a cache miss and a best-effort Pydantic re-validation on a hit (when the handler's declared return type exposesmodel_validate). Callers that return raw dicts get dicts back. - Positional
class LazyBackend (factory: LazyBackendFactory, *, allow_clear_all: bool = False)-
Expand source code
class LazyBackend(IdempotencyBackend): """Resolve an :class:`IdempotencyBackend` lazily on first use. :param factory: A zero-arg callable returning an :class:`IdempotencyBackend` or an awaitable that resolves to one. Invoked at most once across the wrapper's lifetime once it succeeds; re-invoked only if a prior attempt raised. :param allow_clear_all: When ``True``, expose :meth:`clear_all`, delegating to the resolved backend's ``clear_all`` or ``clear`` method. Defaults to ``False`` because bulk clearing is dangerous on shared production stores and the SDK uses method presence as the reset-safety contract. Concurrency: the first ``get``/``put``/``delete_expired`` triggers resolution. Multiple concurrent first operations share a single factory invocation via an :class:`asyncio.Lock`; later callers reuse the cached instance without locking on the hot path. """ def __init__( self, factory: LazyBackendFactory, *, allow_clear_all: bool = False, ) -> None: self._factory = factory self._allow_clear_all = allow_clear_all self._backend: IdempotencyBackend | None = None self._lock = asyncio.Lock() async def _resolve(self) -> IdempotencyBackend: """Return the resolved backend, invoking the factory once on first use. Resolve-once + concurrency-safe: the fast path returns the memoized instance without locking. The slow path holds ``_lock`` so concurrent first callers share a single factory invocation; the double-check inside the lock means a caller that waited on the lock sees the instance the winner resolved. A factory that raises is not memoized — the next call retries. """ cached = self._backend if cached is not None: return cached async with self._lock: # Re-read under the lock: a task that lost the race to acquire it # must observe the winner's resolved instance, not re-run the # factory. (Read into a local so the narrowing is on the local, # not the instance attribute another task may have mutated.) cached = self._backend if cached is not None: return cached result = self._factory() resolved = await result if isinstance(result, Awaitable) else result if not isinstance(resolved, IdempotencyBackend): raise TypeError( "LazyBackend factory must resolve to an IdempotencyBackend, " f"got {type(resolved).__name__}" ) self._backend = resolved return resolved async def get(self, scope_key: str, key: str) -> CachedResponse | None: return await (await self._resolve()).get(scope_key, key) async def put(self, scope_key: str, key: str, entry: CachedResponse) -> None: await (await self._resolve()).put(scope_key, key, entry) async def replace(self, scope_key: str, key: str, entry: CachedResponse) -> None: await (await self._resolve()).replace(scope_key, key, entry) async def current_time(self) -> float: return await (await self._resolve()).current_time() @asynccontextmanager async def hold(self, scope_key: str, key: str) -> AsyncIterator[None]: async with (await self._resolve()).hold(scope_key, key): yield async def put_if_absent(self, scope_key: str, key: str, entry: CachedResponse) -> bool: return await (await self._resolve()).put_if_absent(scope_key, key, entry) async def supports_atomic_hold(self) -> bool: return await (await self._resolve()).supports_atomic_hold() async def supports_atomic_put_if_absent(self) -> bool: return await (await self._resolve()).supports_atomic_put_if_absent() async def delete_expired(self, now_epoch: float | None = None) -> int: return await (await self._resolve()).delete_expired(now_epoch) async def _clear_all(self) -> None: """Delegate a bulk clear to the resolved backend. Resolves the backend (so the factory runs if it hasn't yet) and delegates to its ``clear_all`` or ``clear`` method, raising if the resolved backend supports neither. Exposed as ``clear_all`` only when the wrapper is constructed with ``allow_clear_all=True`` (see :meth:`__getattr__`). """ backend = await self._resolve() clear = getattr(backend, "clear_all", None) or getattr(backend, "clear", None) if clear is None: raise NotImplementedError( f"Resolved backend {type(backend).__name__} does not support " "clear_all() or clear()." ) await clear() def __getattr__(self, name: str) -> object: """Expose ``clear_all`` only when opted in. Mirrors the JS wrapper, which attaches ``clearAll`` to the returned object solely when ``{ clearAll: true }`` — reset-safety code uses ``hasattr``/method presence as the "safe to flush" contract, so the attribute must genuinely be absent otherwise. ``__getattr__`` is only consulted for names not found normally, so this never shadows the delegating methods above. """ if name == "clear_all" and self.__dict__.get("_allow_clear_all"): return self._clear_all raise AttributeError(f"{type(self).__name__!r} object has no attribute {name!r}")Resolve an :class:
IdempotencyBackendlazily on first use.:param factory: A zero-arg callable returning an :class:
IdempotencyBackendor an awaitable that resolves to one. Invoked at most once across the wrapper's lifetime once it succeeds; re-invoked only if a prior attempt raised. :param allow_clear_all: WhenTrue, expose :meth:clear_all, delegating to the resolved backend'sclear_allorclearmethod. Defaults toFalsebecause bulk clearing is dangerous on shared production stores and the SDK uses method presence as the reset-safety contract.Concurrency: the first
get/put/delete_expiredtriggers resolution. Multiple concurrent first operations share a single factory invocation via an :class:asyncio.Lock; later callers reuse the cached instance without locking on the hot path.Ancestors
- IdempotencyBackend
- abc.ABC
Inherited members
class MemoryBackend (*, clock: Callable[[], float] = <built-in function time>)-
Expand source code
class MemoryBackend(IdempotencyBackend): """In-process dict-backed store. Suitable for tests, single-process reference implementations, and local development. **Not suitable for multi-process deployments** — each worker has its own cache, so a retry that lands on a different worker is treated as a fresh request. Thread safety: the backend uses an :class:`asyncio.Lock` to serialize mutations of the shared dict. Reads go through the lock too; for a pure in-process backend this is cheap and prevents torn reads across concurrent ``get``/``put`` interleaving. :param clock: Callable returning the current epoch seconds. Override for tests that need to advance time deterministically without monkeypatching :mod:`time`. Defaults to :func:`time.time`. """ def __init__(self, *, clock: Callable[[], float] = time.time) -> None: self._store: dict[tuple[str, str], CachedResponse] = {} self._lock = asyncio.Lock() self._key_locks: weakref.WeakValueDictionary[tuple[str, str], asyncio.Lock] = ( weakref.WeakValueDictionary() ) self._clock = clock async def get(self, scope_key: str, key: str) -> CachedResponse | None: async with self._lock: entry = self._store.get((scope_key, key)) if entry is None: return None if entry.expires_at_epoch <= self._clock(): # Lazy expiry — drop the stale entry so the next request # treats the slot as fresh and races to repopulate. del self._store[(scope_key, key)] return None return entry async def put( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: async with self._lock: self._store[(scope_key, key)] = entry async def current_time(self) -> float: return self._clock() @asynccontextmanager async def hold(self, scope_key: str, key: str) -> AsyncIterator[None]: """Serialize one idempotent handler execution in this process.""" slot = (scope_key, key) async with self._lock: key_lock = self._key_locks.get(slot) if key_lock is None: key_lock = asyncio.Lock() self._key_locks[slot] = key_lock async with key_lock: yield async def put_if_absent(self, scope_key: str, key: str, entry: CachedResponse) -> bool: """Atomically claim a missing or expired slot.""" slot = (scope_key, key) async with self._lock: existing = self._store.get(slot) if existing is not None and existing.expires_at_epoch > self._clock(): return False self._store[slot] = entry return True async def delete_expired(self, now_epoch: float | None = None) -> int: cutoff = now_epoch if now_epoch is not None else self._clock() async with self._lock: stale = [k for k, v in self._store.items() if v.expires_at_epoch <= cutoff] for k in stale: del self._store[k] return len(stale) async def clear(self) -> None: """Remove all cached entries. Test-suite hook — handy for resetting state between fixtures when a single :class:`MemoryBackend` is shared across multiple tests. """ async with self._lock: self._store.clear() async def _size(self) -> int: """Test-only: return the current entry count.""" async with self._lock: return len(self._store)In-process dict-backed store.
Suitable for tests, single-process reference implementations, and local development. Not suitable for multi-process deployments — each worker has its own cache, so a retry that lands on a different worker is treated as a fresh request.
Thread safety: the backend uses an :class:
asyncio.Lockto serialize mutations of the shared dict. Reads go through the lock too; for a pure in-process backend this is cheap and prevents torn reads across concurrentget/putinterleaving.:param clock: Callable returning the current epoch seconds. Override for tests that need to advance time deterministically without monkeypatching :mod:
time. Defaults to :func:time.time.Ancestors
- IdempotencyBackend
- abc.ABC
Methods
async def clear(self) ‑> None-
Expand source code
async def clear(self) -> None: """Remove all cached entries. Test-suite hook — handy for resetting state between fixtures when a single :class:`MemoryBackend` is shared across multiple tests. """ async with self._lock: self._store.clear()Remove all cached entries.
Test-suite hook — handy for resetting state between fixtures when a single :class:
MemoryBackendis shared across multiple tests. async def hold(self, scope_key: str, key: str) ‑> AsyncIterator[None]-
Expand source code
@asynccontextmanager async def hold(self, scope_key: str, key: str) -> AsyncIterator[None]: """Serialize one idempotent handler execution in this process.""" slot = (scope_key, key) async with self._lock: key_lock = self._key_locks.get(slot) if key_lock is None: key_lock = asyncio.Lock() self._key_locks[slot] = key_lock async with key_lock: yieldSerialize one idempotent handler execution in this process.
async def put_if_absent(self,
scope_key: str,
key: str,
entry: CachedResponse) ‑> bool-
Expand source code
async def put_if_absent(self, scope_key: str, key: str, entry: CachedResponse) -> bool: """Atomically claim a missing or expired slot.""" slot = (scope_key, key) async with self._lock: existing = self._store.get(slot) if existing is not None and existing.expires_at_epoch > self._clock(): return False self._store[slot] = entry return TrueAtomically claim a missing or expired slot.
Inherited members
class PgBackend (*, pool: Any, lock_pool: Any, table_name: str = 'adcp_idempotency')-
Expand source code
class PgBackend(IdempotencyBackend): """PostgreSQL-backed :class:`IdempotencyBackend`. Multi-worker durable replay cache. Adopters running ≥2 processes wire this in place of :class:`MemoryBackend` so a retry that lands on a different worker still replays the cached response. Example:: from psycopg_pool import AsyncConnectionPool from adcp.server.idempotency import IdempotencyStore, PgBackend pool = AsyncConnectionPool("postgresql://...", min_size=2, max_size=10) lock_pool = AsyncConnectionPool("postgresql://...", min_size=2, max_size=10) backend = PgBackend(pool=pool, lock_pool=lock_pool) await backend.create_schema() # idempotent; safe to call on every boot store = IdempotencyStore(backend=backend, ttl_seconds=86400) **Atomicity caveat.** The backend holds a Postgres advisory lock and writes the cache in its transaction, but the cache write is NOT automatically in the same transaction as the handler's unrelated business writes. A crash after an external side effect but before cache commit can still leave the slot empty; the next retry re-executes the handler. Idempotent handlers absorb this without harm. **Handlers with non-idempotent side effects** (e.g., ``INSERT INTO media_buys`` without a unique constraint on the buyer's idempotency_key) need either handler-level dedupe via a database unique constraint or the explicit :meth:`reserve` API. A reservation owns the business transaction, records the response before its commit, and releases the advisory lock only after that commit. The existing decorator keeps its original contract. **Schema bootstrap caveat.** :meth:`create_schema` uses ``CREATE TABLE IF NOT EXISTS`` — if a table with the same name but a different shape already exists (Alembic migration drift, manual DDL with ``response JSON`` instead of ``JSONB``, missing ``COLLATE "C"``), this method is a no-op and the backend will run against the wrong column types. If you manage the schema with Alembic / dbmate, copy the DDL inside :meth:`create_schema` verbatim into a migration revision — keep ``COLLATE "C"`` and ``JSONB`` identical — and skip calling :meth:`create_schema` at boot. **Response payload contract.** :attr:`CachedResponse.response` is serialized via ``json.dumps`` for the JSONB column. Values must be JSON-safe — no ``datetime``, ``Decimal``, ``set``, or ``bytes``. Coerce in your handler before returning. **Cardinality / DoS.** This backend has no row cap; only TTL bounds the table size. Per AdCP spec, per-principal rate limiting at the auth tier is required — the backend trusts that. Schedule :meth:`delete_expired` as a cron / pg_cron / app-loop sweep (``get`` self-filters expired rows, but they accumulate on disk until something deletes them). **Schema.** Created idempotently by :meth:`create_schema`: .. code-block:: sql CREATE TABLE IF NOT EXISTS adcp_idempotency ( scope_key TEXT COLLATE "C" NOT NULL, key TEXT COLLATE "C" NOT NULL, payload_hash TEXT NOT NULL, response JSONB NOT NULL, expires_at TIMESTAMPTZ NOT NULL, PRIMARY KEY (scope_key, key) ); CREATE INDEX IF NOT EXISTS adcp_idempotency_expires_idx ON adcp_idempotency (expires_at); ``COLLATE "C"`` on identifier columns avoids locale-driven equivalence (``Principal-A`` ≡ ``principal-a`` under Turkish/locale-aware collations) collapsing distinct tenants into the same cache slot. :param pool: ``psycopg_pool.AsyncConnectionPool`` owned by the caller. Each operation acquires a short-lived connection. We don't open, own, or close the pool. :param lock_pool: A distinct caller-owned pool reserved for advisory-lock transactions. It MUST NOT be the business/cache ``pool``: ``hold`` keeps one connection checked out while adopter code runs, and sharing that pool with handler SQL can deadlock under saturation. :param table_name: Override the default table name. Useful for multi-tenant schema scoping. Default ``adcp_idempotency``. :raises ImportError: when psycopg/psycopg-pool are not installed. Install via the ``pg`` extra: ``pip install 'adcp[pg]'``. :raises ValueError: when ``table_name`` is not a safe ASCII identifier (``[a-z_][a-z0-9_]{0,62}``). """ def __init__( self, *, pool: Any, # psycopg_pool.AsyncConnectionPool — Any avoids runtime psycopg import lock_pool: Any, table_name: str = DEFAULT_IDEMPOTENCY_TABLE, ) -> None: if not _PG_AVAILABLE: raise ImportError(_PG_INSTALL_HINT) if lock_pool is pool: raise ValueError("lock_pool must be distinct from pool to prevent handler deadlocks") self._pool = pool self._lock_pool = lock_pool self._table = _safe_identifier(table_name) self._hold_guard = asyncio.Lock() self._key_locks: weakref.WeakValueDictionary[tuple[str, str], asyncio.Lock] = ( weakref.WeakValueDictionary() ) self._active_connection: ContextVar[tuple[Any, asyncio.Task[Any] | None] | None] = ( ContextVar(f"adcp_idempotency_connection_{id(self)}", default=None) ) # Pre-format SQL once. Validated identifier so f-string interpolation # is byte-safe; values always go through %s parameterization. Same # convention as PgWebhookDeliverySupervisor / PgReplayStore. t = self._table self._sql_get = ( f"SELECT payload_hash, response, expires_at " # noqa: S608 f"FROM {t} WHERE scope_key = %s AND key = %s " f"AND expires_at > clock_timestamp()" ) # First-writer-wins under concurrent put. The store's pre-check # ("slot is empty or expired") is NOT a lock — two workers can # both see an empty slot and race into put. With a naive # last-writer-wins ON CONFLICT, the second put would overwrite # the first's payload_hash, violating the cache invariant # "same (scope, key) → same hash". The WHERE on the UPDATE # arm restricts the overwrite to actually-expired rows: a # concurrent fresh write becomes a no-op, both callers # observe an equivalent cached entry from the first writer. self._sql_put = ( f"INSERT INTO {t} " # noqa: S608 f"(scope_key, key, payload_hash, response, expires_at) " f"VALUES (%s, %s, %s, %s::jsonb, %s) " f"ON CONFLICT (scope_key, key) DO UPDATE SET " f" payload_hash = EXCLUDED.payload_hash, " f" response = EXCLUDED.response, " f" expires_at = EXCLUDED.expires_at " f"WHERE {t}.expires_at <= clock_timestamp()" ) self._sql_delete_expired = f"DELETE FROM {t} WHERE expires_at <= %s" # noqa: S608 self._sql_delete_expired_now = ( f"DELETE FROM {t} WHERE expires_at <= clock_timestamp()" # noqa: S608 ) self._sql_lock = "SELECT pg_advisory_xact_lock(hashtextextended(%s, 6217))" self._sql_now = "SELECT EXTRACT(EPOCH FROM clock_timestamp())" self._sql_put_if_absent = ( f"INSERT INTO {t} " # noqa: S608 f"(scope_key, key, payload_hash, response, expires_at) " f"VALUES (%s, %s, %s, %s::jsonb, %s) " f"ON CONFLICT (scope_key, key) DO UPDATE SET " f" payload_hash = EXCLUDED.payload_hash, response = EXCLUDED.response, " f" expires_at = EXCLUDED.expires_at " f"WHERE {t}.expires_at <= clock_timestamp() RETURNING 1" ) self._sql_replace = ( f"UPDATE {t} SET payload_hash = %s, response = %s::jsonb, " # noqa: S608 f"expires_at = %s WHERE scope_key = %s AND key = %s" ) def reserve( self, scope_key: str, key: str, payload_hash: str, *, ttl_seconds: int = 86400, connection: AsyncConnection[Any] | None = None, operation: str = "handler", ) -> AbstractAsyncContextManager[PgReservation]: """Reserve a key and own a business/replay transaction through commit. Prefer ``IdempotencyStore.reserve`` for authenticated request scoping and canonical hashing. An optional caller-owned psycopg connection must be idle and use the same database as the lock pool. Existing outer transactions are rejected. Pool connections are borrowed after locking. """ from adcp.server.idempotency.reservation import reserve_transaction return reserve_transaction( self, scope_key, key, payload_hash, ttl_seconds=ttl_seconds, connection=connection, operation=operation, ) async def create_schema(self) -> None: """Bootstrap the table + index. Idempotent. Safe to call on every app boot. Each DDL statement is executed separately — psycopg does not split on ``;``. """ t = self._table statements = [ f"""CREATE TABLE IF NOT EXISTS {t} ( scope_key TEXT COLLATE "C" NOT NULL, key TEXT COLLATE "C" NOT NULL, payload_hash TEXT NOT NULL, response JSONB NOT NULL, expires_at TIMESTAMPTZ NOT NULL, PRIMARY KEY (scope_key, key) )""", # Partial-free expiry index — cheap eviction sweep. f"""CREATE INDEX IF NOT EXISTS {t}_expires_idx ON {t} (expires_at)""", ] async with self._pool.connection() as conn: for stmt in statements: await conn.execute(stmt) async def get(self, scope_key: str, key: str) -> CachedResponse | None: """Read the cached entry, filtering expired rows in the WHERE clause. Lazy expiry — expired rows stay on disk until ``delete_expired`` sweeps them. ``get`` self-filters via ``clock_timestamp()`` so a stale row never replays. """ active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): return await self._get_on_connection(active[0], scope_key, key) async with self._pool.connection() as conn: return await self._get_on_connection(conn, scope_key, key) async def _get_on_connection( self, conn: Any, scope_key: str, key: str ) -> CachedResponse | None: cur = await conn.execute(self._sql_get, (scope_key, key)) row = await cur.fetchone() if row is None: return None payload_hash, response, expires_at = row return CachedResponse( payload_hash=payload_hash, response=response if isinstance(response, dict) else json.loads(response), expires_at_epoch=_to_epoch(expires_at), ) async def put( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Atomic upsert under ``(scope_key, key)``. ``ON CONFLICT DO UPDATE`` because the store only calls ``put`` after verifying the slot is empty or expired — an overwrite in that window is a legitimate retry of the write itself. """ expires_at_dt = datetime.fromtimestamp(entry.expires_at_epoch, tz=timezone.utc) params = ( scope_key, key, entry.payload_hash, json.dumps(entry.response), expires_at_dt, ) active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): await active[0].execute(self._sql_put, params) return async with self._pool.connection() as conn: await conn.execute(self._sql_put, params) async def replace( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Replace one live row while the caller holds its advisory lock.""" params = ( entry.payload_hash, json.dumps(entry.response), datetime.fromtimestamp(entry.expires_at_epoch, tz=timezone.utc), scope_key, key, ) active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): await active[0].execute(self._sql_replace, params) return raise RuntimeError("PgBackend.replace() requires hold(scope_key, key)") async def current_time(self) -> float: """Return PostgreSQL time so leases and row expiry share one clock.""" active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): cur = await active[0].execute(self._sql_now) row = await cur.fetchone() return float(row[0]) async with self._pool.connection() as conn: cur = await conn.execute(self._sql_now) row = await cur.fetchone() return float(row[0]) @asynccontextmanager async def hold(self, scope_key: str, key: str) -> AsyncIterator[None]: """Hold a cross-process transaction advisory lock for this key. The same pooled connection remains checked out while the handler runs; nested ``get``/``put`` calls reuse it via a context-local binding. Local same-key waiters coalesce before acquiring a pool connection. """ slot = (scope_key, key) async with self._hold_guard: key_lock = self._key_locks.get(slot) if key_lock is None: key_lock = asyncio.Lock() self._key_locks[slot] = key_lock # Every holder/waiter keeps a strong reference to key_lock; weak # bookkeeping drops it after the last task exits, including cancellation. # Same-process retries wait here before borrowing a scarce connection. async with key_lock: lock_identity = json.dumps([scope_key, key], separators=(",", ":")) async with self._lock_pool.connection() as conn, conn.transaction(): await conn.execute(self._sql_lock, (lock_identity,)) token = self._active_connection.set((conn, asyncio.current_task())) try: yield finally: self._active_connection.reset(token) async def put_if_absent(self, scope_key: str, key: str, entry: CachedResponse) -> bool: """Atomically insert a webhook dedup marker, including stale replace.""" expires_at_dt = datetime.fromtimestamp(entry.expires_at_epoch, tz=timezone.utc) async with self._pool.connection() as conn: cur = await conn.execute( self._sql_put_if_absent, ( scope_key, key, entry.payload_hash, json.dumps(entry.response), expires_at_dt, ), ) return await cur.fetchone() is not None async def delete_expired(self, now_epoch: float | None = None) -> int: """Best-effort sweep of expired entries. Returns rows removed.""" async with self._pool.connection() as conn: if now_epoch is None: cur = await conn.execute(self._sql_delete_expired_now) else: cutoff_dt = datetime.fromtimestamp(now_epoch, tz=timezone.utc) cur = await conn.execute(self._sql_delete_expired, (cutoff_dt,)) return cur.rowcount or 0PostgreSQL-backed :class:
IdempotencyBackend.Multi-worker durable replay cache. Adopters running ≥2 processes wire this in place of :class:
MemoryBackendso a retry that lands on a different worker still replays the cached response.Example::
from psycopg_pool import AsyncConnectionPool from adcp.server.idempotency import IdempotencyStore, PgBackend pool = AsyncConnectionPool("postgresql://...", min_size=2, max_size=10) lock_pool = AsyncConnectionPool("postgresql://...", min_size=2, max_size=10) backend = PgBackend(pool=pool, lock_pool=lock_pool) await backend.create_schema() # idempotent; safe to call on every boot store = IdempotencyStore(backend=backend, ttl_seconds=86400)Atomicity caveat. The backend holds a Postgres advisory lock and writes the cache in its transaction, but the cache write is NOT automatically in the same transaction as the handler's unrelated business writes. A crash after an external side effect but before cache commit can still leave the slot empty; the next retry re-executes the handler. Idempotent handlers absorb this without harm. Handlers with non-idempotent side effects (e.g.,
INSERT INTO media_buyswithout a unique constraint on the buyer's idempotency_key) need either handler-level dedupe via a database unique constraint or the explicit :meth:reserveAPI. A reservation owns the business transaction, records the response before its commit, and releases the advisory lock only after that commit. The existing decorator keeps its original contract.Schema bootstrap caveat. :meth:
create_schemausesCREATE TABLE IF NOT EXISTS— if a table with the same name but a different shape already exists (Alembic migration drift, manual DDL withresponse JSONinstead ofJSONB, missingCOLLATE "C"), this method is a no-op and the backend will run against the wrong column types. If you manage the schema with Alembic / dbmate, copy the DDL inside :meth:create_schemaverbatim into a migration revision — keepCOLLATE "C"andJSONBidentical — and skip calling :meth:create_schemaat boot.Response payload contract. :attr:
CachedResponse.responseis serialized viajson.dumpsfor the JSONB column. Values must be JSON-safe — nodatetime,Decimal,set, orbytes. Coerce in your handler before returning.Cardinality / DoS. This backend has no row cap; only TTL bounds the table size. Per AdCP spec, per-principal rate limiting at the auth tier is required — the backend trusts that. Schedule :meth:
delete_expiredas a cron / pg_cron / app-loop sweep (getself-filters expired rows, but they accumulate on disk until something deletes them).Schema. Created idempotently by :meth:
create_schema:.. code-block:: sql
CREATE TABLE IF NOT EXISTS adcp_idempotency ( scope_key TEXT COLLATE "C" NOT NULL, key TEXT COLLATE "C" NOT NULL, payload_hash TEXT NOT NULL, response JSONB NOT NULL, expires_at TIMESTAMPTZ NOT NULL, PRIMARY KEY (scope_key, key) ); CREATE INDEX IF NOT EXISTS adcp_idempotency_expires_idx ON adcp_idempotency (expires_at);COLLATE "C"on identifier columns avoids locale-driven equivalence (Principal-A≡principal-aunder Turkish/locale-aware collations) collapsing distinct tenants into the same cache slot.:param pool:
psycopg_pool.AsyncConnectionPoolowned by the caller. Each operation acquires a short-lived connection. We don't open, own, or close the pool. :param lock_pool: A distinct caller-owned pool reserved for advisory-lock transactions. It MUST NOT be the business/cachepool:holdkeeps one connection checked out while adopter code runs, and sharing that pool with handler SQL can deadlock under saturation. :param table_name: Override the default table name. Useful for multi-tenant schema scoping. Defaultadcp_idempotency.:raises ImportError: when psycopg/psycopg-pool are not installed. Install via the
pgextra:pip install 'adcp[pg]'. :raises ValueError: whentable_nameis not a safe ASCII identifier ([a-z_][a-z0-9_]{0,62}).Ancestors
- IdempotencyBackend
- abc.ABC
Methods
async def create_schema(self) ‑> None-
Expand source code
async def create_schema(self) -> None: """Bootstrap the table + index. Idempotent. Safe to call on every app boot. Each DDL statement is executed separately — psycopg does not split on ``;``. """ t = self._table statements = [ f"""CREATE TABLE IF NOT EXISTS {t} ( scope_key TEXT COLLATE "C" NOT NULL, key TEXT COLLATE "C" NOT NULL, payload_hash TEXT NOT NULL, response JSONB NOT NULL, expires_at TIMESTAMPTZ NOT NULL, PRIMARY KEY (scope_key, key) )""", # Partial-free expiry index — cheap eviction sweep. f"""CREATE INDEX IF NOT EXISTS {t}_expires_idx ON {t} (expires_at)""", ] async with self._pool.connection() as conn: for stmt in statements: await conn.execute(stmt)Bootstrap the table + index. Idempotent.
Safe to call on every app boot. Each DDL statement is executed separately — psycopg does not split on
;. async def current_time(self) ‑> float-
Expand source code
async def current_time(self) -> float: """Return PostgreSQL time so leases and row expiry share one clock.""" active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): cur = await active[0].execute(self._sql_now) row = await cur.fetchone() return float(row[0]) async with self._pool.connection() as conn: cur = await conn.execute(self._sql_now) row = await cur.fetchone() return float(row[0])Return PostgreSQL time so leases and row expiry share one clock.
async def delete_expired(self, now_epoch: float | None = None) ‑> int-
Expand source code
async def delete_expired(self, now_epoch: float | None = None) -> int: """Best-effort sweep of expired entries. Returns rows removed.""" async with self._pool.connection() as conn: if now_epoch is None: cur = await conn.execute(self._sql_delete_expired_now) else: cutoff_dt = datetime.fromtimestamp(now_epoch, tz=timezone.utc) cur = await conn.execute(self._sql_delete_expired, (cutoff_dt,)) return cur.rowcount or 0Best-effort sweep of expired entries. Returns rows removed.
async def get(self, scope_key: str, key: str) ‑> CachedResponse | None-
Expand source code
async def get(self, scope_key: str, key: str) -> CachedResponse | None: """Read the cached entry, filtering expired rows in the WHERE clause. Lazy expiry — expired rows stay on disk until ``delete_expired`` sweeps them. ``get`` self-filters via ``clock_timestamp()`` so a stale row never replays. """ active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): return await self._get_on_connection(active[0], scope_key, key) async with self._pool.connection() as conn: return await self._get_on_connection(conn, scope_key, key)Read the cached entry, filtering expired rows in the WHERE clause.
Lazy expiry — expired rows stay on disk until
delete_expiredsweeps them.getself-filters viaclock_timestamp()so a stale row never replays. async def hold(self, scope_key: str, key: str) ‑> AsyncIterator[None]-
Expand source code
@asynccontextmanager async def hold(self, scope_key: str, key: str) -> AsyncIterator[None]: """Hold a cross-process transaction advisory lock for this key. The same pooled connection remains checked out while the handler runs; nested ``get``/``put`` calls reuse it via a context-local binding. Local same-key waiters coalesce before acquiring a pool connection. """ slot = (scope_key, key) async with self._hold_guard: key_lock = self._key_locks.get(slot) if key_lock is None: key_lock = asyncio.Lock() self._key_locks[slot] = key_lock # Every holder/waiter keeps a strong reference to key_lock; weak # bookkeeping drops it after the last task exits, including cancellation. # Same-process retries wait here before borrowing a scarce connection. async with key_lock: lock_identity = json.dumps([scope_key, key], separators=(",", ":")) async with self._lock_pool.connection() as conn, conn.transaction(): await conn.execute(self._sql_lock, (lock_identity,)) token = self._active_connection.set((conn, asyncio.current_task())) try: yield finally: self._active_connection.reset(token)Hold a cross-process transaction advisory lock for this key.
The same pooled connection remains checked out while the handler runs; nested
get/putcalls reuse it via a context-local binding. Local same-key waiters coalesce before acquiring a pool connection. async def put(self,
scope_key: str,
key: str,
entry: CachedResponse) ‑> None-
Expand source code
async def put( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Atomic upsert under ``(scope_key, key)``. ``ON CONFLICT DO UPDATE`` because the store only calls ``put`` after verifying the slot is empty or expired — an overwrite in that window is a legitimate retry of the write itself. """ expires_at_dt = datetime.fromtimestamp(entry.expires_at_epoch, tz=timezone.utc) params = ( scope_key, key, entry.payload_hash, json.dumps(entry.response), expires_at_dt, ) active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): await active[0].execute(self._sql_put, params) return async with self._pool.connection() as conn: await conn.execute(self._sql_put, params)Atomic upsert under
(scope_key, key).ON CONFLICT DO UPDATEbecause the store only callsputafter verifying the slot is empty or expired — an overwrite in that window is a legitimate retry of the write itself. async def put_if_absent(self,
scope_key: str,
key: str,
entry: CachedResponse) ‑> bool-
Expand source code
async def put_if_absent(self, scope_key: str, key: str, entry: CachedResponse) -> bool: """Atomically insert a webhook dedup marker, including stale replace.""" expires_at_dt = datetime.fromtimestamp(entry.expires_at_epoch, tz=timezone.utc) async with self._pool.connection() as conn: cur = await conn.execute( self._sql_put_if_absent, ( scope_key, key, entry.payload_hash, json.dumps(entry.response), expires_at_dt, ), ) return await cur.fetchone() is not NoneAtomically insert a webhook dedup marker, including stale replace.
async def replace(self,
scope_key: str,
key: str,
entry: CachedResponse) ‑> None-
Expand source code
async def replace( self, scope_key: str, key: str, entry: CachedResponse, ) -> None: """Replace one live row while the caller holds its advisory lock.""" params = ( entry.payload_hash, json.dumps(entry.response), datetime.fromtimestamp(entry.expires_at_epoch, tz=timezone.utc), scope_key, key, ) active = self._active_connection.get() if active is not None and active[1] is asyncio.current_task(): await active[0].execute(self._sql_replace, params) return raise RuntimeError("PgBackend.replace() requires hold(scope_key, key)")Replace one live row while the caller holds its advisory lock.
def reserve(self,
scope_key: str,
key: str,
payload_hash: str,
*,
ttl_seconds: int = 86400,
connection: AsyncConnection[Any] | None = None,
operation: str = 'handler') ‑> AbstractAsyncContextManager[PgReservation]-
Expand source code
def reserve( self, scope_key: str, key: str, payload_hash: str, *, ttl_seconds: int = 86400, connection: AsyncConnection[Any] | None = None, operation: str = "handler", ) -> AbstractAsyncContextManager[PgReservation]: """Reserve a key and own a business/replay transaction through commit. Prefer ``IdempotencyStore.reserve`` for authenticated request scoping and canonical hashing. An optional caller-owned psycopg connection must be idle and use the same database as the lock pool. Existing outer transactions are rejected. Pool connections are borrowed after locking. """ from adcp.server.idempotency.reservation import reserve_transaction return reserve_transaction( self, scope_key, key, payload_hash, ttl_seconds=ttl_seconds, connection=connection, operation=operation, )Reserve a key and own a business/replay transaction through commit.
Prefer
IdempotencyStore.reserve()for authenticated request scoping and canonical hashing. An optional caller-owned psycopg connection must be idle and use the same database as the lock pool. Existing outer transactions are rejected. Pool connections are borrowed after locking.
Inherited members
class PgReservation (backend: _ReservationBackend,
connection: AsyncConnection[Any],
scope_key: str,
key: str,
payload_hash: str,
ttl_seconds: int,
response: dict[str, Any] | None,
replay_table_oid: int)-
Expand source code
class PgReservation: """A task-owned reservation yielded by ``PgBackend.reserve``. On a miss, execute business SQL on :attr:`connection` and call :meth:`record` exactly once before leaving the context. Successful context exit commits both writes before unlocking. On a hit, :attr:`response` is a deep copied replay envelope and business SQL is forbidden. Do not commit, roll back, or issue transaction-control SQL yourself; the reservation owns the transaction. """ def __init__( self, backend: _ReservationBackend, connection: AsyncConnection[Any], scope_key: str, key: str, payload_hash: str, ttl_seconds: int, response: dict[str, Any] | None, replay_table_oid: int, ) -> None: self._backend = backend self._connection = connection self._scope_key = scope_key self._key = key self._payload_hash = payload_hash self._ttl_seconds = ttl_seconds self._replay_table_oid = replay_table_oid self._response = copy.deepcopy(response) self._replayed = response is not None self._recorded = False self._failed = False self._active = True self._owner = asyncio.current_task() @property def replayed(self) -> bool: """Whether this reservation found a committed matching replay.""" return self._replayed @property def recorded(self) -> bool: """Whether ``record`` wrote a response (commit occurs on context exit).""" return self._recorded @property def response(self) -> dict[str, Any] | None: """Return an independent response copy, adding ``replayed`` on a hit.""" response = copy.deepcopy(self._response) if response is not None and self._replayed: response["replayed"] = True return response def _check_active(self) -> None: if not self._active or asyncio.current_task() is not self._owner: raise IdempotencyReservationError("Reservation must be used by its active owning task") @property def connection(self) -> AsyncConnection[Any]: """The business connection; available only on an active cache miss.""" self._check_active() if self._replayed: raise IdempotencyReservationError("A replay reservation cannot execute business writes") return self._connection async def record(self, response: dict[str, Any]) -> None: """Write a JSON response in the open business transaction, without committing. Errors propagate and poison the reservation even if caught by the handler. The response must remain present through transaction exit; rolling it back in a nested savepoint makes the entire reservation fail closed. """ self._check_active() if self._replayed or self._recorded or self._failed: raise IdempotencyReservationError("Only a fresh reservation can record once") try: await self._check_replay_table() response_copy = copy.deepcopy(response) encoded = json.dumps(response_copy, allow_nan=False) response_copy = json.loads(encoded) cursor = await self._connection.execute(self._backend._sql_now) row = await cursor.fetchone() assert row is not None expires_at = datetime.fromtimestamp(float(row[0]) + self._ttl_seconds, timezone.utc) cursor = await self._connection.execute( self._backend._sql_put_if_absent, (self._scope_key, self._key, self._payload_hash, encoded, expires_at), ) if await cursor.fetchone() is None: raise IdempotencyReservationError("Replay slot changed during reservation") self._response = response_copy self._recorded = True except BaseException: self._failed = True raise async def _check_replay_table(self) -> None: cursor = await self._connection.execute( "SELECT to_regclass(%s)::oid", (self._backend._table,) ) row = await cursor.fetchone() if row is None or row[0] != self._replay_table_oid: raise IdempotencyReservationError("Reservation replay table changed during transaction") async def _verify_record(self) -> None: if self._failed or not self._recorded: raise IdempotencyReservationError("Reservation exited without a successful record") await self._check_replay_table() cached = await self._backend._get_on_connection( self._connection, self._scope_key, self._key ) if ( cached is None or cached.payload_hash != self._payload_hash or cached.response != self._response ): raise IdempotencyReservationError("Recorded response was rolled back or changed")A task-owned reservation yielded by
PgBackend.reserve().On a miss, execute business SQL on :attr:
connectionand call :meth:recordexactly once before leaving the context. Successful context exit commits both writes before unlocking. On a hit, :attr:responseis a deep copied replay envelope and business SQL is forbidden. Do not commit, roll back, or issue transaction-control SQL yourself; the reservation owns the transaction.Instance variables
prop connection : AsyncConnection[Any]-
Expand source code
@property def connection(self) -> AsyncConnection[Any]: """The business connection; available only on an active cache miss.""" self._check_active() if self._replayed: raise IdempotencyReservationError("A replay reservation cannot execute business writes") return self._connectionThe business connection; available only on an active cache miss.
prop recorded : bool-
Expand source code
@property def recorded(self) -> bool: """Whether ``record`` wrote a response (commit occurs on context exit).""" return self._recordedWhether
recordwrote a response (commit occurs on context exit). prop replayed : bool-
Expand source code
@property def replayed(self) -> bool: """Whether this reservation found a committed matching replay.""" return self._replayedWhether this reservation found a committed matching replay.
prop response : dict[str, Any] | None-
Expand source code
@property def response(self) -> dict[str, Any] | None: """Return an independent response copy, adding ``replayed`` on a hit.""" response = copy.deepcopy(self._response) if response is not None and self._replayed: response["replayed"] = True return responseReturn an independent response copy, adding
replayedon a hit.
Methods
async def record(self, response: dict[str, Any]) ‑> None-
Expand source code
async def record(self, response: dict[str, Any]) -> None: """Write a JSON response in the open business transaction, without committing. Errors propagate and poison the reservation even if caught by the handler. The response must remain present through transaction exit; rolling it back in a nested savepoint makes the entire reservation fail closed. """ self._check_active() if self._replayed or self._recorded or self._failed: raise IdempotencyReservationError("Only a fresh reservation can record once") try: await self._check_replay_table() response_copy = copy.deepcopy(response) encoded = json.dumps(response_copy, allow_nan=False) response_copy = json.loads(encoded) cursor = await self._connection.execute(self._backend._sql_now) row = await cursor.fetchone() assert row is not None expires_at = datetime.fromtimestamp(float(row[0]) + self._ttl_seconds, timezone.utc) cursor = await self._connection.execute( self._backend._sql_put_if_absent, (self._scope_key, self._key, self._payload_hash, encoded, expires_at), ) if await cursor.fetchone() is None: raise IdempotencyReservationError("Replay slot changed during reservation") self._response = response_copy self._recorded = True except BaseException: self._failed = True raiseWrite a JSON response in the open business transaction, without committing.
Errors propagate and poison the reservation even if caught by the handler. The response must remain present through transaction exit; rolling it back in a nested savepoint makes the entire reservation fail closed.
class WebhookDedupClaim (status: WebhookClaimStatus, claim_token: str | None = None)-
Expand source code
@dataclass(frozen=True) class WebhookDedupClaim: """Outcome of one atomic delivery claim. ``claim_token`` is present only for ``claimed`` and is an owner capability consumed by :meth:`WebhookDedupStore.complete` or :meth:`WebhookDedupStore.release`. It must not be logged. """ status: WebhookClaimStatus claim_token: str | None = NoneOutcome of one atomic delivery claim.
claim_tokenis present only forclaimedand is an owner capability consumed by :meth:WebhookDedupStore.complete()or :meth:WebhookDedupStore.release(). It must not be logged.Instance variables
var claim_token : str | Nonevar status : Literal['claimed', 'handled', 'in_progress', 'conflict']
class WebhookDedupOwnershipError (*args, **kwargs)-
Expand source code
class WebhookDedupOwnershipError(RuntimeError): """A handler attempted to settle a claim it no longer owns."""A handler attempted to settle a claim it no longer owns.
Ancestors
- builtins.RuntimeError
- builtins.Exception
- builtins.BaseException
class WebhookDedupStore (backend: IdempotencyBackend,
ttl_seconds: int = 86400,
*,
processing_ttl_seconds: int | None = None,
namespace: str = 'webhook',
clock: Callable[[], float] = <built-in function time>)-
Expand source code
class WebhookDedupStore: """Dedup ``(sender_id, idempotency_key)`` pairs to suppress retried webhooks. :param backend: any :class:`IdempotencyBackend`. Same MemoryBackend or PgBackend type used by :class:`IdempotencyStore` is fine — the ``namespace`` parameter prefixes all sender IDs so request-side and webhook-side scopes can't alias even when sharing one backend instance. :param ttl_seconds: payload-binding and replay window. Must be within ``[86400, 604800]`` per the spec minimum. Defaults to 86400 (24h). :param processing_ttl_seconds: processing-owner lease. An exact retry receives ``in_progress`` while the lease is live and may claim after it expires. Defaults to 300 seconds so a crashed handler cannot fence retries for the complete advertised delivery horizon. Configure it above the receiver's normal publication timeout; shorter leases increase the chance of overlapping owners and therefore require idempotent application publication. Must be positive and no longer than ``ttl_seconds``. :param namespace: prefix applied to every ``sender_id`` before it hits the backend. Defaults to ``"webhook"``, which is safe when the same backend is shared with :class:`IdempotencyStore` (request-side keys are scoped by a principal_id that isn't wrapped in this namespace, so collisions are impossible). Override only if you run multiple webhook scopes against one backend (e.g., separate dedup spaces for task webhooks vs list-change webhooks). :param clock: Deprecated compatibility argument. Expiry and lease decisions use :meth:`IdempotencyBackend.current_time` so durable backends and replicas share one authoritative clock. """ def __init__( self, backend: IdempotencyBackend, ttl_seconds: int = _MIN_TTL_SECONDS, *, processing_ttl_seconds: int | None = None, namespace: str = "webhook", clock: Callable[[], float] = time.time, ) -> None: if not _MIN_TTL_SECONDS <= ttl_seconds <= _MAX_TTL_SECONDS: raise ValueError( f"ttl_seconds must be in [{_MIN_TTL_SECONDS}, {_MAX_TTL_SECONDS}] " f"per webhook spec minimum, got {ttl_seconds}" ) if not namespace: raise ValueError("namespace must be a non-empty string") effective_processing_ttl = ( _DEFAULT_PROCESSING_TTL_SECONDS if processing_ttl_seconds is None else processing_ttl_seconds ) if not 1 <= effective_processing_ttl <= ttl_seconds: raise ValueError( "processing_ttl_seconds must be positive and no greater than " f"ttl_seconds, got {effective_processing_ttl}" ) self.backend = backend self.ttl_seconds = ttl_seconds self.processing_ttl_seconds = effective_processing_ttl self.namespace = namespace _ = clock self._warned_legacy_backend = False @asynccontextmanager async def _hold(self, scope_key: str, key: str) -> AsyncIterator[None]: """Use the backend lock or a warned process-local compatibility lock.""" if await self.backend.supports_atomic_hold(): async with self.backend.hold(scope_key, key): yield return if not self._warned_legacy_backend: warnings.warn( f"{type(self.backend).__name__} does not implement hold(); using " "process-local webhook claim locking. This does not satisfy the " "multi-process AdCP 3.2 receiver contract; implement a durable " "atomic hold operation.", DeprecationWarning, stacklevel=3, ) self._warned_legacy_backend = True slot = (scope_key, key) state = _legacy_backend_lock_state(self.backend) async with state.guard: lock = state.locks.get(slot) if lock is None: lock = asyncio.Lock() state.locks[slot] = lock async with lock: yield async def claim( self, sender_id: str, idempotency_key: str, payload_hash: str, ) -> WebhookDedupClaim: """Atomically claim a delivery or classify its retained state.""" if not sender_id: raise ValueError("sender_id must be a non-empty string") if not idempotency_key: raise ValueError("idempotency_key must be a non-empty string") if not payload_hash: raise ValueError("payload_hash must be a non-empty string") scoped_sender = f"{self.namespace}:{sender_id}" async with self._hold(scoped_sender, idempotency_key): now = await self.backend.current_time() existing = await self.backend.get(scoped_sender, idempotency_key) if existing is not None: if existing.payload_hash != payload_hash: return WebhookDedupClaim(status="conflict") state = existing.response.get("webhook_state") if state == "handled": return WebhookDedupClaim(status="handled") lease_expires_at = existing.response.get("lease_expires_at") if ( state == "processing" and isinstance(lease_expires_at, (int, float)) and lease_expires_at > now ): return WebhookDedupClaim(status="in_progress") retain_until = existing.expires_at_epoch else: retain_until = now + self.ttl_seconds owner = secrets.token_urlsafe(32) claim_entry = CachedResponse( payload_hash=payload_hash, response={ "webhook_state": "processing", "owner": owner, "lease_expires_at": now + self.processing_ttl_seconds, }, expires_at_epoch=retain_until, ) if existing is None: await self.backend.put(scoped_sender, idempotency_key, claim_entry) else: await self.backend.replace(scoped_sender, idempotency_key, claim_entry) return WebhookDedupClaim(status="claimed", claim_token=owner) async def complete( self, sender_id: str, idempotency_key: str, payload_hash: str, claim_token: str, ) -> bool: """Durably mark an owned delivery handled. Returns ``True`` for the owner transition and ``False`` when the exact payload was already handled. Any missing, changed, or superseded claim fails closed with :class:`WebhookDedupOwnershipError`. """ return await self._settle( sender_id, idempotency_key, payload_hash, claim_token, handled=True, ) async def release( self, sender_id: str, idempotency_key: str, payload_hash: str, claim_token: str, ) -> bool: """Release an owned failed claim while retaining payload binding.""" return await self._settle( sender_id, idempotency_key, payload_hash, claim_token, handled=False, ) async def _settle( self, sender_id: str, idempotency_key: str, payload_hash: str, claim_token: str, *, handled: bool, ) -> bool: if not all((sender_id, idempotency_key, payload_hash, claim_token)): raise ValueError( "sender_id, idempotency_key, payload_hash, and claim_token are required" ) scoped_sender = f"{self.namespace}:{sender_id}" async with self._hold(scoped_sender, idempotency_key): existing = await self.backend.get(scoped_sender, idempotency_key) if existing is None or existing.payload_hash != payload_hash: raise WebhookDedupOwnershipError("webhook delivery claim is missing or changed") state = existing.response.get("webhook_state") if handled and state == "handled": return False if state != "processing" or existing.response.get("owner") != claim_token: raise WebhookDedupOwnershipError("webhook delivery claim is not owned") response: dict[str, object] if handled: response = {"webhook_state": "handled"} else: response = {"webhook_state": "retryable"} await self.backend.replace( scoped_sender, idempotency_key, CachedResponse( payload_hash=payload_hash, response=response, expires_at_epoch=existing.expires_at_epoch, ), ) return TrueDedup
(sender_id, idempotency_key)pairs to suppress retried webhooks.:param backend: any :class:
IdempotencyBackend. Same MemoryBackend or PgBackend type used by :class:IdempotencyStoreis fine — thenamespaceparameter prefixes all sender IDs so request-side and webhook-side scopes can't alias even when sharing one backend instance. :param ttl_seconds: payload-binding and replay window. Must be within[86400, 604800]per the spec minimum. Defaults to 86400 (24h). :param processing_ttl_seconds: processing-owner lease. An exact retry receivesin_progresswhile the lease is live and may claim after it expires. Defaults to 300 seconds so a crashed handler cannot fence retries for the complete advertised delivery horizon. Configure it above the receiver's normal publication timeout; shorter leases increase the chance of overlapping owners and therefore require idempotent application publication. Must be positive and no longer thanttl_seconds. :param namespace: prefix applied to everysender_idbefore it hits the backend. Defaults to"webhook", which is safe when the same backend is shared with :class:IdempotencyStore(request-side keys are scoped by a principal_id that isn't wrapped in this namespace, so collisions are impossible). Override only if you run multiple webhook scopes against one backend (e.g., separate dedup spaces for task webhooks vs list-change webhooks). :param clock: Deprecated compatibility argument. Expiry and lease decisions use :meth:IdempotencyBackend.current_time()so durable backends and replicas share one authoritative clock.Methods
async def claim(self, sender_id: str, idempotency_key: str, payload_hash: str) ‑> WebhookDedupClaim-
Expand source code
async def claim( self, sender_id: str, idempotency_key: str, payload_hash: str, ) -> WebhookDedupClaim: """Atomically claim a delivery or classify its retained state.""" if not sender_id: raise ValueError("sender_id must be a non-empty string") if not idempotency_key: raise ValueError("idempotency_key must be a non-empty string") if not payload_hash: raise ValueError("payload_hash must be a non-empty string") scoped_sender = f"{self.namespace}:{sender_id}" async with self._hold(scoped_sender, idempotency_key): now = await self.backend.current_time() existing = await self.backend.get(scoped_sender, idempotency_key) if existing is not None: if existing.payload_hash != payload_hash: return WebhookDedupClaim(status="conflict") state = existing.response.get("webhook_state") if state == "handled": return WebhookDedupClaim(status="handled") lease_expires_at = existing.response.get("lease_expires_at") if ( state == "processing" and isinstance(lease_expires_at, (int, float)) and lease_expires_at > now ): return WebhookDedupClaim(status="in_progress") retain_until = existing.expires_at_epoch else: retain_until = now + self.ttl_seconds owner = secrets.token_urlsafe(32) claim_entry = CachedResponse( payload_hash=payload_hash, response={ "webhook_state": "processing", "owner": owner, "lease_expires_at": now + self.processing_ttl_seconds, }, expires_at_epoch=retain_until, ) if existing is None: await self.backend.put(scoped_sender, idempotency_key, claim_entry) else: await self.backend.replace(scoped_sender, idempotency_key, claim_entry) return WebhookDedupClaim(status="claimed", claim_token=owner)Atomically claim a delivery or classify its retained state.
async def complete(self, sender_id: str, idempotency_key: str, payload_hash: str, claim_token: str) ‑> bool-
Expand source code
async def complete( self, sender_id: str, idempotency_key: str, payload_hash: str, claim_token: str, ) -> bool: """Durably mark an owned delivery handled. Returns ``True`` for the owner transition and ``False`` when the exact payload was already handled. Any missing, changed, or superseded claim fails closed with :class:`WebhookDedupOwnershipError`. """ return await self._settle( sender_id, idempotency_key, payload_hash, claim_token, handled=True, )Durably mark an owned delivery handled.
Returns
Truefor the owner transition andFalsewhen the exact payload was already handled. Any missing, changed, or superseded claim fails closed with :class:WebhookDedupOwnershipError. async def release(self, sender_id: str, idempotency_key: str, payload_hash: str, claim_token: str) ‑> bool-
Expand source code
async def release( self, sender_id: str, idempotency_key: str, payload_hash: str, claim_token: str, ) -> bool: """Release an owned failed claim while retaining payload binding.""" return await self._settle( sender_id, idempotency_key, payload_hash, claim_token, handled=False, )Release an owned failed claim while retaining payload binding.