Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/_experimental/mcp_server/gateway_dcr_flow.py: 32%
641 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 12:01 +0000
1"""The gateway-level DCR flow for the aggregate ``/mcp`` endpoint (``mcp_gateway_dcr``).
3An OAuth-only DCR client (Claude Desktop, Claude Code, MCP Inspector) pointed at the
4aggregate ``/mcp`` endpoint discovers the gateway as its authorization server (PR 1 of
5this track) and then walks the flow implemented here:
71. ``POST /register``: stateless dynamic client registration. The ``client_id`` IS the
8 registration: the client's redirect URIs are sealed into it with the repo's
9 authenticated symmetric helper, so nothing is persisted and a forged or tampered
10 client_id simply fails to open. Clients are always public (``token_endpoint_auth_method
11 "none"``); PKCE S256 is what protects the code.
122. ``GET /authorize``: validates the client and redirect URI, requires S256 PKCE, and
13 interposes LiteLLM sign-in. Without a session cookie the browser is sent through
14 ``/sso/key/generate`` with a same-origin ``return_to`` so it lands back here after
15 login. With a session, the flow parameters and the SSO user are sealed into a per-flow
16 HttpOnly cookie (the same pattern as the upstream OAuth state relay) and the browser is
17 sent to the connect page, where the user authorizes individual servers (vaulting those
18 tokens server-side) before finishing.
193. ``POST /authorize/complete``: the deliberate finish step. A POST (not GET) bound to the
20 SameSite=Lax flow cookie, so a cross-site link cannot silently mint a code with the
21 victim's session, and the signed-in user must match the user sealed into the flow.
22 Mints a short-lived, single-use, gateway-sealed authorization code and redirects to the
23 client's registered redirect URI.
244. ``POST /token``: exchanges the code (PKCE-verified, client- and redirect-bound,
25 single-use) for the identity-only session tokens of
26 :mod:`.outbound_credentials.session_token`, re-validating that the litellm user is
27 still active first; the ``refresh_token`` grant rotates the pair the same way.
29Nothing here stores state server-side except the single-use code guard (a TTL cache
30entry). Every sealed value is authenticated encryption over the proxy salt/master key
31family, opened totally (bad input maps to an OAuth error, never a raise), and every
32identity is a stable reference re-validated live at mint, refresh, and (in the admission
33PR) tool-call time. Upstream server credentials never appear anywhere in this flow; they
34are vaulted per user by the existing ``/v1/mcp`` authorize endpoints and resolved at
35egress by user id.
36"""
38from __future__ import annotations
40import hashlib
41import hmac
42import html
43import secrets
44from base64 import urlsafe_b64encode
45from collections.abc import Awaitable, Callable, Iterable, Mapping
46from datetime import datetime, timezone
47from types import MappingProxyType
48from typing import Final, Literal, Protocol, TypeVar
49from urllib.parse import parse_qsl, urlencode, urlparse, urlunparse
51from fastapi import HTTPException, Request
52from fastapi.responses import HTMLResponse, JSONResponse, RedirectResponse, Response
53from pydantic import BaseModel, ConfigDict, Field, ValidationError
54from typing_extensions import NotRequired, ReadOnly, TypedDict, assert_never
56from litellm._logging import verbose_logger
57from litellm.caching.caching import DualCache
58from litellm.proxy._experimental.mcp_server.oauth_utils import (
59 TOKEN_NO_CACHE_HEADERS,
60 canonical_resource_uri,
61 canonicalize_url_identity,
62 get_request_base_url,
63 is_loopback_redirect_host,
64 validate_redirect_uri_shape,
65)
66from litellm.proxy._experimental.mcp_server.outbound_credentials.session_credentials import (
67 SessionRefreshOpened,
68 SessionSigningConfigError,
69 active_session_signing_keys,
70 open_session_refresh_bearer,
71)
72from litellm.proxy._experimental.mcp_server.outbound_credentials.session_token import (
73 SESSION_ISSUER,
74 SESSION_REFRESH_TTL_SECONDS,
75 MintedSessionToken,
76 OpenedSessionToken,
77 SessionAudience,
78 SessionPrincipal,
79 SessionSigningKeys,
80 is_session_refresh_token,
81 is_session_token,
82 mint_session_refresh_token,
83 mint_session_token,
84 open_session_refresh_token,
85 open_session_token,
86)
87from litellm.proxy.common_utils.encrypt_decrypt_utils import (
88 decrypt_value_helper,
89 encrypt_value_helper,
90)
91from litellm.proxy.common_utils.html_forms.native_client_consent import (
92 render_native_client_consent_page,
93)
94from litellm.types.mcp_server.mcp_server_manager import MCPServer
96GATEWAY_DCR_CLIENT_ID_PREFIX: Final = "llm_dcrc_"
97"""Marker prefix on every gateway-issued DCR client_id so the root authorize/token
98endpoints can route an aggregate-flow request without decrypting, and existing per-server
99flows (whose client_ids are upstream-issued) are never captured by the aggregate arm."""
101GATEWAY_AUTH_CODE_PREFIX: Final = "llm_gcode_"
102"""Marker prefix on the gateway-sealed authorization code, distinct from the bridge
103``llm_bcode_`` so neither flow can consume the other's codes."""
105CONNECT_FLOW_COOKIE_PREFIX: Final = "mcp_connect_flow_"
106"""Per-flow HttpOnly cookie holding the sealed connect flow, keyed by a short random
107handle carried in the connect-page URL (the same handle-plus-cookie pattern as the
108``mcp_oauth_state_`` upstream relay, for the same reasons: replica-safe with no
109server-side session store, and the sealed value never appears in a URL)."""
111CONNECT_FLOW_TTL_SECONDS: Final = 600
112GATEWAY_AUTH_CODE_TTL_SECONDS: Final = 120
113MANUAL_DELIVERY_AUTH_CODE_TTL_SECONDS: Final = 300
114"""Lifetime of a code the user delivers by hand (headless/remote client, LIT-4863 class):
115copy-pasting a callback URL from a laptop browser to an SSH session is slower than a
116browser redirect, so manual-delivery codes get 5 minutes instead of 2, still well under
117the 10-minute ceiling RFC 6749 section 4.1.2 recommends. Single-use and PKCE binding are
118unchanged, so the longer window only extends how long the legitimate holder has to paste
119it, not what an observer could do with it."""
120_CLAIM_TTL_BUFFER_SECONDS: Final = 60
121_USED_CODE_CACHE_PREFIX: Final = "mcp_gateway_dcr_code_used:"
122_USED_FLOW_CACHE_PREFIX: Final = "mcp_gateway_dcr_flow_used:"
123_USED_REFRESH_CACHE_PREFIX: Final = "mcp_gateway_dcr_refresh_used:"
125MAX_REDIRECT_URIS: Final = 4
126MAX_REDIRECT_URI_LENGTH: Final = 256
127MAX_CLIENT_ID_LENGTH: Final = 2048
128"""Registration bounds. They exist to bound the sealed client_id, which rides inside
129every session-token claim set. Four 256-character ASCII URIs seal to roughly 1.5KB;
130the encoded client_id is checked against its own cap before registration succeeds.
131VS Code registers four callbacks for its web and desktop environments."""
133MAX_STATE_LENGTH: Final = 1024
134"""Bound on the client ``state`` sealed into the flow cookie and echoed on the auth-code
135redirect. An unbounded ``state`` can push the sealed cookie past the browser's ~4KB cap
136(silently dropped, breaking the flow); spec clients send a short opaque value."""
138MIN_CODE_VERIFIER_LENGTH: Final = 43
139MAX_CODE_VERIFIER_LENGTH: Final = 128
140"""RFC 7636 section 4.1 bounds for the PKCE ``code_verifier``. Enforced so an out-of-range
141verifier gets a clean ``invalid_request`` instead of an opaque PKCE-mismatch."""
143_UNPREFIXED: Final = ""
144"""Prefix for a sealed value that carries no wire marker because it is never routed by
145prefix (the connect flow lives only in its own per-handle cookie, opened by that one
146handle). Named so the empty-string argument to ``_seal`` / ``_open_sealed`` reads as
147deliberate rather than a typo."""
149_CLIENT_RECORD_DEBUG_KEY: Final = "gateway_dcr_client"
150_CONNECT_FLOW_DEBUG_KEY: Final = "gateway_connect_flow"
151_AUTH_CODE_DEBUG_KEY: Final = "gateway_authorization_code"
153ReloadUserFailure = Literal["unresolvable", "unavailable", "faulted", "no_active_key"]
154ReloadUser = Callable[[str], Awaitable[ReloadUserFailure | None]]
155VendorCredentialState = Literal["present", "absent", "unavailable"]
156"""The per-user vendor credential read has three outcomes: present, absent, or unavailable."""
158_DB_UNAVAILABLE_DESCRIPTION: Final = "the gateway database is unavailable; retry"
159_DB_FAULTED_DESCRIPTION: Final = (
160 "the gateway database reported a fault that is not a transient outage; "
161 "retrying will not help until the gateway deployment is repaired"
162)
164PROXY_API_AUDIENCE: Final[SessionAudience] = "proxy_api"
165"""The audience a native client (``lite login --pkce``, a Go CLI) asks for by sending the
166proxy base URL itself as the RFC 8707 ``resource``: the grant then mints the proxy-API CLI
167credential that LLM routes accept, instead of the MCP-only session pair."""
169ProxyCredentialMintFailure = Literal[ReloadUserFailure, "not_a_member", "team_required"]
172class MintedProxyCredential(BaseModel):
173 model_config = ConfigDict(frozen=True)
174 key: str = Field(min_length=1)
175 expires_in: int = Field(gt=0)
176 user_id: str = Field(min_length=1)
177 team_id: str | None = None
180class MintProxyCredential(Protocol):
181 """Injected proxy-API credential minter ``(user_id, team_id)``: reloads the user live,
182 checks team membership, refuses a teamless grant for a user who has teams to pick from,
183 and mints the same credential ``lite login`` mints."""
185 def __call__( 185 ↛ exitline 185 didn't return from function '__call__' because
186 self, user_id: str, team_id: str | None, /
187 ) -> Awaitable[MintedProxyCredential | ProxyCredentialMintFailure]: ...
190TOKEN_EXCHANGE_GRANT_TYPE: Final = "urn:ietf:params:oauth:grant-type:token-exchange"
193def supported_grant_types(token_exchange_available: bool) -> tuple[str, ...]:
194 """The grants ``/token`` can serve on this deployment. The RFC 8693 exchange is listed
195 only where the JWT auth that proves a subject token is on, backed by a database, and
196 licensed, so a client never selects a grant the gateway would then refuse."""
197 if token_exchange_available: 197 ↛ 198line 197 didn't jump to line 198 because the condition on line 197 was never true
198 return ("authorization_code", "refresh_token", TOKEN_EXCHANGE_GRANT_TYPE)
199 return ("authorization_code", "refresh_token")
202"""RFC 8693: a native client that already holds a token from the customer's identity
203provider trades it for the proxy-API credential without a browser round trip."""
205_IssuedTokenType = Literal["urn:ietf:params:oauth:token-type:access_token"]
206ACCESS_TOKEN_TOKEN_TYPE: Final[_IssuedTokenType] = "urn:ietf:params:oauth:token-type:access_token"
207SUBJECT_TOKEN_TYPES: Final = frozenset(
208 {
209 "urn:ietf:params:oauth:token-type:jwt",
210 "urn:ietf:params:oauth:token-type:id_token",
211 ACCESS_TOKEN_TOKEN_TYPE,
212 }
213)
216class SubjectIdentity(BaseModel):
217 model_config = ConfigDict(frozen=True)
218 user_id: str = Field(min_length=1)
219 team_id: str | None = None
222class SubjectTokenRefusal(BaseModel):
223 model_config = ConfigDict(frozen=True)
224 error: Literal["unsupported_grant_type", "invalid_request", "temporarily_unavailable"]
225 description: str = Field(min_length=1)
228class ExchangeSubjectToken(Protocol):
229 """Injected RFC 8693 subject-token verifier ``(subject_token, request)``: proves the
230 IdP token the way the proxy's own JWT auth does and names the litellm user and team it
231 stands for, or says why this gateway will not take it."""
233 def __call__(self, subject_token: str, request: Request, /) -> Awaitable[SubjectIdentity | SubjectTokenRefusal]: ... 233 ↛ exitline 233 didn't return from function '__call__' because
236class ConsentTeam(BaseModel):
237 model_config = ConfigDict(frozen=True)
238 team_id: str = Field(min_length=1)
239 team_alias: str | None = None
242class LookupVendorCredential(Protocol):
243 """Injected read of a user's vendor credential for one server."""
245 def __call__(self, user_id: str, server_id: str, /) -> Awaitable[VendorCredentialState]: ... 245 ↛ exitline 245 didn't return from function '__call__' because
248class LookupServerReachability(Protocol):
249 def __call__(self, user_id: str, server_id: str, /) -> Awaitable[bool]: ... 249 ↛ exitline 249 didn't return from function '__call__' because
252class LookupConsentTeams(Protocol):
253 """Injected lookup of the teams a signed-in user may bind a proxy-API credential to."""
255 def __call__(self, user_id: str, /) -> Awaitable[tuple[ConsentTeam, ...] | ReloadUserFailure]: ... 255 ↛ exitline 255 didn't return from function '__call__' because
258async def _refuse_proxy_credential(user_id: str, team_id: str | None) -> ProxyCredentialMintFailure:
259 return "unresolvable"
262async def _refuse_subject_token(subject_token: str, request: Request) -> SubjectTokenRefusal:
263 return SubjectTokenRefusal(
264 error="unsupported_grant_type", description="this gateway is not configured to exchange IdP tokens"
265 )
268async def _unavailable_vendor_credential(user_id: str, server_id: str) -> VendorCredentialState:
269 return "unavailable"
272async def _unreachable_server(user_id: str, server_id: str) -> bool:
273 return False
276class GatewayDcrClient(BaseModel):
277 """The registration record sealed into a gateway DCR ``client_id``.
279 ``extra="forbid"`` so a sealed value of another type (an auth code, a connect flow)
280 that happened to decrypt under the shared key can never validate as a client record:
281 cross-type confusion is rejected at the model boundary, not left to differing required
282 fields."""
284 model_config = ConfigDict(frozen=True, extra="forbid")
285 redirect_uris: tuple[str, ...] = Field(min_length=1, max_length=MAX_REDIRECT_URIS)
286 iat: int
289class _ConnectFlow(BaseModel):
290 """One in-flight authorize: the SSO user it belongs to and the client parameters
291 needed to mint the code at the finish step. Sealed into the per-flow cookie. ``jti``
292 makes the flow single-use at complete; ``extra="forbid"`` rejects cross-type
293 confusion."""
295 model_config = ConfigDict(frozen=True, extra="forbid")
296 user_id: str = Field(min_length=1)
297 client_id: str = Field(min_length=1)
298 redirect_uri: str = Field(min_length=1)
299 state: str
300 code_challenge: str = Field(min_length=1)
301 jti: str = Field(min_length=1)
302 exp: int
303 resource_server_id: str | None = None
304 audience: SessionAudience | None = None
307class _GatewayAuthCode(BaseModel):
308 """The gateway-sealed authorization code: the user consent it represents and the
309 bindings the token endpoint must verify (client, redirect URI, PKCE challenge),
310 plus a ``jti`` for the single-use guard. ``extra="forbid"`` rejects cross-type
311 confusion."""
313 model_config = ConfigDict(frozen=True, extra="forbid")
314 user_id: str = Field(min_length=1)
315 client_id: str = Field(min_length=1)
316 redirect_uri: str = Field(min_length=1)
317 code_challenge: str = Field(min_length=1)
318 jti: str = Field(min_length=1)
319 iat: int
320 exp: int
321 resource_server_id: str | None = None
322 audience: SessionAudience | None = None
323 team_id: str | None = None
326def is_gateway_dcr_client_id(client_id: str | None) -> bool:
327 """Cheap prefix routing test so the root endpoints only enter the aggregate arm for
328 clients this flow registered; every other client_id keeps today's behavior."""
329 return client_id is not None and client_id.startswith(GATEWAY_DCR_CLIENT_ID_PREFIX)
332def _oauth_error(status_code: int, error: str, description: str) -> JSONResponse:
333 """RFC 6749 section 5.2 / RFC 7591 section 3.2.2 error body. Descriptions carry no
334 token, code, or URL material so they are safe to relay to any client."""
335 return JSONResponse(
336 status_code=status_code,
337 content={"error": error, "error_description": description},
338 headers=TOKEN_NO_CACHE_HEADERS,
339 )
342def _seal(prefix: str, payload: BaseModel) -> str:
343 """Serialized ``exclude_none`` for the same reason session JWTs are minted that way: an
344 optional claim that is unset never reaches the wire, so during a rolling deploy a blob
345 sealed by a new pod without the new claim set stays byte-compatible with predating pods
346 whose strict models forbid unknown keys. This holds for every sealed artifact and every
347 future optional claim by construction; it requires each optional field to default to
348 ``None`` so reopening restores exactly what was sealed."""
349 return prefix + encrypt_value_helper(payload.model_dump_json(exclude_none=True))
352_SealedModelT = TypeVar("_SealedModelT", bound=BaseModel)
355def _open_sealed(value: str, prefix: str, model: type[_SealedModelT], debug_key: str) -> _SealedModelT | None:
356 """Open a sealed value totally: anything that is not prefix-shaped, does not decrypt,
357 or does not validate returns ``None`` for the caller to map onto an OAuth error."""
358 if not value.startswith(prefix): 358 ↛ 360line 358 didn't jump to line 360 because the condition on line 358 was always true
359 return None
360 decrypted: Final = decrypt_value_helper(value[len(prefix) :], debug_key, return_original_value=False)
361 if not isinstance(decrypted, str):
362 return None
363 try:
364 return model.model_validate_json(decrypted)
365 except ValidationError:
366 return None
369def open_gateway_dcr_client(client_id: str) -> GatewayDcrClient | None:
370 return _open_sealed(client_id, GATEWAY_DCR_CLIENT_ID_PREFIX, GatewayDcrClient, _CLIENT_RECORD_DEBUG_KEY)
373async def register_aggregate_client(
374 request: Request, request_body: Mapping[str, object], token_exchange_available: bool
375) -> Response:
376 """RFC 7591 dynamic registration against the gateway itself, statelessly.
378 Only ``redirect_uris`` is authoritative; every client is registered as a public
379 ``token_endpoint_auth_method "none"`` client regardless of what it asked for (RFC
380 7591 lets the server override metadata), because the gateway never issues client
381 secrets: possession of a secret would add nothing over the mandatory S256 PKCE, and a
382 stateless registration has nowhere to keep one. Nothing is persisted, so open
383 registration cannot be used to fill storage.
385 Redirect-URI *hygiene* is not decided here: :func:`validate_redirect_uri_shape` is
386 the single owner of that rule across the MCP OAuth surface, so allowlisted native
387 callbacks (``cursor://``) are accepted and fragments, missing hosts, userinfo
388 (``https://claude.ai@attacker.example/cb``) and backslash hosts are rejected exactly
389 as they are on /authorize and /callback.
391 What this endpoint does decide is its own trust policy, which is deliberately wider
392 than :func:`validate_trusted_redirect_uri`'s: registration is *public*, so any https
393 client may register (that is what lets a hosted MCP client register at all), and the
394 controls are mandatory S256 PKCE plus the consent screen showing the client origin.
395 http is confined to loopback per RFC 8252 section 7.3.
396 """
397 raw_uris: Final = request_body.get("redirect_uris")
398 if not isinstance(raw_uris, list) or not raw_uris or len(raw_uris) > MAX_REDIRECT_URIS:
399 return _oauth_error(
400 400,
401 "invalid_redirect_uri",
402 f"redirect_uris must be a list of 1 to {MAX_REDIRECT_URIS} URIs",
403 )
404 if not all(isinstance(uri, str) and len(uri) <= MAX_REDIRECT_URI_LENGTH for uri in raw_uris):
405 return _oauth_error(
406 400,
407 "invalid_redirect_uri",
408 f"each redirect URI must be a string of at most {MAX_REDIRECT_URI_LENGTH} characters",
409 )
410 for uri in raw_uris:
411 parsed = urlparse(uri)
412 try:
413 if validate_redirect_uri_shape(parsed):
414 continue # allowlisted native callback, e.g. cursor://
415 except HTTPException as exc:
416 # The shared validator speaks HTTP; RFC 7591 registration answers with an OAuth
417 # error object, so translate the shape without re-deciding the rule.
418 return _oauth_error(400, "invalid_redirect_uri", str(exc.detail))
419 if parsed.scheme == "https" or (parsed.scheme == "http" and is_loopback_redirect_host(parsed)):
420 continue
421 return _oauth_error(
422 400,
423 "invalid_redirect_uri",
424 "each redirect URI must be https, http on a loopback host, or a registered native callback",
425 )
426 now: Final = datetime.now(timezone.utc)
427 client_id: Final = _seal(
428 GATEWAY_DCR_CLIENT_ID_PREFIX, GatewayDcrClient(redirect_uris=tuple(raw_uris), iat=int(now.timestamp()))
429 )
430 if len(client_id) > MAX_CLIENT_ID_LENGTH:
431 return _oauth_error(400, "invalid_client_metadata", "registered metadata is too large")
432 return JSONResponse(
433 status_code=201,
434 content={
435 "client_id": client_id,
436 "client_id_issued_at": int(now.timestamp()),
437 "redirect_uris": list(raw_uris),
438 "token_endpoint_auth_method": "none",
439 "grant_types": list(supported_grant_types(token_exchange_available)),
440 "response_types": ["code"],
441 },
442 )
445def _flow_cookie_name(handle: str) -> str:
446 return f"{CONNECT_FLOW_COOKIE_PREFIX}{handle}"
449def _cookie_path_and_secure(request: Request) -> tuple[str, bool]:
450 parsed: Final = urlparse(get_request_base_url(request))
451 return parsed.path or "/", parsed.scheme == "https"
454def _append_query_params(url: str, params: Iterable[tuple[str, str]]) -> str:
455 parsed: Final = urlparse(url)
456 query: Final = (*parse_qsl(parsed.query, keep_blank_values=True), *params)
457 return urlunparse(parsed._replace(query=urlencode(query)))
460def relative_request_url(request: Request) -> str:
461 """The request's own path and query as a same-origin ``return_to`` target for the
462 login round-trip; relative by construction, so it can never leave the gateway."""
463 path: Final = request.url.path
464 return f"{path}?{request.url.query}" if request.url.query else path
467def resolve_scoped_resource_server(request: Request, resource: str | None) -> MCPServer | None:
468 """Resolve an RFC 8707 ``resource`` value to the single gateway-owned server it
469 names, or ``None`` for every other shape: absent, the aggregate resource, a foreign
470 host, an unparseable value, a multi-server path, an unknown name, or any server mode the
471 keyless gateway flow does not serve (whose protected-resource metadata never directs a
472 client here). ``None`` means the flow stays unscoped and byte-identical to today, so a
473 hostile or confused ``resource`` can never widen anything; a resolved server only ever
474 NARROWS the session via the sealed scope.
476 Resolution is an IDENTITY question, deliberately free of the per-IP visibility filter:
477 access is enforced where it belongs (grant intersection at admission, IP checks on the
478 MCP routes), while filtering here would mint an entitlement-wide UNSCOPED bearer exactly
479 when the caller asked to narrow, and would let authorize-time vs token-time IP drift
480 turn a matching redemption into a spurious ``invalid_target``."""
481 if resource is None:
482 return None
483 canonical: Final = canonical_resource_uri(resource)
484 if canonical is None:
485 return None
486 base: Final = canonicalize_url_identity(get_request_base_url(request))
487 if canonical == f"{base}/mcp" or not canonical.startswith(f"{base}/"):
488 return None
489 from litellm.proxy._experimental.mcp_server.auth.user_api_key_auth_mcp import ( # noqa: PLC0415 # proxy import cycle
490 MCPRequestHandler,
491 )
492 from litellm.proxy._experimental.mcp_server.mcp_server_manager import ( # noqa: PLC0415 # proxy import cycle
493 global_mcp_server_manager,
494 )
496 names: Final = MCPRequestHandler.extract_target_server_names_from_path(canonical[len(base) :])
497 if len(names) != 1:
498 return None
499 server: Final = global_mcp_server_manager.get_mcp_server_by_name(names[0])
500 if server is None or not (server.is_gateway_managed_oauth2 or server.advertises_gateway_authorization_server):
501 return None
502 return server
505def aggregate_authorize(
506 request: Request,
507 client_id: str,
508 redirect_uri: str,
509 state: str,
510 code_challenge: str | None,
511 code_challenge_method: str | None,
512 response_type: str | None,
513 session_user_id: str | None,
514 resource: str | None = None,
515) -> Response:
516 """The aggregate authorize verb: validate the client, require S256 PKCE, interpose
517 LiteLLM sign-in, and hand the browser to the connect page with the flow sealed into a
518 per-flow cookie.
520 A per-server RFC 8707 ``resource`` naming a gateway-managed oauth2 server scopes the
521 flow to that one server: the scope is sealed into the flow, carried into the code, and
522 bound into the session token. The connect URL carries only the flow handle; the page
523 learns the client origin, the scoped server, and whether its vendor OAuth is done from
524 :func:`describe_connect_flow`, which reads the sealed flow, so nothing a link can carry
525 steers which server the page authorizes or names on the confirmation.
527 Validation failures respond directly with 400 and never redirect: per RFC 6749
528 section 4.1.2.1 an unvalidated redirect URI must not receive an error redirect, and
529 once the client is at fault there is no trusted place to send the browser.
530 """
531 rejected: Final = _rejected_authorize_request(
532 client_id, redirect_uri, state, code_challenge, code_challenge_method, response_type
533 )
534 if rejected is not None: 534 ↛ 536line 534 didn't jump to line 536 because the condition on line 534 was always true
535 return rejected
536 base_url: Final = get_request_base_url(request)
537 if session_user_id is None:
538 return _login_redirect(base_url, request)
539 scoped_server: Final = resolve_scoped_resource_server(request, resource)
540 handle: Final = secrets.token_urlsafe(24)
541 flow: Final = _new_connect_flow(
542 session_user_id=session_user_id,
543 client_id=client_id,
544 redirect_uri=redirect_uri,
545 state=state,
546 code_challenge=code_challenge or "",
547 resource_server_id=scoped_server.server_id if scoped_server is not None else None,
548 audience=None,
549 )
550 connect_url: Final = _append_query_params(f"{base_url}/ui/connect", (("connect_flow", handle),))
551 response: Final = RedirectResponse(connect_url, status_code=303)
552 _set_flow_cookie(response, request, handle, flow)
553 return response
556async def native_client_authorize(
557 request: Request,
558 client_id: str,
559 redirect_uri: str,
560 state: str,
561 code_challenge: str | None,
562 code_challenge_method: str | None,
563 response_type: str | None,
564 session_user_id: str | None,
565 lookup_consent_teams: LookupConsentTeams,
566) -> Response:
567 """The authorize verb for a native client that named the proxy API itself as its
568 RFC 8707 ``resource``: the same client, redirect, PKCE, and sign-in checks as the
569 aggregate verb plus a loopback-only redirect (the credential this grant mints is the
570 user's personal proxy key, which belongs on their own machine and never behind a hosted
571 callback), then the consent page rendered right here (no connect-page interlude, since
572 there is no per-server vaulting to do) with the flow sealed into the per-flow cookie
573 and its handle carried only in the form, never in a URL."""
574 rejected: Final = _rejected_authorize_request(
575 client_id, redirect_uri, state, code_challenge, code_challenge_method, response_type
576 )
577 if rejected is not None:
578 return rejected
579 if not is_loopback_redirect_host(urlparse(redirect_uri)):
580 return _oauth_error(400, "invalid_request", "a proxy-API grant may only redirect to a loopback address")
581 base_url: Final = get_request_base_url(request)
582 if session_user_id is None:
583 return _login_redirect(base_url, request)
584 teams: Final = await lookup_consent_teams(session_user_id)
585 if not isinstance(teams, tuple):
586 return _consent_lookup_failure_response(teams)
587 handle: Final = secrets.token_urlsafe(24)
588 flow: Final = _new_connect_flow(
589 session_user_id=session_user_id,
590 client_id=client_id,
591 redirect_uri=redirect_uri,
592 state=state,
593 code_challenge=code_challenge or "",
594 resource_server_id=None,
595 audience=PROXY_API_AUDIENCE,
596 )
597 page: Final = render_native_client_consent_page(
598 client_origin=_origin_only(redirect_uri),
599 user_id=session_user_id,
600 teams=tuple((team.team_id, team.team_alias or team.team_id) for team in teams),
601 flow_handle=handle,
602 complete_url=f"{base_url}/authorize/complete",
603 )
604 response: Final = HTMLResponse(page, headers=_CONSENT_PAGE_HEADERS)
605 _set_flow_cookie(response, request, handle, flow)
606 return response
609_CONSENT_PAGE_HEADERS: Final = MappingProxyType(
610 {
611 **TOKEN_NO_CACHE_HEADERS,
612 "X-Frame-Options": "DENY",
613 "Content-Security-Policy": "frame-ancestors 'none'",
614 }
615)
617NATIVE_CLIENT_AUTH_CONTRACT_VERSION: Final = 1
618"""The version a native client checks before trusting the rest of the discovery document.
619Bump it only when an existing field changes meaning or goes away; adding fields is free."""
622class NativeClientAuthContract(TypedDict):
623 contract_version: ReadOnly[int]
624 issuer: ReadOnly[str]
625 authorization_endpoint: ReadOnly[str]
626 token_endpoint: ReadOnly[str]
627 registration_endpoint: ReadOnly[str]
628 revocation_endpoint: ReadOnly[str]
629 resource: ReadOnly[str]
630 response_types_supported: ReadOnly[tuple[str, ...]]
631 grant_types_supported: ReadOnly[tuple[str, ...]]
632 code_challenge_methods_supported: ReadOnly[tuple[str, ...]]
633 token_endpoint_auth_methods_supported: ReadOnly[tuple[str, ...]]
634 revocation_endpoint_auth_methods_supported: ReadOnly[tuple[str, ...]]
637def native_client_auth_contract(request: Request, token_exchange_available: bool) -> NativeClientAuthContract:
638 """The versioned discovery document at ``/.well-known/litellm-cli-auth``: everything a
639 native client (in any language) needs to run the sign-in without reading LiteLLM
640 source. ``resource`` is the exact value to send as the RFC 8707 ``resource`` parameter
641 on authorize and token requests so the grant is issued for the proxy API."""
642 base_url: Final = get_request_base_url(request)
643 contract: Final[NativeClientAuthContract] = {
644 "contract_version": NATIVE_CLIENT_AUTH_CONTRACT_VERSION,
645 "issuer": base_url,
646 "authorization_endpoint": f"{base_url}/authorize",
647 "token_endpoint": f"{base_url}/token",
648 "registration_endpoint": f"{base_url}/register",
649 "revocation_endpoint": f"{base_url}/revoke",
650 "resource": base_url,
651 "response_types_supported": ("code",),
652 "grant_types_supported": supported_grant_types(token_exchange_available),
653 "code_challenge_methods_supported": ("S256",),
654 "token_endpoint_auth_methods_supported": ("none",),
655 "revocation_endpoint_auth_methods_supported": ("none",),
656 }
657 return contract
660def is_proxy_api_resource(request: Request, resource: str | None) -> bool:
661 """True when the RFC 8707 ``resource`` names the proxy itself (its base URL), which is
662 how a native client asks for the proxy-API audience rather than an MCP session."""
663 if resource is None:
664 return False
665 canonical: Final = canonical_resource_uri(resource)
666 return canonical is not None and canonical == canonicalize_url_identity(get_request_base_url(request))
669def _rejected_authorize_request(
670 client_id: str,
671 redirect_uri: str,
672 state: str,
673 code_challenge: str | None,
674 code_challenge_method: str | None,
675 response_type: str | None,
676) -> Response | None:
677 client: Final = open_gateway_dcr_client(client_id)
678 if client is None: 678 ↛ 680line 678 didn't jump to line 680 because the condition on line 678 was always true
679 return _oauth_error(400, "invalid_client", "unknown or malformed client_id")
680 if redirect_uri not in client.redirect_uris:
681 return _oauth_error(400, "invalid_request", "redirect_uri is not registered for this client")
682 if response_type != "code":
683 return _oauth_error(400, "unsupported_response_type", "response_type must be 'code'")
684 if not code_challenge or code_challenge_method != "S256":
685 return _oauth_error(
686 400,
687 "invalid_request",
688 "PKCE is required: send code_challenge with code_challenge_method=S256",
689 )
690 if len(state) > MAX_STATE_LENGTH:
691 return _oauth_error(400, "invalid_request", f"state must be at most {MAX_STATE_LENGTH} characters")
692 return None
695def _login_redirect(base_url: str, request: Request) -> Response:
696 return_to: Final = urlencode((("return_to", relative_request_url(request)),))
697 return RedirectResponse(f"{base_url}/sso/key/generate?{return_to}", status_code=303)
700def _new_connect_flow(
701 session_user_id: str,
702 client_id: str,
703 redirect_uri: str,
704 state: str,
705 code_challenge: str,
706 resource_server_id: str | None,
707 audience: SessionAudience | None,
708) -> _ConnectFlow:
709 now: Final = datetime.now(timezone.utc)
710 return _ConnectFlow(
711 user_id=session_user_id,
712 client_id=client_id,
713 redirect_uri=redirect_uri,
714 state=state,
715 code_challenge=code_challenge,
716 jti=secrets.token_urlsafe(24),
717 exp=int(now.timestamp()) + CONNECT_FLOW_TTL_SECONDS,
718 resource_server_id=resource_server_id,
719 audience=audience,
720 )
723def _set_flow_cookie(response: Response, request: Request, handle: str, flow: _ConnectFlow) -> None:
724 path, secure = _cookie_path_and_secure(request)
725 response.set_cookie(
726 key=_flow_cookie_name(handle),
727 value=_seal(_UNPREFIXED, flow),
728 max_age=CONNECT_FLOW_TTL_SECONDS,
729 path=path,
730 secure=secure,
731 httponly=True,
732 samesite="lax",
733 )
736def _consent_lookup_failure_response(failure: ReloadUserFailure) -> Response:
737 match failure:
738 case "unavailable":
739 return _oauth_error(503, "temporarily_unavailable", _DB_UNAVAILABLE_DESCRIPTION)
740 case "faulted":
741 return _oauth_error(503, "temporarily_unavailable", _DB_FAULTED_DESCRIPTION)
742 case "unresolvable":
743 return _oauth_error(500, "server_error", "the gateway is not configured to resolve users")
744 case "no_active_key":
745 return _oauth_error(403, "access_denied", "the signed-in user is not active")
746 case _:
747 assert_never(failure)
750def _origin_only(url: str) -> str:
751 """Scheme+host for display on the connect page; never the full redirect URI, whose
752 path or query could carry values that do not belong in a page URL or logs."""
753 parsed: Final = urlparse(url)
754 return f"{parsed.scheme}://{parsed.netloc}" if parsed.netloc else ""
757def _open_flow_for(
758 request: Request, flow_handle: str, session_user_id: str | None, now: datetime
759) -> _ConnectFlow | Response:
760 sealed_flow: Final = request.cookies.get(_flow_cookie_name(flow_handle))
761 if sealed_flow is None: 761 ↛ 763line 761 didn't jump to line 763 because the condition on line 761 was always true
762 return _oauth_error(400, "invalid_request", "unknown or expired connect flow")
763 flow: Final = _open_sealed(sealed_flow, _UNPREFIXED, _ConnectFlow, _CONNECT_FLOW_DEBUG_KEY)
764 if flow is None or now.timestamp() >= flow.exp:
765 return _oauth_error(400, "invalid_request", "unknown or expired connect flow")
766 if session_user_id is None:
767 return _oauth_error(401, "login_required", "sign in to LiteLLM to finish connecting")
768 if session_user_id != flow.user_id:
769 return _oauth_error(403, "access_denied", "the signed-in user does not match this connect flow")
770 return flow
773async def _flow_target(
774 flow: _ConnectFlow, lookup_server_reachability: LookupServerReachability
775) -> tuple[Literal["unscoped", "interactive", "m2m", "stale"], MCPServer | None]:
776 if flow.resource_server_id is None:
777 return "unscoped", None
778 from litellm.proxy._experimental.mcp_server.mcp_server_manager import ( # noqa: PLC0415 # import cycle
779 MCPServerManager,
780 global_mcp_server_manager,
781 )
783 server: Final = global_mcp_server_manager.get_mcp_server_by_id(flow.resource_server_id)
784 if (
785 server is None
786 or not (server.is_gateway_managed_oauth2 or server.advertises_gateway_authorization_server)
787 or not await lookup_server_reachability(flow.user_id, server.server_id)
788 ):
789 return "stale", None
790 state: Final = (
791 "interactive"
792 if server.is_gateway_managed_oauth2 and MCPServerManager.effective_oauth2_flow(server) != "client_credentials"
793 else "m2m"
794 )
795 return state, server
798class ConnectFlowDescription(TypedDict):
799 """What the connect page is allowed to know about one in-flight flow."""
801 state: ReadOnly[Literal["unscoped", "interactive", "m2m", "stale"]]
802 client_origin: ReadOnly[str]
803 server_id: ReadOnly[str | None]
804 server_name: ReadOnly[str | None]
805 connected: ReadOnly[bool | None]
808async def _describe_opened_flow(
809 flow: _ConnectFlow,
810 lookup_vendor_credential: LookupVendorCredential,
811 lookup_server_reachability: LookupServerReachability,
812) -> ConnectFlowDescription | Response:
813 state, server = await _flow_target(flow, lookup_server_reachability)
814 if state == "interactive" and server is not None:
815 credential: Final = await lookup_vendor_credential(flow.user_id, server.server_id)
816 if credential == "unavailable":
817 return _oauth_error(503, "temporarily_unavailable", _DB_UNAVAILABLE_DESCRIPTION)
818 interactive_description: Final[ConnectFlowDescription] = {
819 "state": state,
820 "client_origin": _origin_only(flow.redirect_uri),
821 "server_id": server.server_id,
822 "server_name": server.server_name or server.alias or server.name,
823 "connected": credential == "present",
824 }
825 return interactive_description
826 described: Final[ConnectFlowDescription] = {
827 "state": state,
828 "client_origin": _origin_only(flow.redirect_uri),
829 "server_id": None if server is None else server.server_id,
830 "server_name": None if server is None else (server.server_name or server.alias or server.name),
831 "connected": state == "m2m" or None,
832 }
833 return described
836async def describe_connect_flow(
837 request: Request,
838 flow_handle: str,
839 session_user_id: str | None,
840 lookup_vendor_credential: LookupVendorCredential,
841 lookup_server_reachability: LookupServerReachability,
842) -> Response:
843 opened: Final = _open_flow_for(request, flow_handle, session_user_id, datetime.now(timezone.utc))
844 if isinstance(opened, Response): 844 ↛ 846line 844 didn't jump to line 846 because the condition on line 844 was always true
845 return opened
846 described: Final = await _describe_opened_flow(opened, lookup_vendor_credential, lookup_server_reachability)
847 return (
848 described
849 if isinstance(described, Response)
850 else JSONResponse(content=described, headers=TOKEN_NO_CACHE_HEADERS)
851 )
854async def complete_connect_flow(
855 request: Request,
856 flow_handle: str,
857 session_user_id: str | None,
858 cache: DualCache,
859 delivery: str | None = None,
860 team_id: str | None = None,
861 decision: str | None = None,
862 lookup_vendor_credential: LookupVendorCredential = _unavailable_vendor_credential,
863 lookup_server_reachability: LookupServerReachability = _unreachable_server,
864) -> Response:
865 """Mint the code only after a deliberate POST by the sealed user.
867 A scoped flow additionally requires its sealed server to have a live vendor credential
868 before a code can be minted. The check happens before the single-use claim, so a
869 premature submit can be retried after authorization; denial deliberately bypasses it.
870 """
871 if delivery not in (None, "redirect", "manual"):
872 return _oauth_error(400, "invalid_request", "delivery must be 'redirect' or 'manual'")
873 if decision not in (None, "approve", "deny"):
874 return _oauth_error(400, "invalid_request", "decision must be 'approve' or 'deny'")
875 now: Final = datetime.now(timezone.utc)
876 opened: Final = _open_flow_for(request, flow_handle, session_user_id, now)
877 if isinstance(opened, Response): 877 ↛ 879line 877 didn't jump to line 879 because the condition on line 877 was always true
878 return opened
879 if decision != "deny":
880 described: Final = await _describe_opened_flow(opened, lookup_vendor_credential, lookup_server_reachability)
881 if isinstance(described, Response):
882 return described
883 if described["state"] == "stale":
884 return _oauth_error(400, "invalid_request", "the requested MCP server is no longer available")
885 if described["connected"] is False:
886 return _oauth_error(400, "invalid_request", "authorize the requested MCP server before finishing")
887 flow_refusal: Final = _claim_refusal(
888 await _SingleUseGuard(cache).claim(
889 f"{_USED_FLOW_CACHE_PREFIX}{opened.jti}", CONNECT_FLOW_TTL_SECONDS + _CLAIM_TTL_BUFFER_SECONDS
890 ),
891 replayed=_oauth_error(
892 400, "invalid_request", "this connect flow was already completed; restart the connection"
893 ),
894 )
895 if flow_refusal is not None:
896 return flow_refusal
897 response: Final = (
898 _denied_flow_response(opened) if decision == "deny" else _approved_flow_response(opened, delivery, team_id, now)
899 )
900 path, secure = _cookie_path_and_secure(request)
901 response.delete_cookie(key=_flow_cookie_name(flow_handle), path=path, secure=secure, httponly=True, samesite="lax")
902 return response
905def _state_param(flow: _ConnectFlow) -> tuple[tuple[str, str], ...]:
906 return (("state", flow.state),) if flow.state else ()
909def _denied_flow_response(flow: _ConnectFlow) -> Response:
910 params: Final = (("error", "access_denied"), *_state_param(flow))
911 return RedirectResponse(_append_query_params(flow.redirect_uri, params), status_code=303)
914def _approved_flow_response(flow: _ConnectFlow, delivery: str | None, team_id: str | None, now: datetime) -> Response:
915 manual_delivery: Final = delivery == "manual" and is_loopback_redirect_host(urlparse(flow.redirect_uri))
916 code_ttl: Final = MANUAL_DELIVERY_AUTH_CODE_TTL_SECONDS if manual_delivery else GATEWAY_AUTH_CODE_TTL_SECONDS
917 code: Final = _seal(
918 GATEWAY_AUTH_CODE_PREFIX,
919 _GatewayAuthCode(
920 user_id=flow.user_id,
921 client_id=flow.client_id,
922 redirect_uri=flow.redirect_uri,
923 code_challenge=flow.code_challenge,
924 jti=secrets.token_urlsafe(24),
925 iat=int(now.timestamp()),
926 exp=int(now.timestamp()) + code_ttl,
927 resource_server_id=flow.resource_server_id,
928 audience=flow.audience,
929 team_id=(team_id or None) if flow.audience == PROXY_API_AUDIENCE else None,
930 ),
931 )
932 callback_url: Final = _append_query_params(flow.redirect_uri, (("code", code), *_state_param(flow)))
933 if manual_delivery:
934 return _manual_delivery_response(callback_url)
935 return RedirectResponse(callback_url, status_code=303)
938def _manual_delivery_response(callback_url: str) -> Response:
939 """The manual code-delivery page: the callback URL the 303 would have followed,
940 rendered for the user to carry to the machine the client actually runs on (paste into
941 the client's prompt, or fetch with curl from that machine's terminal). Served
942 no-store because the body holds a live single-use code, and the URL is HTML-escaped
943 because it is client-influenced. The page renders the URL as data only, never as a
944 ready-to-paste shell command: no single quoting of an attacker-influenced string is
945 correct across POSIX shells, cmd.exe, and PowerShell (cmd.exe ignores single quotes
946 and percent-expands inside double quotes), so any command string this page suggested
947 would be wrong for some shell the user might paste it into."""
948 safe_url: Final = html.escape(callback_url, quote=True)
949 minutes: Final = MANUAL_DELIVERY_AUTH_CODE_TTL_SECONDS // 60
950 body: Final = (
951 "<html><head><title>Finish connecting</title></head><body>"
952 "<h2>Almost done</h2>"
953 "<p>Your MCP client runs on a different machine, so this browser cannot deliver the"
954 " authorization code to it. On the machine where the client runs, paste this URL into"
955 " the client's prompt (Claude Code accepts the pasted callback URL), or pass it as the"
956 " quoted argument of a curl command from that machine's terminal:</p>"
957 f'<p><input type="text" value="{safe_url}" readonly size="100" onclick="this.select()"></p>'
958 f"<p>The code is single-use and expires in {minutes} minutes. You can close this window"
959 " once the client confirms it is connected.</p>"
960 "</body></html>"
961 )
962 return HTMLResponse(body, headers=TOKEN_NO_CACHE_HEADERS)
965def _pkce_verifier_matches(code_verifier: str, code_challenge: str) -> bool:
966 """RFC 7636 S256 verification, total over hostile input. The comparison is over bytes
967 so a non-ASCII ``code_challenge`` (which reaches here unvalidated from the client's
968 authorize request) simply fails to match instead of raising ``TypeError`` the way
969 ``hmac.compare_digest`` does on two ``str`` with non-ASCII content. The verifier is
970 ASCII per spec; a compliant client's challenge is base64url and matches."""
971 digest: Final = hashlib.sha256(code_verifier.encode("ascii", "replace")).digest()
972 computed: Final = urlsafe_b64encode(digest).rstrip(b"=")
973 return hmac.compare_digest(computed, code_challenge.encode("utf-8"))
976ClaimOutcome = Literal["first", "replayed", "unavailable"]
978_CLAIM_UNAVAILABLE_DESCRIPTION: Final = "the single-use record is unavailable right now; try again shortly"
981def _claim_refusal(outcome: ClaimOutcome, replayed: Response) -> Response | None:
982 """A claim that is not the first caller's is refused, but the two reasons must stay apart on the
983 wire: a replay is the grant's own 4xx, while a shared backend that could not record the claim is
984 a 503 (RFC 7009 section 2.2.1, RFC 6749 section 5.2 ``temporarily_unavailable``), so the client
985 keeps the still-valid token and retries instead of being told it was already used."""
986 match outcome:
987 case "first":
988 return None
989 case "replayed":
990 return replayed
991 case "unavailable":
992 return _oauth_error(503, "temporarily_unavailable", _CLAIM_UNAVAILABLE_DESCRIPTION)
993 case _:
994 assert_never(outcome)
997class _SingleUseGuard:
998 """Atomic single-use claim for a one-time id (an auth-code, connect-flow ``jti``, or refresh-token
999 ``jti``) over the injected proxy cache.
1001 Uses an atomic increment rather than a get-then-set: two concurrent redemptions of the same id
1002 cannot both observe "unused", because exactly one increment returns 1. The claim IS the gate, so it
1003 fails closed. Crucially, the increment must be recorded in a backend SHARED across replicas, or the
1004 single-use property is per-worker only (each replica's in-memory counter returns 1, so a captured
1005 id replays through a different worker):
1007 - When a Redis backend is configured it is the SOLE authority: the claim goes straight to Redis
1008 (``INCR`` is atomic across replicas), and any Redis fault fails the claim CLOSED — it never falls
1009 back to the per-worker in-memory count (``DualCache.async_increment_cache`` does fall back, which
1010 is exactly the replay window this avoids).
1011 - With no Redis configured (single-replica) the in-memory increment is authoritative within the one
1012 process. A multi-worker deployment must run Redis for the guarantee to hold across workers.
1014 The id's own TTL is the outer bound. For the auth code, PKCE binding is the primary defense against
1015 interception; this makes the RFC 6749 4.1.2 single-use property reliable on top of it."""
1017 def __init__(self, cache: DualCache) -> None:
1018 self._cache = cache
1020 async def claim(self, key: str, ttl_seconds: int) -> ClaimOutcome:
1021 """Atomically claim ``key``. ``"first"`` iff this caller is the first (increment to 1),
1022 ``"replayed"`` on a replay (>1), and ``"unavailable"`` when the claim could not be recorded in
1023 the shared backend, which every caller treats as a refusal (fail closed)."""
1024 from litellm.proxy.proxy_server import redis_usage_cache # noqa: PLC0415 # circular import at module load
1026 # Resolve the shared authority HERE rather than trusting the injected cache: callers pass
1027 # user_api_key_cache, which only carries a redis_cache when enable_redis_auth_cache is set
1028 # (off by default), so a guard that read its injected cache silently degraded every claim to
1029 # a per-worker count on a stock multi-worker deployment. redis_usage_cache is the store the
1030 # proxy already treats as cross-worker, so no call site can wire the guarantee away.
1031 redis_cache: Final = redis_usage_cache or getattr(self._cache, "redis_cache", None)
1032 if redis_cache is not None:
1033 # Shared, atomic authority for multi-replica deployments. Claim ONLY against Redis and fail
1034 # CLOSED on any Redis fault (async_increment re-raises) rather than fall back to the
1035 # per-worker in-memory count, which would let each replica observe count==1 and replay the id.
1036 try:
1037 count = await redis_cache.async_increment(key, 1, ttl=ttl_seconds)
1038 except Exception as e: # noqa: BLE001 # ANY Redis fault fails the single-use claim closed
1039 verbose_logger.warning(
1040 "mcp gateway single-use claim: shared cache backend unavailable, failing closed: %s", e
1041 )
1042 return "unavailable"
1043 return "first" if count == 1 else "replayed"
1044 # No shared backend configured (single-replica): the in-memory increment is authoritative.
1045 count = await self._cache.async_increment_cache(key, 1, ttl=ttl_seconds, local_only=True)
1046 return "first" if count == 1 else "replayed"
1048 async def peek(self, key: str) -> Literal["unclaimed", "claimed", "unavailable"]:
1049 """Read-only view of a single-use marker, resolved against the same shared authority as
1050 :meth:`claim` so introspection observes exactly the record redemption and revocation wrote.
1051 A backend fault is ``"unavailable"`` (fail closed) rather than a guess either way."""
1052 from litellm.proxy.proxy_server import redis_usage_cache # noqa: PLC0415 # circular import at module load
1054 redis_cache: Final = redis_usage_cache or getattr(self._cache, "redis_cache", None)
1055 if redis_cache is not None:
1056 try:
1057 value = await redis_cache.async_get_cache(key)
1058 except Exception as e: # noqa: BLE001 # ANY Redis fault fails the read closed
1059 verbose_logger.warning("mcp gateway single-use peek: shared cache backend unavailable: %s", e)
1060 return "unavailable"
1061 return "unclaimed" if value is None else "claimed"
1062 local: Final = await self._cache.async_get_cache(key, local_only=True)
1063 return "unclaimed" if local is None else "claimed"
1066def _session_token_pair(principal: SessionPrincipal, keys: SessionSigningKeys, now: datetime) -> Response:
1067 access: Final = mint_session_token(principal, keys, now)
1068 refresh: Final = mint_session_refresh_token(principal, keys, now)
1069 if not isinstance(access, MintedSessionToken) or not isinstance(refresh, MintedSessionToken):
1070 return _oauth_error(500, "server_error", "failed to mint the session credential")
1071 return JSONResponse(
1072 status_code=200,
1073 content={
1074 "access_token": access.token.get_secret_value(),
1075 "token_type": "Bearer",
1076 "expires_in": int((access.expires_at - now).total_seconds()),
1077 "refresh_token": refresh.token.get_secret_value(),
1078 },
1079 headers=TOKEN_NO_CACHE_HEADERS,
1080 )
1083class _ProxyCredentialTokenResponse(TypedDict):
1084 access_token: ReadOnly[str]
1085 token_type: ReadOnly[Literal["Bearer"]]
1086 expires_in: ReadOnly[int]
1087 refresh_token: ReadOnly[str]
1088 user_id: ReadOnly[str]
1089 team_id: ReadOnly[str | None]
1090 issued_token_type: NotRequired[ReadOnly[_IssuedTokenType]]
1093def _proxy_credential_response(
1094 minted: MintedProxyCredential,
1095 principal: SessionPrincipal,
1096 keys: SessionSigningKeys,
1097 now: datetime,
1098 issued_token_type: _IssuedTokenType | None = None,
1099) -> Response:
1100 """The proxy-API token response: the access token is the very credential ``lite
1101 login`` stores (accepted on every proxy route with user and team attribution), and
1102 the refresh token is a gateway-sealed rotating token bound to the team the credential
1103 was minted for, so a renewal keeps the team the user consented to. A token exchange
1104 also states ``issued_token_type``, which RFC 8693 section 2.2.1 requires."""
1105 bound_principal: Final = principal.model_copy(update=MappingProxyType({"team_id": minted.team_id}))
1106 refresh: Final = mint_session_refresh_token(bound_principal, keys, now)
1107 if not isinstance(refresh, MintedSessionToken):
1108 return _oauth_error(500, "server_error", "failed to mint the session credential")
1109 credential: Final[_ProxyCredentialTokenResponse] = {
1110 "access_token": minted.key,
1111 "token_type": "Bearer",
1112 "expires_in": minted.expires_in,
1113 "refresh_token": refresh.token.get_secret_value(),
1114 "user_id": minted.user_id,
1115 "team_id": minted.team_id,
1116 }
1117 if issued_token_type is None:
1118 return JSONResponse(status_code=200, content=credential, headers=TOKEN_NO_CACHE_HEADERS)
1119 exchanged: Final[_ProxyCredentialTokenResponse] = {**credential, "issued_token_type": issued_token_type}
1120 return JSONResponse(status_code=200, content=exchanged, headers=TOKEN_NO_CACHE_HEADERS)
1123def _reload_failure_response(failure: ReloadUserFailure) -> Response:
1124 """Map the live-user revalidation failure onto its OAuth error, exhaustively, so a new
1125 ``ReloadUserFailure`` member is a type error here rather than silently 400ing."""
1126 match failure:
1127 case "unavailable":
1128 return _oauth_error(503, "temporarily_unavailable", _DB_UNAVAILABLE_DESCRIPTION)
1129 case "faulted":
1130 return _oauth_error(503, "temporarily_unavailable", _DB_FAULTED_DESCRIPTION)
1131 case "unresolvable":
1132 return _oauth_error(500, "server_error", "the gateway is not configured to resolve users")
1133 case "no_active_key":
1134 return _oauth_error(400, "invalid_grant", "the user for this grant is no longer active")
1135 case _:
1136 assert_never(failure)
1139def _subject_token_refusal_response(refusal: SubjectTokenRefusal) -> Response:
1140 match refusal.error:
1141 case "temporarily_unavailable":
1142 return _oauth_error(503, refusal.error, refusal.description)
1143 case "unsupported_grant_type" | "invalid_request":
1144 return _oauth_error(400, refusal.error, refusal.description)
1145 case _:
1146 assert_never(refusal.error)
1149def _mint_failure_response(failure: ProxyCredentialMintFailure) -> Response:
1150 match failure:
1151 case "not_a_member":
1152 return _oauth_error(
1153 400, "invalid_grant", "the user is no longer a member of the team this grant was issued for"
1154 )
1155 case "team_required":
1156 return _oauth_error(
1157 400, "invalid_grant", "this user belongs to a team; sign in again and pick the team for this credential"
1158 )
1159 case "unavailable" | "faulted" | "unresolvable" | "no_active_key":
1160 return _reload_failure_response(failure)
1161 case _:
1162 assert_never(failure)
1165def _resource_conflicts_with_scope(
1166 request: Request, resource: str | None, sealed_resource_server_id: str | None
1167) -> bool:
1168 """True when a scoped grant is being redeemed for a DIFFERENT resource than the one
1169 sealed into it (RFC 8707 section 2.2: reject with ``invalid_target``). An absent
1170 ``resource`` never conflicts (the sealed scope still binds the minted session), and an
1171 unscoped grant ignores the parameter entirely, exactly as the endpoint always has, so
1172 no pre-existing client breaks."""
1173 if sealed_resource_server_id is None or resource is None:
1174 return False
1175 resolved: Final = resolve_scoped_resource_server(request, resource)
1176 return resolved is None or resolved.server_id != sealed_resource_server_id
1179async def aggregate_token(
1180 request: Request,
1181 grant_type: str,
1182 code: str | None,
1183 redirect_uri: str | None,
1184 client_id: str,
1185 code_verifier: str | None,
1186 refresh_token: str | None,
1187 master_key: str | None,
1188 reload_user: ReloadUser,
1189 cache: DualCache,
1190 resource: str | None = None,
1191 mint_proxy_credential: MintProxyCredential = _refuse_proxy_credential,
1192 subject_token: str | None = None,
1193 subject_token_type: str | None = None,
1194 requested_token_type: str | None = None,
1195 exchange_subject_token: ExchangeSubjectToken = _refuse_subject_token,
1196) -> Response:
1197 """The aggregate token verb: authorization_code and refresh_token grants for the
1198 identity-only session pair, or for the proxy-API credential when the grant was issued
1199 with that audience, and the RFC 8693 token exchange that turns an IdP token straight
1200 into the proxy-API credential. Every path re-validates the litellm user live before
1201 minting, so a deactivated user cannot obtain or renew a session."""
1202 if master_key is None:
1203 verbose_logger.error("mcp_gateway_dcr token grant rejected: no master_key configured")
1204 return _oauth_error(500, "server_error", "the gateway has no master key configured")
1205 keys: Final = active_session_signing_keys(master_key)
1206 if isinstance(keys, SessionSigningConfigError):
1207 verbose_logger.error("mcp_gateway_dcr token grant rejected: %s", keys.detail)
1208 return _oauth_error(500, "server_error", "the gateway session signing configuration is invalid")
1209 now: Final = datetime.now(timezone.utc)
1210 issue: Final = _GrantIssuer(
1211 request=request,
1212 resource=resource,
1213 keys=keys,
1214 now=now,
1215 reload_user=reload_user,
1216 mint_proxy_credential=mint_proxy_credential,
1217 guard=_SingleUseGuard(cache),
1218 )
1219 if grant_type == "authorization_code":
1220 return await _authorization_code_grant(
1221 request=request,
1222 code=code,
1223 redirect_uri=redirect_uri,
1224 client_id=client_id,
1225 code_verifier=code_verifier,
1226 resource=resource,
1227 now=now,
1228 issue=issue,
1229 )
1230 if grant_type == "refresh_token":
1231 return await _refresh_token_grant(
1232 request=request,
1233 refresh_token=refresh_token,
1234 client_id=client_id,
1235 resource=resource,
1236 keys=keys,
1237 now=now,
1238 issue=issue,
1239 )
1240 if grant_type == TOKEN_EXCHANGE_GRANT_TYPE:
1241 return await _token_exchange_grant(
1242 subject_token=subject_token,
1243 subject_token_type=subject_token_type,
1244 requested_token_type=requested_token_type,
1245 client_id=client_id,
1246 exchange_subject_token=exchange_subject_token,
1247 issue=issue,
1248 )
1249 return _oauth_error(
1250 400,
1251 "unsupported_grant_type",
1252 f"grant_type must be authorization_code, refresh_token, or {TOKEN_EXCHANGE_GRANT_TYPE}",
1253 )
1256class _GrantIssuer:
1257 """The tail every grant shares once its own proof (code + PKCE, or a refresh token)
1258 has checked out: revalidate the user live, claim the single-use marker, mint. The
1259 claim comes AFTER revalidation and minting so a transient DB 503 never burns a
1260 still-valid code or refresh token, and fails closed when it cannot be recorded."""
1262 def __init__(
1263 self,
1264 request: Request,
1265 resource: str | None,
1266 keys: SessionSigningKeys,
1267 now: datetime,
1268 reload_user: ReloadUser,
1269 mint_proxy_credential: MintProxyCredential,
1270 guard: _SingleUseGuard,
1271 ) -> None:
1272 self._request: Final = request
1273 self._resource: Final = resource
1274 self._keys: Final = keys
1275 self._now: Final = now
1276 self._reload_user: Final = reload_user
1277 self._mint_proxy_credential: Final = mint_proxy_credential
1278 self._guard: Final = guard
1280 async def __call__(
1281 self, principal: SessionPrincipal, claim_key: str, claim_ttl_seconds: int, replayed: str
1282 ) -> Response:
1283 match principal.audience:
1284 case None:
1285 return await self._issue_session_pair(principal, claim_key, claim_ttl_seconds, replayed)
1286 case "proxy_api":
1287 return await self._issue_proxy_credential(principal, claim_key, claim_ttl_seconds, replayed)
1288 case _:
1289 assert_never(principal.audience)
1291 async def _issue_session_pair(
1292 self, principal: SessionPrincipal, claim_key: str, claim_ttl_seconds: int, replayed: str
1293 ) -> Response:
1294 failure: Final = await self._reload_user(principal.user_id)
1295 if failure is not None:
1296 return _reload_failure_response(failure)
1297 refusal: Final = await self._claim_refusal(claim_key, claim_ttl_seconds, replayed)
1298 if refusal is not None:
1299 return refusal
1300 return _session_token_pair(principal, self._keys, self._now)
1302 async def _issue_proxy_credential(
1303 self, principal: SessionPrincipal, claim_key: str, claim_ttl_seconds: int, replayed: str
1304 ) -> Response:
1305 target_refusal: Final = self._proxy_api_target_refusal()
1306 if target_refusal is not None:
1307 return target_refusal
1308 minted: Final = await self._mint_proxy_credential(principal.user_id, principal.team_id)
1309 if not isinstance(minted, MintedProxyCredential):
1310 return _mint_failure_response(minted)
1311 refusal: Final = await self._claim_refusal(claim_key, claim_ttl_seconds, replayed)
1312 if refusal is not None:
1313 return refusal
1314 return _proxy_credential_response(minted, principal, self._keys, self._now)
1316 async def exchange(
1317 self, subject_token: str, client_id: str, exchange_subject_token: ExchangeSubjectToken
1318 ) -> Response:
1319 """The RFC 8693 tail: prove the IdP token, then mint. No single-use marker, because
1320 the subject token stays a valid proof for as long as the IdP says it is and every
1321 exchange mints a fresh credential and refresh token of its own."""
1322 target_refusal: Final = self._proxy_api_target_refusal()
1323 if target_refusal is not None:
1324 return target_refusal
1325 identity: Final = await exchange_subject_token(subject_token, self._request)
1326 if isinstance(identity, SubjectTokenRefusal):
1327 return _subject_token_refusal_response(identity)
1328 principal: Final = SessionPrincipal(
1329 user_id=identity.user_id, client_id=client_id, audience=PROXY_API_AUDIENCE, team_id=identity.team_id
1330 )
1331 minted: Final = await self._mint_proxy_credential(principal.user_id, principal.team_id)
1332 if not isinstance(minted, MintedProxyCredential):
1333 return _mint_failure_response(minted)
1334 return _proxy_credential_response(
1335 minted, principal, self._keys, self._now, issued_token_type=ACCESS_TOKEN_TOKEN_TYPE
1336 )
1338 def _proxy_api_target_refusal(self) -> Response | None:
1339 if self._resource is None or is_proxy_api_resource(self._request, self._resource):
1340 return None
1341 return _oauth_error(400, "invalid_target", "resource does not match the proxy API this grant was issued for")
1343 async def _claim_refusal(self, claim_key: str, claim_ttl_seconds: int, replayed: str) -> Response | None:
1344 return _claim_refusal(
1345 await self._guard.claim(claim_key, claim_ttl_seconds), replayed=_oauth_error(400, "invalid_grant", replayed)
1346 )
1349async def _authorization_code_grant(
1350 request: Request,
1351 code: str | None,
1352 redirect_uri: str | None,
1353 client_id: str,
1354 code_verifier: str | None,
1355 resource: str | None,
1356 now: datetime,
1357 issue: _GrantIssuer,
1358) -> Response:
1359 if not code or not redirect_uri or not code_verifier:
1360 return _oauth_error(400, "invalid_request", "code, redirect_uri, and code_verifier are required")
1361 if not MIN_CODE_VERIFIER_LENGTH <= len(code_verifier) <= MAX_CODE_VERIFIER_LENGTH:
1362 return _oauth_error(400, "invalid_request", "code_verifier must be 43 to 128 characters (RFC 7636)")
1363 parsed: Final = _open_sealed(code, GATEWAY_AUTH_CODE_PREFIX, _GatewayAuthCode, _AUTH_CODE_DEBUG_KEY)
1364 if parsed is None:
1365 return _oauth_error(400, "invalid_grant", "the authorization code is invalid")
1366 if now.timestamp() >= parsed.exp:
1367 return _oauth_error(400, "invalid_grant", "the authorization code has expired")
1368 if client_id != parsed.client_id or redirect_uri != parsed.redirect_uri:
1369 return _oauth_error(400, "invalid_grant", "the authorization code was issued to a different client")
1370 if _resource_conflicts_with_scope(request, resource, parsed.resource_server_id):
1371 return _oauth_error(400, "invalid_target", "resource does not match the scope this code was issued for")
1372 if not _pkce_verifier_matches(code_verifier, parsed.code_challenge):
1373 return _oauth_error(400, "invalid_grant", "PKCE verification failed")
1374 # The marker's TTL derives from the code's own remaining lifetime so it outlives
1375 # whichever lifetime the code was minted with.
1376 return await issue(
1377 SessionPrincipal(
1378 user_id=parsed.user_id,
1379 client_id=client_id,
1380 resource_server_id=parsed.resource_server_id,
1381 audience=parsed.audience,
1382 team_id=parsed.team_id,
1383 ),
1384 claim_key=f"{_USED_CODE_CACHE_PREFIX}{parsed.jti}",
1385 claim_ttl_seconds=parsed.exp - int(now.timestamp()) + _CLAIM_TTL_BUFFER_SECONDS,
1386 replayed="the authorization code was already used",
1387 )
1390async def _refresh_token_grant(
1391 request: Request,
1392 refresh_token: str | None,
1393 client_id: str,
1394 resource: str | None,
1395 keys: SessionSigningKeys,
1396 now: datetime,
1397 issue: _GrantIssuer,
1398) -> Response:
1399 if not refresh_token:
1400 return _oauth_error(400, "invalid_request", "refresh_token is required")
1401 opened: Final = open_session_refresh_bearer(refresh_token, keys, now, expected_client_id=client_id)
1402 if not isinstance(opened, SessionRefreshOpened):
1403 return _oauth_error(400, "invalid_grant", "the refresh token is invalid for this client")
1404 if _resource_conflicts_with_scope(request, resource, opened.principal.resource_server_id):
1405 return _oauth_error(400, "invalid_target", "resource does not match the scope this token was issued for")
1406 # Refresh-token rotation (OAuth 2.0 Security BCP section 4.13): the presented refresh token is
1407 # single-use, so a captured or replayed refresh token cannot mint a second pair after the
1408 # legitimate holder rotated.
1409 return await issue(
1410 opened.principal,
1411 claim_key=f"{_USED_REFRESH_CACHE_PREFIX}{opened.jti}",
1412 claim_ttl_seconds=SESSION_REFRESH_TTL_SECONDS + _CLAIM_TTL_BUFFER_SECONDS,
1413 replayed="the refresh token was already used",
1414 )
1417async def _token_exchange_grant(
1418 subject_token: str | None,
1419 subject_token_type: str | None,
1420 requested_token_type: str | None,
1421 client_id: str,
1422 exchange_subject_token: ExchangeSubjectToken,
1423 issue: _GrantIssuer,
1424) -> Response:
1425 """RFC 8693 token exchange for a registered native client that already holds an IdP
1426 token: the gateway proves the token the way its JWT auth does and answers with the
1427 proxy-API credential, so a fresh laptop with only an IdP login gets a gateway key
1428 without a browser round trip. The client must be registered because the refresh token
1429 in the answer is bound to it."""
1430 if not is_gateway_dcr_client_id(client_id) or open_gateway_dcr_client(client_id) is None:
1431 return _oauth_error(401, "invalid_client", "unknown or malformed client_id")
1432 if not subject_token or not subject_token_type:
1433 return _oauth_error(400, "invalid_request", "subject_token and subject_token_type are required")
1434 if subject_token_type not in SUBJECT_TOKEN_TYPES:
1435 return _oauth_error(
1436 400, "invalid_request", f"subject_token_type must be one of {', '.join(sorted(SUBJECT_TOKEN_TYPES))}"
1437 )
1438 if requested_token_type is not None and requested_token_type != ACCESS_TOKEN_TOKEN_TYPE:
1439 return _oauth_error(400, "invalid_request", f"requested_token_type must be {ACCESS_TOKEN_TOKEN_TYPE}")
1440 return await issue.exchange(subject_token, client_id, exchange_subject_token)
1443async def revoke_refresh_token(token: str, client_id: str, master_key: str | None, cache: DualCache) -> Response:
1444 """RFC 7009 revocation for the gateway's refresh tokens: burn the presented token's
1445 ``jti`` so neither the holder nor a thief can rotate it again. Access tokens are
1446 stateless and expire on their own (the proxy-API credential within
1447 ``CLI_JWT_EXPIRATION_HOURS``), so per RFC 7009 section 2.2 an unrecognized or already
1448 dead token still answers 200; only an unknown client is refused. A live token whose
1449 burn could not be recorded in the shared backend answers 503 (section 2.2.1), so the
1450 client knows the token still stands and retries instead of reporting a logout that
1451 never happened."""
1452 if not is_gateway_dcr_client_id(client_id) or open_gateway_dcr_client(client_id) is None: 1452 ↛ 1454line 1452 didn't jump to line 1454 because the condition on line 1452 was always true
1453 return _oauth_error(401, "invalid_client", "unknown or malformed client_id")
1454 if master_key is None:
1455 verbose_logger.error("mcp_gateway_dcr revoke rejected: no master_key configured")
1456 return _oauth_error(500, "server_error", "the gateway has no master key configured")
1457 keys: Final = active_session_signing_keys(master_key)
1458 if isinstance(keys, SessionSigningConfigError):
1459 verbose_logger.error("mcp_gateway_dcr revoke rejected: %s", keys.detail)
1460 return _oauth_error(500, "server_error", "the gateway session signing configuration is invalid")
1461 now: Final = datetime.now(timezone.utc)
1462 opened: Final = open_session_refresh_bearer(token, keys, now, expected_client_id=client_id)
1463 if isinstance(opened, SessionRefreshOpened):
1464 burned: Final = await _SingleUseGuard(cache).claim(
1465 f"{_USED_REFRESH_CACHE_PREFIX}{opened.jti}", SESSION_REFRESH_TTL_SECONDS + _CLAIM_TTL_BUFFER_SECONDS
1466 )
1467 if burned == "unavailable":
1468 return _oauth_error(503, "temporarily_unavailable", _CLAIM_UNAVAILABLE_DESCRIPTION)
1469 return Response(content="{}", media_type="application/json", headers=TOKEN_NO_CACHE_HEADERS)
1472def _inactive_introspection_response() -> Response:
1473 """RFC 7662 section 2.2: any token the gateway cannot vouch for, whatever the reason
1474 (wrong family, bad signature, expired, revoked, or a deactivated user), answers 200
1475 with ``active: false`` and nothing else, so introspection is not a token oracle."""
1476 return JSONResponse(status_code=200, content={"active": False}, headers=TOKEN_NO_CACHE_HEADERS)
1479def _active_introspection_response(opened: OpenedSessionToken) -> Response:
1480 principal: Final = opened.principal
1481 optional_claims: Final = {
1482 key: value
1483 for key, value in (
1484 ("token_type", "Bearer" if opened.kind == "session" else None),
1485 ("team_id", principal.team_id),
1486 ("resource_server_id", principal.resource_server_id),
1487 ("audience", principal.audience),
1488 )
1489 if value is not None
1490 }
1491 return JSONResponse(
1492 status_code=200,
1493 content={
1494 "active": True,
1495 "iss": SESSION_ISSUER,
1496 "sub": principal.user_id,
1497 "client_id": principal.client_id,
1498 "jti": opened.jti,
1499 "iat": opened.iat,
1500 "exp": opened.exp,
1501 "kind": opened.kind,
1502 **optional_claims,
1503 },
1504 headers=TOKEN_NO_CACHE_HEADERS,
1505 )
1508async def introspect_gateway_token(
1509 token: str,
1510 master_key: str | None,
1511 reload_user: ReloadUser,
1512 cache: DualCache,
1513) -> Response:
1514 """RFC 7662 introspection for the gateway's session tokens, so an external gateway
1515 (Kong, an API management layer) can validate a LiteLLM-issued MCP session credential
1516 without holding the signing secret. The caller is already authenticated by the route
1517 (section 2.1). Active means everything admission itself would require: valid signature
1518 under the configured session signing keys, unexpired, not a revoked or rotated refresh
1519 token, and a litellm user that is still live, so a deactivated user's outstanding
1520 tokens introspect as inactive immediately. A shared-backend or DB outage answers 503
1521 rather than guessing in either direction."""
1522 if master_key is None: 1522 ↛ 1523line 1522 didn't jump to line 1523 because the condition on line 1522 was never true
1523 verbose_logger.error("mcp_gateway_dcr introspect rejected: no master_key configured")
1524 return _oauth_error(500, "server_error", "the gateway has no master key configured")
1525 keys: Final = active_session_signing_keys(master_key)
1526 if isinstance(keys, SessionSigningConfigError): 1526 ↛ 1527line 1526 didn't jump to line 1527 because the condition on line 1526 was never true
1527 verbose_logger.error("mcp_gateway_dcr introspect rejected: %s", keys.detail)
1528 return _oauth_error(500, "server_error", keys.detail)
1529 now: Final = datetime.now(timezone.utc)
1530 if is_session_token(token): 1530 ↛ 1531line 1530 didn't jump to line 1531 because the condition on line 1530 was never true
1531 opened = open_session_token(token, keys, now)
1532 elif is_session_refresh_token(token): 1532 ↛ 1533line 1532 didn't jump to line 1533 because the condition on line 1532 was never true
1533 opened = open_session_refresh_token(token, keys, now)
1534 else:
1535 return _inactive_introspection_response()
1536 if not isinstance(opened, OpenedSessionToken):
1537 return _inactive_introspection_response()
1538 if opened.kind == "session_refresh":
1539 peeked: Final = await _SingleUseGuard(cache).peek(f"{_USED_REFRESH_CACHE_PREFIX}{opened.jti}")
1540 if peeked == "unavailable":
1541 return _oauth_error(503, "temporarily_unavailable", _CLAIM_UNAVAILABLE_DESCRIPTION)
1542 if peeked == "claimed":
1543 return _inactive_introspection_response()
1544 failure: Final = await reload_user(opened.principal.user_id)
1545 if failure == "unavailable" or failure == "faulted":
1546 return _reload_failure_response(failure)
1547 if failure is not None:
1548 return _inactive_introspection_response()
1549 return _active_introspection_response(opened)