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 the adcp[pg] extra. await backend.create_schema() once at boot; commits go through a fresh pool connection. Use IdempotencyStore.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:IdempotencyStore coordinator: 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.

  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.

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 JS createLazyBackend factory shape.

See :class:LazyBackend for 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 False

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.

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 out

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.

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: float

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.

Instance variables

var expires_at_epoch : float
var payload_hash : str
var 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. hold must 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:get MUST self-filter expired entries. Backends that have natural TTL primitives (Redis EXPIRE, 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_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.

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 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 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 True

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 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: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.

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.hold

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.

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_absent is 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 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.

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.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.

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 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.

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 _wrapped

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.

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: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.

Ancestors

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.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.

Ancestors

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:MemoryBackend is 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:
        yield

Serialize 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 True

Atomically 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 0

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}).

Ancestors

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 0

Best-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_expired sweeps them. get self-filters via clock_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/put calls 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 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.

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 None

Atomically 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: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.

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._connection

The 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._recorded

Whether record wrote 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._replayed

Whether 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 response

Return an independent response copy, adding replayed on 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
        raise

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.

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 = None

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.

Instance variables

var claim_token : str | None
var 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 True

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.

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 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.

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.