Module adcp.server.idempotency.store

The :class:IdempotencyStore coordinator: canonical hashing + backend + decorator.

Responsibilities:

  1. Extract idempotency_key from the incoming request.
  2. Scope lookups by (scope_key, key) via the backend, where scope_key composes tenant_id (when present) with caller_identity.
  3. On cache hit with matching canonical payload hash: return the cached response and mark replayed=True on the envelope.
  4. On cache hit with a different hash: raise :class:IdempotencyConflictError.
  5. On miss: run the wrapped handler, then commit (hash, response) to the backend.

Per-scope scoping is a hard security requirement (AdCP #2315): a key from principal A on tenant T has no meaning for principal B or tenant T'. The store pulls both tenant_id and caller_identity from :class:ToolContext and composes them into a single scope key — sellers whose principal ids are only unique within a tenant (Okta group-scoped IDs, seller-internal employee IDs, SCIM per-tenant IDs) must populate tenant_id so the store can keep those tenants isolated. When no tenant_id is set, the scope collapses to caller_identity alone (safe for single-tenant deployments).

If no context / no caller_identity is supplied, the store refuses to proceed — fail-closed rather than collapse every buyer into a shared namespace.

Functions

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.

Classes

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.