Module adcp.oauth
Security-focused OAuth 2.0 authorization-code helpers for buyer clients.
This module deliberately implements a small, explicit surface:
- RFC 8414 authorization-server discovery from a trusted issuer URL;
- pre-registered public clients (no dynamic registration or client secrets);
- authorization code with RFC 7636 S256 PKCE;
- one-time, atomically consumed pending flows; and
- bounded, DNS-pinned metadata and token HTTP requests.
It does not discover an authorization server from an MCP resource URL.
A
resource-to-issuer flow first needs RFC 9728 protected-resource metadata; pass
the resulting trusted issuer URL to :func:discover_oauth_metadata().
Functions
-
Expand source code
async def complete_oauth_authorization( *, code: str | None, callback_state: str, expected_state: str, store: PendingOAuthFlowStore, callback_issuer: str | None = None, callback_error: str | None = None, timeout: float = 10.0, clock: Callable[[], datetime] | None = None, ) -> OAuthTokenSet: """Consume one callback and exchange its code for public-client tokens. ``expected_state`` must come from a separate browser-session binding, not from the callback query itself. Once the two states match, the pending flow is consumed before any issuer/error/code handling or network I/O. Any failure therefore requires starting a new authorization flow. """ if isinstance(timeout, bool) or not math.isfinite(timeout) or timeout <= 0 or timeout > 60: raise ValueError("timeout must be greater than 0 and at most 60 seconds") if not ( 43 <= len(callback_state) <= _MAX_STATE_LENGTH and 43 <= len(expected_state) <= _MAX_STATE_LENGTH and callback_state.isascii() and expected_state.isascii() and re.fullmatch(r"[A-Za-z0-9_-]+", callback_state) and re.fullmatch(r"[A-Za-z0-9_-]+", expected_state) ) or not hmac.compare_digest(callback_state, expected_state): raise OAuthAuthorizationError("state_mismatch", phase="callback") pending = await store.consume(callback_state) if pending is None: raise OAuthAuthorizationError("flow_not_found", phase="callback") if not hmac.compare_digest(pending.state, callback_state): raise OAuthAuthorizationError("state_mismatch", phase="callback") now = (clock or (lambda: datetime.now(timezone.utc)))() if now.tzinfo is None: raise ValueError("clock must return a timezone-aware datetime") if pending.expires_at <= now: raise OAuthAuthorizationError("flow_expired", phase="callback") if callback_issuer is not None and ( len(callback_issuer) > _MAX_URL_LENGTH or callback_issuer != pending.issuer ): raise OAuthAuthorizationError("issuer_mismatch", phase="callback") if ( pending.issuer_binding is OAuthIssuerBinding.AUTHORIZATION_RESPONSE_ISS and callback_issuer is None ): raise OAuthAuthorizationError("issuer_mismatch", phase="callback") if callback_error is not None: raise OAuthAuthorizationError( "authorization_failed", phase="callback", oauth_error=_safe_oauth_error(callback_error), ) if code is None or not 1 <= len(code) <= _MAX_CODE_LENGTH or not _VSCHAR_RE.fullmatch(code): raise OAuthAuthorizationError("invalid_code", phase="callback") return await _exchange_code( pending, code=code, timeout=timeout, )Consume one callback and exchange its code for public-client tokens.
expected_statemust come from a separate browser-session binding, not from the callback query itself. Once the two states match, the pending flow is consumed before any issuer/error/code handling or network I/O. Any failure therefore requires starting a new authorization flow. async def discover_oauth_metadata(issuer_url: str, *, timeout: float = 5.0, allow_loopback_http: bool = False) ‑> OAuthAuthorizationServerMetadata-
Expand source code
async def discover_oauth_metadata( issuer_url: str, *, timeout: float = 5.0, allow_loopback_http: bool = False, ) -> OAuthAuthorizationServerMetadata: """Discover and validate RFC 8414 metadata for an exact issuer URL. ``issuer_url`` is an authorization-server issuer, not an MCP/HTTP resource URL. The returned ``issuer`` must match it exactly. Plain HTTP is available only for literal loopback IPs when ``allow_loopback_http=True``. """ if isinstance(timeout, bool) or not math.isfinite(timeout) or timeout <= 0 or timeout > 60: raise ValueError("timeout must be greater than 0 and at most 60 seconds") discovery_url = _metadata_url(issuer_url, allow_loopback_http=allow_loopback_http) async def fetch() -> Mapping[str, Any]: transport = await _transport_for( discovery_url, allow_loopback_http=allow_loopback_http, ) async with httpx.AsyncClient( transport=transport, timeout=timeout, follow_redirects=False, trust_env=False, ) as client: async with client.stream( "GET", discovery_url, headers={"accept": "application/json", "accept-encoding": "identity"}, ) as response: if response.status_code != 200: raise OAuthDiscoveryError( "http_error", phase="discovery", status_code=response.status_code ) return await _read_json_response(response, phase="discovery") try: raw = await _run_with_deadline(fetch, timeout) except _OAuthDeadlineExpired: raise OAuthDiscoveryError("timeout", phase="discovery") from None except OAuthDiscoveryError: raise except OAuthClientError as exc: raise OAuthDiscoveryError( exc.code, phase="discovery", status_code=exc.status_code ) from None except (httpx.HTTPError, OSError, SSRFValidationError, ValueError): raise OAuthDiscoveryError("network_error", phase="discovery") from None try: metadata = OAuthAuthorizationServerMetadata.model_validate(raw) except ValidationError: raise OAuthDiscoveryError("invalid_metadata", phase="discovery") from None if metadata.issuer != issuer_url: raise OAuthDiscoveryError("issuer_mismatch", phase="discovery") _validate_public_client_metadata(metadata, allow_loopback_http=allow_loopback_http) return metadataDiscover and validate RFC 8414 metadata for an exact issuer URL.
issuer_urlis an authorization-server issuer, not an MCP/HTTP resource URL. The returnedissuermust match it exactly. Plain HTTP is available only for literal loopback IPs whenallow_loopback_http=True. -
Expand source code
async def start_oauth_authorization( metadata: OAuthAuthorizationServerMetadata, *, client_id: str, redirect_uri: str, store: PendingOAuthFlowStore, issuer_binding: OAuthIssuerBinding, scopes: Sequence[str] = (), resource: str | None = None, ttl_seconds: int = _DEFAULT_FLOW_TTL_SECONDS, clock: Callable[[], datetime] | None = None, allow_loopback_http: bool = False, ) -> OAuthAuthorizationRequest: """Create, persist, and return a browser authorization request. ``DISTINCT_REDIRECT_URI`` is an explicit assertion that this redirect URI is not shared with any other issuer. Prefer ``AUTHORIZATION_RESPONSE_ISS`` when the server advertises RFC 9207 support. """ _validate_public_client_metadata(metadata, allow_loopback_http=allow_loopback_http) if not 1 <= len(client_id) <= _MAX_CLIENT_ID_LENGTH or not _VSCHAR_RE.fullmatch(client_id): raise ValueError("client_id must use bounded RFC 6749 VSCHAR syntax") if isinstance(ttl_seconds, bool) or not 1 <= ttl_seconds <= 3600: raise ValueError("ttl_seconds must be between 1 and 3600") if not isinstance(issuer_binding, OAuthIssuerBinding): raise ValueError("issuer_binding must be an OAuthIssuerBinding value") try: _validate_absolute_url( redirect_uri, field="redirect_uri", allow_loopback_http=allow_loopback_http, ) except ValueError: policy = "HTTPS or a literal loopback HTTP address" if allow_loopback_http else "HTTPS" raise ValueError(f"redirect_uri must use {policy}") from None if ( issuer_binding is OAuthIssuerBinding.AUTHORIZATION_RESPONSE_ISS and not metadata.authorization_response_iss_parameter_supported ): raise ValueError("authorization server does not advertise RFC 9207 issuer responses") scope = _validate_scope(scopes) resource = _validate_resource(resource) now = (clock or (lambda: datetime.now(timezone.utc)))() if now.tzinfo is None: raise ValueError("clock must return a timezone-aware datetime") expires_at = now + timedelta(seconds=ttl_seconds) state = _random_urlsafe(32) verifier = _random_urlsafe(64) challenge = _pkce_challenge(verifier) pending = PendingOAuthAuthorization( state=state, code_verifier=SecretStr(verifier), issuer=metadata.issuer, authorization_endpoint=metadata.authorization_endpoint, token_endpoint=metadata.token_endpoint, client_id=client_id, redirect_uri=redirect_uri, scope=scope, resource=resource, issuer_binding=issuer_binding, allow_loopback_http=allow_loopback_http, created_at=now, expires_at=expires_at, ) if not await store.insert_if_absent(pending): raise OAuthFlowStoreError("state_collision", phase="store") parts = urlsplit(metadata.authorization_endpoint) query = parse_qsl(parts.query, keep_blank_values=True) query.extend( [ ("response_type", "code"), ("client_id", client_id), ("redirect_uri", redirect_uri), ("state", state), ("code_challenge", challenge), ("code_challenge_method", "S256"), ] ) if scope is not None: query.append(("scope", scope)) if resource is not None: query.append(("resource", resource)) authorization_url = urlunsplit((parts.scheme, parts.netloc, parts.path, urlencode(query), "")) if len(authorization_url) > _MAX_URL_LENGTH: # The flow remains stored and must not be reused. Consume it best-effort. await store.consume(state) raise ValueError("authorization URL exceeds the accepted length") return OAuthAuthorizationRequest( authorization_url=authorization_url, state=state, expires_at=expires_at, )Create, persist, and return a browser authorization request.
DISTINCT_REDIRECT_URIis an explicit assertion that this redirect URI is not shared with any other issuer. PreferAUTHORIZATION_RESPONSE_ISSwhen the server advertises RFC 9207 support.
Classes
class InMemoryPendingOAuthFlowStore (*, capacity: int = 1024, clock: Callable[[], datetime] | None = None)-
Expand source code
class InMemoryPendingOAuthFlowStore: """Lock-protected reference store for one process only. Multi-process and multi-replica applications must provide a shared store whose insert-if-absent and consume operations are atomic across replicas. """ def __init__( self, *, capacity: int = _DEFAULT_STORE_CAPACITY, clock: Callable[[], datetime] | None = None, ) -> None: if isinstance(capacity, bool) or not 1 <= capacity <= 1_000_000: raise ValueError("capacity must be between 1 and 1000000") self._capacity = capacity self._clock = clock or (lambda: datetime.now(timezone.utc)) self._entries: dict[str, PendingOAuthAuthorization] = {} self._lock = asyncio.Lock() def _now(self) -> datetime: now = self._clock() if now.tzinfo is None: raise ValueError("OAuth store clock must return a timezone-aware datetime") return now def _prune_expired(self, now: datetime) -> None: expired = [state for state, pending in self._entries.items() if pending.expires_at <= now] for state in expired: del self._entries[state] async def insert_if_absent(self, pending: PendingOAuthAuthorization) -> bool: async with self._lock: now = self._now() self._prune_expired(now) if pending.expires_at <= now or pending.state in self._entries: return False if len(self._entries) >= self._capacity: raise OAuthFlowStoreError("capacity_exceeded", phase="store") self._entries[pending.state] = pending return True async def consume(self, state: str) -> PendingOAuthAuthorization | None: async with self._lock: now = self._now() self._prune_expired(now) return self._entries.pop(state, None)Lock-protected reference store for one process only.
Multi-process and multi-replica applications must provide a shared store whose insert-if-absent and consume operations are atomic across replicas.
Methods
async def consume(self, state: str) ‑> PendingOAuthAuthorization | None-
Expand source code
async def consume(self, state: str) -> PendingOAuthAuthorization | None: async with self._lock: now = self._now() self._prune_expired(now) return self._entries.pop(state, None) async def insert_if_absent(self,
pending: PendingOAuthAuthorization) ‑> bool-
Expand source code
async def insert_if_absent(self, pending: PendingOAuthAuthorization) -> bool: async with self._lock: now = self._now() self._prune_expired(now) if pending.expires_at <= now or pending.state in self._entries: return False if len(self._entries) >= self._capacity: raise OAuthFlowStoreError("capacity_exceeded", phase="store") self._entries[pending.state] = pending return True
class OAuthAuthorizationError (code: str,
*,
phase: str,
status_code: int | None = None,
oauth_error: str | None = None,
retry_after_seconds: float | None = None)-
Expand source code
class OAuthAuthorizationError(OAuthClientError): """Authorization callback validation failed; start a new flow."""Authorization callback validation failed; start a new flow.
Ancestors
- OAuthClientError
- builtins.Exception
- builtins.BaseException
class OAuthAuthorizationRequest (**data: Any)-
Expand source code
class OAuthAuthorizationRequest(BaseModel): """Browser-facing result. The PKCE verifier is intentionally absent.""" model_config = ConfigDict(extra="forbid", frozen=True) authorization_url: BoundedUrl state: BoundedState expires_at: datetime @field_validator("expires_at") @classmethod def _aware_expiry(cls, value: datetime) -> datetime: if value.tzinfo is None: raise ValueError("expires_at must be timezone-aware") return valueBrowser-facing result. The PKCE verifier is intentionally absent.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- pydantic.main.BaseModel
Class variables
var expires_at : datetime.datetimevar model_configvar state : str
class OAuthAuthorizationServerMetadata (**data: Any)-
Expand source code
class OAuthAuthorizationServerMetadata(BaseModel): """Bounded subset of RFC 8414 metadata used by this helper.""" model_config = ConfigDict(extra="ignore", frozen=True) issuer: BoundedUrl authorization_endpoint: BoundedUrl token_endpoint: BoundedUrl response_types_supported: tuple[ResponseType, ...] = Field(min_length=1, max_length=32) grant_types_supported: tuple[BoundedToken, ...] | None = Field( default=None, min_length=1, max_length=32 ) token_endpoint_auth_methods_supported: tuple[BoundedToken, ...] = Field( min_length=1, max_length=32 ) code_challenge_methods_supported: tuple[BoundedToken, ...] = Field(min_length=1, max_length=16) scopes_supported: tuple[ScopeToken, ...] | None = Field( default=None, min_length=1, max_length=128 ) authorization_response_iss_parameter_supported: bool = False @field_validator("issuer") @classmethod def _issuer_url(cls, value: str) -> str: return _validate_absolute_url( value, field="issuer", issuer=True, allow_loopback_http=True, ) @field_validator("authorization_endpoint", "token_endpoint") @classmethod def _endpoint_url(cls, value: str) -> str: return _validate_absolute_url( value, field="OAuth endpoint", allow_loopback_http=True, ) @field_validator("authorization_endpoint") @classmethod def _reserved_authorization_query(cls, value: str) -> str: query_keys = [key for key, _ in parse_qsl(urlsplit(value).query, keep_blank_values=True)] if any(key in _RESERVED_AUTHORIZATION_QUERY_KEYS for key in query_keys): raise ValueError("authorization endpoint contains a reserved OAuth query parameter") if len(query_keys) != len(set(query_keys)): raise ValueError("authorization endpoint contains duplicate query parameters") return valueBounded subset of RFC 8414 metadata used by this helper.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- pydantic.main.BaseModel
Class variables
var code_challenge_methods_supported : tuple[str, ...]var grant_types_supported : tuple[str, ...] | Nonevar issuer : strvar model_configvar response_types_supported : tuple[str, ...]var scopes_supported : tuple[str, ...] | Nonevar token_endpoint : strvar token_endpoint_auth_methods_supported : tuple[str, ...]
class OAuthClientError (code: str,
*,
phase: str,
status_code: int | None = None,
oauth_error: str | None = None,
retry_after_seconds: float | None = None)-
Expand source code
class OAuthClientError(Exception): """Base error whose message never contains remote response content.""" def __init__( self, code: str, *, phase: str, status_code: int | None = None, oauth_error: str | None = None, retry_after_seconds: float | None = None, ) -> None: self.code = code self.phase = phase self.status_code = status_code self.oauth_error = oauth_error self.retry_after_seconds = retry_after_seconds message = f"OAuth {phase} failed ({code})" if status_code is not None: message += f": HTTP {status_code}" super().__init__(message)Base error whose message never contains remote response content.
Ancestors
- builtins.Exception
- builtins.BaseException
Subclasses
class OAuthDiscoveryError (code: str,
*,
phase: str,
status_code: int | None = None,
oauth_error: str | None = None,
retry_after_seconds: float | None = None)-
Expand source code
class OAuthDiscoveryError(OAuthClientError): """Authorization-server discovery failed safely."""Authorization-server discovery failed safely.
Ancestors
- OAuthClientError
- builtins.Exception
- builtins.BaseException
class OAuthFlowStoreError (code: str,
*,
phase: str,
status_code: int | None = None,
oauth_error: str | None = None,
retry_after_seconds: float | None = None)-
Expand source code
class OAuthFlowStoreError(OAuthClientError): """The pending-flow store rejected an insert or lookup."""The pending-flow store rejected an insert or lookup.
Ancestors
- OAuthClientError
- builtins.Exception
- builtins.BaseException
class OAuthIssuerBinding (*args, **kwds)-
Expand source code
class OAuthIssuerBinding(str, Enum): """How a redirect URI is bound to one authorization-server issuer.""" AUTHORIZATION_RESPONSE_ISS = "authorization_response_iss" DISTINCT_REDIRECT_URI = "distinct_redirect_uri"How a redirect URI is bound to one authorization-server issuer.
Ancestors
- builtins.str
- enum.Enum
Class variables
var AUTHORIZATION_RESPONSE_ISSvar DISTINCT_REDIRECT_URI
class OAuthTokenExchangeError (code: str,
*,
phase: str,
status_code: int | None = None,
oauth_error: str | None = None,
retry_after_seconds: float | None = None)-
Expand source code
class OAuthTokenExchangeError(OAuthClientError): """Token exchange failed; the consumed flow must not be retried."""Token exchange failed; the consumed flow must not be retried.
Ancestors
- OAuthClientError
- builtins.Exception
- builtins.BaseException
class OAuthTokenSet (**data: Any)-
Expand source code
class OAuthTokenSet(BaseModel): """Closed, secret-safe token response returned after exchange.""" model_config = ConfigDict(extra="ignore", frozen=True) access_token: SecretStr token_type: Annotated[ str, Field(min_length=1, max_length=64, pattern=_TOKEN_TYPE_RE.pattern), ] expires_in: Annotated[int | None, Field(default=None, ge=0, le=315_360_000)] refresh_token: SecretStr | None = None id_token: SecretStr | None = None scope: Annotated[str | None, Field(default=None, max_length=_MAX_SCOPE_LENGTH)] @field_validator("access_token", "refresh_token", "id_token") @classmethod def _bounded_secret(cls, value: SecretStr | None) -> SecretStr | None: if value is not None: secret = value.get_secret_value() if not 1 <= len(secret) <= _MAX_TOKEN_LENGTH or not _VSCHAR_RE.fullmatch(secret): raise ValueError("token must use bounded RFC 6749 VSCHAR syntax") return value @field_validator("scope") @classmethod def _valid_scope(cls, value: str | None) -> str | None: if value is not None and not _SCOPE_RE.fullmatch(value): raise ValueError("scope must contain space-delimited RFC 6749 scope tokens") return valueClosed, secret-safe token response returned after exchange.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- pydantic.main.BaseModel
Class variables
var access_token : pydantic.types.SecretStrvar expires_in : int | Nonevar id_token : pydantic.types.SecretStr | Nonevar model_configvar refresh_token : pydantic.types.SecretStr | Nonevar scope : str | Nonevar token_type : str
class PendingOAuthAuthorization (**data: Any)-
Expand source code
class PendingOAuthAuthorization(BaseModel): """Authoritative server-side snapshot for one authorization attempt.""" model_config = ConfigDict(extra="forbid", frozen=True) state: BoundedState code_verifier: SecretStr issuer: BoundedUrl authorization_endpoint: BoundedUrl token_endpoint: BoundedUrl client_id: Annotated[str, Field(min_length=1, max_length=_MAX_CLIENT_ID_LENGTH)] redirect_uri: BoundedUrl scope: Annotated[str | None, Field(max_length=_MAX_SCOPE_LENGTH)] = None resource: Annotated[str | None, Field(max_length=_MAX_RESOURCE_LENGTH)] = None issuer_binding: OAuthIssuerBinding allow_loopback_http: bool = False created_at: datetime expires_at: datetime @field_validator("code_verifier") @classmethod def _valid_verifier(cls, value: SecretStr) -> SecretStr: if not _PKCE_VERIFIER_RE.fullmatch(value.get_secret_value()): raise ValueError("invalid PKCE verifier") return value @model_validator(mode="after") def _valid_lifetime(self) -> PendingOAuthAuthorization: if self.created_at.tzinfo is None or self.expires_at.tzinfo is None: raise ValueError("OAuth flow timestamps must be timezone-aware") if self.expires_at <= self.created_at: raise ValueError("OAuth flow expiry must follow creation") return selfAuthoritative server-side snapshot for one authorization attempt.
Create a new model by parsing and validating input data from keyword arguments.
Raises [
ValidationError][pydantic_core.ValidationError] if the input data cannot be validated to form a valid model.selfis explicitly positional-only to allowselfas a field name.Ancestors
- pydantic.main.BaseModel
Class variables
var allow_loopback_http : boolvar client_id : strvar code_verifier : pydantic.types.SecretStrvar created_at : datetime.datetimevar expires_at : datetime.datetimevar issuer : strvar issuer_binding : OAuthIssuerBindingvar model_configvar redirect_uri : strvar resource : str | Nonevar scope : str | Nonevar state : strvar token_endpoint : str
class PendingOAuthFlowStore (*args, **kwargs)-
Expand source code
@runtime_checkable class PendingOAuthFlowStore(Protocol): """Atomic one-time storage contract for pending OAuth flows.""" async def insert_if_absent(self, pending: PendingOAuthAuthorization) -> bool: """Insert a live state without overwriting; return ``False`` on collision.""" async def consume(self, state: str) -> PendingOAuthAuthorization | None: """Atomically remove and return a live state, or return ``None``."""Atomic one-time storage contract for pending OAuth flows.
Ancestors
- typing.Protocol
- typing.Generic
Methods
async def consume(self, state: str) ‑> PendingOAuthAuthorization | None-
Expand source code
async def consume(self, state: str) -> PendingOAuthAuthorization | None: """Atomically remove and return a live state, or return ``None``."""Atomically remove and return a live state, or return
None. async def insert_if_absent(self,
pending: PendingOAuthAuthorization) ‑> bool-
Expand source code
async def insert_if_absent(self, pending: PendingOAuthAuthorization) -> bool: """Insert a live state without overwriting; return ``False`` on collision."""Insert a live state without overwriting; return
Falseon collision.