#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
pub mod private_key_jwt;
pub mod rpc;
pub mod subscriptions;
#[cfg(feature = "tasks")]
pub mod tasks;
mod acquisition;
mod lifetime;
use lifetime::unless_revoked;
use std::fmt;
use std::future::{Future, poll_fn};
use std::io::{self, Write};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::task::Poll;
use std::time::{Duration, Instant};
use asupersync::Cx;
use asupersync::http::h1::{HttpClient, RedirectPolicy, Request, RetryPolicy};
use asupersync::sync::Mutex;
use asupersync::tls::Certificate;
use asupersync::types::Time;
use fastmcp_core::{AccessToken, CanonicalHttpUrl, McpRequestCancellation};
use fastmcp_protocol::protocol_policy::ProtocolEra;
use fastmcp_protocol::{
CoreRequest, CoreResult, FINAL_CLIENT_CAPABILITIES_META_KEY, FINAL_PROTOCOL_VERSION,
FinalCoreResult, RequestId, decode_strict_jsonrpc_response,
};
use serde::Deserialize;
use serde_json::{Value, json};
use super::{
OAuthDiscoveryError, OAuthDiscoveryPlan, TrustedOAuthIssuer, admit_root, admit_scopes,
check_context, decode_metadata, discovery_deadline, has, present, validate_array,
validate_optional_array, within,
};
use crate::http_auth::BoundBearerCredential;
use crate::http_executor::{
ModernHttpExecutor, ModernHttpRequest, ModernHttpResponseKind, ModernHttpResponseMetadata,
ModernHttpResponseStream, ModernHttpSseResponseStream, ResourceTlsTrust,
};
use crate::sse::SseLimits;
pub const CLIENT_CREDENTIALS_EXTENSION: &str = "io.modelcontextprotocol/oauth-client-credentials";
const MAX_TOKEN_BYTES: usize = 64 * 1024;
const MAX_SECRET_BYTES: usize = 4096;
const MAX_SECRET_GRANT_BYTES: usize = 64 * 1024;
const MAX_REQUEST_BYTES: usize = 8 * 1024 * 1024;
const MAX_ACQUISITIONS: usize = 64;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ClientSecretAuthenticationMethod {
Basic,
Post,
}
#[derive(Debug)]
pub enum ClientCredentialsError {
InvalidPolicy,
InvalidRegistration,
AssertionSigning,
Closed,
Expired,
Saturated,
StateUnavailable,
GenerationExhausted,
ConcurrentAcquisitionFailed,
UnsupportedAuthentication,
InvalidToken,
ExpandedScope,
TokenEndpointRejected,
Transport,
InvalidRequest,
RequestTooLarge,
Negotiation,
UnexpectedResponse,
Redirect {
status: u16,
},
Discovery(OAuthDiscoveryError),
}
impl fmt::Display for ClientCredentialsError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(match self {
Self::InvalidPolicy => "invalid client-credentials policy",
Self::InvalidRegistration => "private-key client registration is invalid, expired or mismatched",
Self::AssertionSigning => "private-key client assertion signing failed",
Self::Closed => "client-credentials owner is closed",
Self::Expired => "client-credentials access token is expired or revoked",
Self::Saturated => "client-credentials acquisition capacity exhausted",
Self::StateUnavailable => "client-credentials state unavailable",
Self::GenerationExhausted => "client-credentials generation exhausted",
Self::ConcurrentAcquisitionFailed => "the joined client-credentials acquisition failed or was abandoned",
Self::UnsupportedAuthentication => "issuer does not admit the selected client-credentials authentication",
Self::InvalidToken => "client-credentials token response rejected",
Self::ExpandedScope => "client-credentials token exceeds requested scopes",
Self::TokenEndpointRejected => "client-credentials token endpoint rejected the grant",
Self::Transport => "client-credentials transport failed",
Self::InvalidRequest => "invalid client-credentials MCP request",
Self::RequestTooLarge => "client-credentials MCP request exceeds its byte bound",
Self::Negotiation => "client-credentials discovery or extension declaration rejected",
Self::UnexpectedResponse => "client-credentials response rejected",
Self::Redirect { .. } => "resource answered the client-credentials request with a redirect, which MCP does not follow",
Self::Discovery(_) => "client-credentials discovery or operation lifetime failed",
})
}
}
impl std::error::Error for ClientCredentialsError {}
impl From<OAuthDiscoveryError> for ClientCredentialsError {
fn from(error: OAuthDiscoveryError) -> Self {
Self::Discovery(error)
}
}
struct ClientSecret(String);
#[derive(Clone)]
enum MachineAuthentication {
Basic(Arc<ClientSecret>),
Post(Arc<ClientSecret>),
#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
PrivateKeyJwt(Arc<private_key_jwt::PrivateKeyJwtAuthentication>),
}
struct PreparedMachineGrant {
body: Vec<u8>,
authorization: Option<String>,
deadline: Time,
}
impl MachineAuthentication {
fn method(&self) -> &'static str {
match self {
Self::Basic(_) => "client_secret_basic",
Self::Post(_) => "client_secret_post",
#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
Self::PrivateKeyJwt(_) => "private_key_jwt",
}
}
#[allow(
clippy::unnecessary_wraps,
reason = "the PrivateKeyJwt arm, compiled on non-wasm targets with builtin-auth-server, \
returns the assertion check's error; the Result is only unnecessary without it"
)]
fn check(&self) -> Result<(), ClientCredentialsError> {
match self {
Self::Basic(_) | Self::Post(_) => Ok(()),
#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
Self::PrivateKeyJwt(authentication) => authentication.check(),
}
}
fn token_expiry_limit(&self) -> Option<Instant> {
match self {
Self::Basic(_) | Self::Post(_) => None,
#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
Self::PrivateKeyJwt(authentication) => Some(authentication.valid_until()),
}
}
#[allow(
clippy::unused_async,
clippy::unused_async_trait_impl,
reason = "the PrivateKeyJwt arm, compiled on non-wasm targets with builtin-auth-server, \
awaits the signed assertion; the function has no await only without it"
)]
async fn prepare(
&self,
cx: &Cx,
deadline: Time,
client_id: &str,
resource: &CanonicalHttpUrl,
scopes: &[String],
) -> Result<PreparedMachineGrant, ClientCredentialsError> {
self.check()?;
check_context(cx, deadline)?;
match self {
Self::Basic(secret) => prepare_secret_grant(
ClientSecretAuthenticationMethod::Basic,
client_id,
&secret.0,
resource,
scopes,
deadline,
),
Self::Post(secret) => prepare_secret_grant(
ClientSecretAuthenticationMethod::Post,
client_id,
&secret.0,
resource,
scopes,
deadline,
),
#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
Self::PrivateKeyJwt(authentication) => {
authentication
.prepare(cx, deadline, client_id, resource, scopes)
.await
}
}
}
}
fn prepare_secret_grant(
method: ClientSecretAuthenticationMethod,
client_id: &str,
secret: &str,
resource: &CanonicalHttpUrl,
scopes: &[String],
deadline: Time,
) -> Result<PreparedMachineGrant, ClientCredentialsError> {
let scope = scopes.join(" ");
let mut fields = vec![
("grant_type", "client_credentials"),
("resource", resource.as_str()),
("scope", scope.as_str()),
];
let authorization = match method {
ClientSecretAuthenticationMethod::Basic => Some(basic(client_id, secret)?),
ClientSecretAuthenticationMethod::Post => {
fields.push(("client_id", client_id));
fields.push(("client_secret", secret));
None
}
};
let body = form(&fields).into_bytes();
if body.len() > MAX_SECRET_GRANT_BYTES {
return Err(ClientCredentialsError::RequestTooLarge);
}
Ok(PreparedMachineGrant {
body,
authorization,
deadline,
})
}
pub struct ClientCredentialsPlan {
discovery: OAuthDiscoveryPlan,
authentication: MachineAuthentication,
maximum_lifetime: Duration,
leeway: Duration,
}
impl ClientCredentialsPlan {
pub fn new(
resource: CanonicalHttpUrl,
issuer: TrustedOAuthIssuer,
client_id: impl Into<String>,
client_secret: impl Into<String>,
scopes: Vec<String>,
) -> Result<Self, ClientCredentialsError> {
let client_secret = client_secret.into();
if client_secret.is_empty()
|| client_secret.len() > MAX_SECRET_BYTES
|| client_secret.chars().any(char::is_control)
{
return Err(ClientCredentialsError::InvalidPolicy);
}
Ok(Self {
discovery: OAuthDiscoveryPlan::new(resource, vec![issuer], client_id, scopes)?,
authentication: MachineAuthentication::Basic(Arc::new(ClientSecret(client_secret))),
maximum_lifetime: Duration::from_secs(3600),
leeway: Duration::from_secs(30),
})
}
pub fn with_secret_authentication(
mut self,
method: ClientSecretAuthenticationMethod,
) -> Result<Self, ClientCredentialsError> {
let secret = match self.authentication {
MachineAuthentication::Basic(secret) | MachineAuthentication::Post(secret) => secret,
#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
MachineAuthentication::PrivateKeyJwt(_) => {
return Err(ClientCredentialsError::InvalidPolicy);
}
};
self.authentication = match method {
ClientSecretAuthenticationMethod::Basic => MachineAuthentication::Basic(secret),
ClientSecretAuthenticationMethod::Post => MachineAuthentication::Post(secret),
};
Ok(self)
}
pub fn with_timeout(mut self, timeout: Duration) -> Result<Self, ClientCredentialsError> {
self.discovery = self.discovery.with_timeout(timeout)?;
Ok(self)
}
pub fn with_maximum_token_lifetime(
mut self,
lifetime: Duration,
) -> Result<Self, ClientCredentialsError> {
if lifetime.is_zero() || lifetime > Duration::from_hours(24) {
return Err(ClientCredentialsError::InvalidPolicy);
}
self.maximum_lifetime = lifetime;
Ok(self)
}
pub fn with_renewal_leeway(mut self, leeway: Duration) -> Result<Self, ClientCredentialsError> {
if leeway > Duration::from_secs(300) {
return Err(ClientCredentialsError::InvalidPolicy);
}
self.leeway = leeway;
Ok(self)
}
pub fn with_resource_root_certificate(
mut self,
root: Certificate,
) -> Result<Self, ClientCredentialsError> {
admit_root(&mut self.discovery.resource_roots, root)?;
Ok(self)
}
pub async fn discover(
&self,
cx: &Cx,
) -> Result<ClientCredentialsClient, ClientCredentialsError> {
self.authentication.check()?;
let mut resource_tls = None;
for root in &self.discovery.resource_roots {
ResourceTlsTrust::add_root(
&mut resource_tls,
self.discovery.resource.clone(),
root.clone(),
)
.map_err(|_| ClientCredentialsError::InvalidPolicy)?;
}
let deadline = discovery_deadline(cx, self.discovery.timeout)?;
let (issuer, body) = self
.discovery
.discover_issuer_document(cx, deadline)
.await?;
let token_endpoint =
admit_machine_issuer(&self.discovery, issuer, &body, &self.authentication)?;
check_context(cx, deadline)?;
self.authentication.check()?;
Ok(ClientCredentialsClient {
inner: Arc::new(ClientInner {
resource: self.discovery.resource.clone(),
token_endpoint,
client_id: self
.discovery
.client_id
.clone()
.ok_or(ClientCredentialsError::InvalidPolicy)?,
scopes: self.discovery.scopes.clone(),
authentication: self.authentication.clone(),
issuer_roots: issuer.roots.clone(),
resource_tls,
timeout: self.discovery.timeout,
maximum_lifetime: self.maximum_lifetime,
leeway: self.leeway,
closed: McpRequestCancellation::new(),
pending: AtomicUsize::new(0),
state: Arc::new(Mutex::new(TokenState::default())),
}),
})
}
}
impl fmt::Debug for ClientCredentialsPlan {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ClientCredentialsPlan")
.field("authentication", &self.authentication.method())
.field("credential", &"<redacted>")
.finish_non_exhaustive()
}
}
#[derive(Deserialize)]
struct MachineIssuerMetadata {
issuer: String,
token_endpoint: String,
grant_types_supported: Vec<String>,
token_endpoint_auth_methods_supported: Vec<String>,
#[serde(default, deserialize_with = "present")]
scopes_supported: Option<Vec<String>>,
#[serde(default, deserialize_with = "present")]
protected_resources: Option<Vec<String>>,
#[serde(default, deserialize_with = "present")]
signed_metadata: Option<String>,
}
fn admit_machine_issuer(
plan: &OAuthDiscoveryPlan,
issuer: &TrustedOAuthIssuer,
body: &[u8],
authentication: &MachineAuthentication,
) -> Result<CanonicalHttpUrl, ClientCredentialsError> {
authentication.check()?;
let metadata: MachineIssuerMetadata = decode_metadata(body)?;
if metadata.issuer != issuer.identifier {
return Err(OAuthDiscoveryError::IssuerMismatch.into());
}
if metadata.signed_metadata.is_some() {
return Err(OAuthDiscoveryError::SignedMetadataUnsupported.into());
}
validate_array(&metadata.grant_types_supported)?;
validate_array(&metadata.token_endpoint_auth_methods_supported)?;
validate_optional_array(metadata.protected_resources.as_deref())?;
if !has(&metadata.grant_types_supported, "client_credentials")
|| !has(
&metadata.token_endpoint_auth_methods_supported,
authentication.method(),
)
{
return Err(ClientCredentialsError::UnsupportedAuthentication);
}
if metadata
.protected_resources
.as_ref()
.is_some_and(|values| !has(values, plan.resource.as_str()))
{
return Err(OAuthDiscoveryError::ResourceMismatch.into());
}
admit_scopes(&plan.scopes, metadata.scopes_supported.as_deref())?;
#[cfg(all(not(target_arch = "wasm32"), feature = "builtin-auth-server"))]
if let MachineAuthentication::PrivateKeyJwt(authentication) = authentication {
authentication.admit_metadata(&metadata, body)?;
}
Ok(issuer.endpoint(&metadata.token_endpoint)?)
}
struct ClientInner {
resource: CanonicalHttpUrl,
token_endpoint: CanonicalHttpUrl,
client_id: String,
scopes: Vec<String>,
authentication: MachineAuthentication,
issuer_roots: Vec<Certificate>,
resource_tls: Option<ResourceTlsTrust>,
timeout: Duration,
maximum_lifetime: Duration,
leeway: Duration,
closed: McpRequestCancellation,
pending: AtomicUsize,
state: Arc<Mutex<TokenState>>,
}
impl Drop for ClientInner {
fn drop(&mut self) {
self.closed.cancel();
}
}
#[derive(Default)]
struct TokenState {
current: Option<ServiceToken>,
generation: u64,
flight: Option<Arc<acquisition::GrantFlight>>,
}
struct ServiceToken {
bearer: BoundBearerCredential,
scopes: Vec<String>,
expires_at: Instant,
renew_after: Instant,
}
#[derive(Clone)]
pub struct ClientCredentialsClient {
inner: Arc<ClientInner>,
}
impl fmt::Debug for ClientCredentialsClient {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ClientCredentialsClient")
.field("closed", &self.inner.closed.is_cancel_requested())
.finish_non_exhaustive()
}
}
pub struct ClientCredentialsSnapshot {
bearer: BoundBearerCredential,
scopes: Vec<String>,
expires_at: Instant,
generation: u64,
}
impl ClientCredentialsSnapshot {
pub fn credential(&self) -> &BoundBearerCredential {
&self.bearer
}
pub fn scopes(&self) -> &[String] {
&self.scopes
}
pub fn expires_at(&self) -> Instant {
self.expires_at
}
pub fn generation(&self) -> u64 {
self.generation
}
}
impl fmt::Debug for ClientCredentialsSnapshot {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ClientCredentialsSnapshot")
.field("generation", &self.generation)
.field("credential", &"<redacted>")
.finish_non_exhaustive()
}
}
impl ClientCredentialsClient {
pub fn resource(&self) -> &CanonicalHttpUrl {
&self.inner.resource
}
fn resource_http_executor(&self) -> ModernHttpExecutor {
ModernHttpExecutor::new().with_resource_tls(self.inner.resource_tls.clone())
}
pub fn close(&self) {
self.inner.closed.cancel();
if let Ok(mut state) = self.inner.state.try_lock_owned() {
state.current = None;
state.flight = None;
}
}
pub async fn credential(
&self,
cx: &Cx,
) -> Result<ClientCredentialsSnapshot, ClientCredentialsError> {
self.credential_with_cancellation(cx, &McpRequestCancellation::new())
.await
}
pub async fn credential_with_cancellation(
&self,
cx: &Cx,
cancellation: &McpRequestCancellation,
) -> Result<ClientCredentialsSnapshot, ClientCredentialsError> {
acquisition::credential(self, cx, cancellation).await
}
pub async fn execute_core(
&self,
cx: &Cx,
request: CoreRequest,
discovery_id: RequestId,
request_id: RequestId,
) -> Result<ClientCredentialsResponse, ClientCredentialsError> {
Box::pin(self.execute_core_with_cancellation(
cx,
&McpRequestCancellation::new(),
request,
discovery_id,
request_id,
))
.await
}
pub async fn execute_core_with_cancellation(
&self,
cx: &Cx,
cancellation: &McpRequestCancellation,
request: CoreRequest,
discovery_id: RequestId,
request_id: RequestId,
) -> Result<ClientCredentialsResponse, ClientCredentialsError> {
if discovery_id.correlates_with(&request_id) {
return Err(ClientCredentialsError::InvalidRequest);
}
let (wire, decoder) = prepare(self.resource(), &request, &request_id)?;
let params = decoder
.encode_params()
.map_err(|_| ClientCredentialsError::InvalidRequest)?
.ok_or(ClientCredentialsError::InvalidRequest)?;
let discovery = CoreRequest::decode(
ProtocolEra::Modern2026,
"server/discover",
Some(&json!({"_meta":params["_meta"]})),
)
.map_err(|_| ClientCredentialsError::InvalidRequest)?;
let (discovery_wire, discovery) = prepare(self.resource(), &discovery, &discovery_id)?;
let deadline = discovery_deadline(cx, self.inner.timeout)?;
Box::pin(active(
cx,
deadline,
&self.inner.closed,
cancellation,
None,
async {
let snapshot = self.credential_with_cancellation(cx, cancellation).await?;
let executor = self.resource_http_executor();
let discovery_wire = authorize(&snapshot, discovery_wire)?;
let response = active(
cx,
deadline,
&self.inner.closed,
cancellation,
Some(&snapshot),
async {
executor
.execute_with_cancellation(cx, cancellation, &discovery_wire)
.await
.map_err(|_| ClientCredentialsError::Transport)
},
)
.await?;
if response.metadata().status() != 200
|| response.metadata().kind() != ModernHttpResponseKind::Json
{
return Err(ClientCredentialsError::Negotiation);
}
let bytes = active(
cx,
deadline,
&self.inner.closed,
cancellation,
Some(&snapshot),
async {
response
.read_to_end_with_cancellation(cx, cancellation, MAX_TOKEN_BYTES)
.await
.map_err(|_| ClientCredentialsError::Transport)
},
)
.await?;
admit_resource(&discovery, &discovery_id, &bytes)?;
check_context(cx, deadline)?;
let wire = authorize(&snapshot, wire)?;
let response = active(
cx,
deadline,
&self.inner.closed,
cancellation,
Some(&snapshot),
async {
executor
.execute_with_cancellation(cx, cancellation, &wire)
.await
.map_err(|_| ClientCredentialsError::Transport)
},
)
.await?;
Ok(ClientCredentialsResponse {
response,
snapshot,
owner: self.inner.closed.clone(),
cancellation: cancellation.clone(),
request: decoder,
request_id,
deadline,
})
},
))
.await
}
}
struct AcquisitionPermit<'a>(&'a AtomicUsize);
impl<'a> AcquisitionPermit<'a> {
fn new(pending: &'a AtomicUsize) -> Result<Self, ClientCredentialsError> {
pending
.try_update(Ordering::AcqRel, Ordering::Acquire, |n| {
(n < MAX_ACQUISITIONS).then(|| n + 1)
})
.map_err(|_| ClientCredentialsError::Saturated)?;
Ok(Self(pending))
}
}
impl Drop for AcquisitionPermit<'_> {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::AcqRel);
}
}
fn token_transport(roots: &[Certificate]) -> HttpClient {
let mut builder = HttpClient::builder()
.redirect_policy(RedirectPolicy::None)
.retry_policy(RetryPolicy::None)
.no_proxy()
.no_cookie_store()
.max_body_size(MAX_TOKEN_BYTES)
.max_total_connections(1);
for root in roots {
builder = builder.add_root_certificate(root.clone());
}
builder.build()
}
#[derive(Deserialize)]
struct TokenDocument {
access_token: String,
token_type: String,
#[serde(default, deserialize_with = "present")]
expires_in: Option<u64>,
#[serde(default, deserialize_with = "present")]
scope: Option<String>,
#[serde(default, deserialize_with = "present")]
resource: Option<String>,
#[serde(default, deserialize_with = "present")]
error: Option<String>,
#[serde(
default,
rename = "refresh_token",
deserialize_with = "reject_token_member"
)]
_refresh_token: (),
#[serde(default, rename = "id_token", deserialize_with = "reject_token_member")]
_id_token: (),
#[serde(default, rename = "cnf", deserialize_with = "reject_token_member")]
_confirmation: (),
#[serde(
default,
rename = "issued_token_type",
deserialize_with = "reject_token_member"
)]
_issued_token_type: (),
#[serde(
default,
rename = "error_description",
deserialize_with = "reject_token_member"
)]
_error_description: (),
#[serde(
default,
rename = "error_uri",
deserialize_with = "reject_token_member"
)]
_error_uri: (),
}
fn reject_token_member<'de, D: serde::Deserializer<'de>>(_: D) -> Result<(), D::Error> {
Err(serde::de::Error::custom(
"unsupported machine token response member",
))
}
fn admit_token(
inner: &ClientInner,
bytes: &[u8],
started: Instant,
) -> Result<ServiceToken, ClientCredentialsError> {
if bytes.len() > MAX_TOKEN_BYTES {
return Err(ClientCredentialsError::InvalidToken);
}
let token: TokenDocument =
decode_metadata(bytes).map_err(|_| ClientCredentialsError::InvalidToken)?;
if !token.token_type.eq_ignore_ascii_case("Bearer")
|| token.access_token.len() > 16 * 1024
|| !AccessToken::is_valid_token68(&token.access_token)
|| token.error.is_some()
|| token
.resource
.as_ref()
.is_some_and(|resource| resource != inner.resource.as_str())
{
return Err(ClientCredentialsError::InvalidToken);
}
let scopes = token.scope.map_or_else(
|| inner.scopes.clone(),
|scope| scope.split(' ').map(str::to_owned).collect(),
);
if scopes.len() > 32
|| scopes.iter().enumerate().any(|(index, scope)| {
scope.is_empty() || !inner.scopes.contains(scope) || scopes[..index].contains(scope)
})
{
return Err(ClientCredentialsError::ExpandedScope);
}
let lifetime = token.expires_in.map_or(inner.maximum_lifetime, |seconds| {
Duration::from_secs(seconds).min(inner.maximum_lifetime)
});
let expires_at = started
.checked_add(lifetime)
.ok_or(ClientCredentialsError::InvalidToken)?;
let expires_at = inner
.authentication
.token_expiry_limit()
.map_or(expires_at, |limit| expires_at.min(limit));
if Instant::now() >= expires_at {
return Err(ClientCredentialsError::Expired);
}
let bearer = BoundBearerCredential::bind_with_expiry(
inner.resource.clone(),
token.access_token,
expires_at,
)
.map_err(|_| ClientCredentialsError::InvalidToken)?
.for_owner(&inner.closed)
.ok_or(ClientCredentialsError::StateUnavailable)?;
let remaining = expires_at.saturating_duration_since(Instant::now());
let renew_after = expires_at
.checked_sub(inner.leeway.min(remaining / 2))
.ok_or(ClientCredentialsError::InvalidToken)?;
Ok(ServiceToken {
bearer,
scopes,
expires_at,
renew_after,
})
}
fn component(value: &str) -> String {
let mut encoded = String::new();
const HEX: &[u8; 16] = b"0123456789ABCDEF";
for byte in value.bytes() {
match byte {
b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => {
encoded.push(char::from(byte));
}
b' ' => encoded.push('+'),
_ => {
encoded.push('%');
encoded.push(char::from(HEX[usize::from(byte >> 4)]));
encoded.push(char::from(HEX[usize::from(byte & 15)]));
}
}
}
encoded
}
fn form(fields: &[(&str, &str)]) -> String {
fields
.iter()
.filter(|(name, value)| *name != "scope" || !value.is_empty())
.map(|(name, value)| format!("{}={}", component(name), component(value)))
.collect::<Vec<_>>()
.join("&")
}
fn basic(id: &str, secret: &str) -> Result<String, ClientCredentialsError> {
Request::post("/")
.basic_auth(component(id), Some(&component(secret)))
.build()
.headers
.into_iter()
.find(|(name, _)| name.eq_ignore_ascii_case("authorization"))
.map(|(_, value)| value)
.ok_or(ClientCredentialsError::StateUnavailable)
}
fn prepare(
resource: &CanonicalHttpUrl,
request: &CoreRequest,
id: &RequestId,
) -> Result<(ModernHttpRequest, CoreRequest), ClientCredentialsError> {
if request.era() != ProtocolEra::Modern2026
|| !matches!(
request.method(),
"server/discover"
| "tools/list"
| "tools/call"
| "resources/list"
| "resources/templates/list"
| "resources/read"
| "prompts/list"
| "prompts/get"
| "completion/complete"
)
{
return Err(ClientCredentialsError::InvalidRequest);
}
id.validate()
.map_err(|_| ClientCredentialsError::InvalidRequest)?;
let mut params = request
.encode_params()
.map_err(|_| ClientCredentialsError::InvalidRequest)?
.ok_or(ClientCredentialsError::InvalidRequest)?;
let capabilities = params
.get_mut("_meta")
.and_then(|meta| meta.get_mut(FINAL_CLIENT_CAPABILITIES_META_KEY))
.and_then(Value::as_object_mut)
.ok_or(ClientCredentialsError::InvalidRequest)?;
if let Some(extensions) = capabilities.get("extensions") {
let extensions = extensions
.as_object()
.ok_or(ClientCredentialsError::InvalidRequest)?;
if extensions.iter().any(|(name, settings)| {
name != CLIENT_CREDENTIALS_EXTENSION
|| !settings.as_object().is_some_and(serde_json::Map::is_empty)
}) {
return Err(ClientCredentialsError::InvalidRequest);
}
}
capabilities.insert(
"extensions".to_owned(),
json!({CLIENT_CREDENTIALS_EXTENSION:{}}),
);
let decoder = CoreRequest::decode(ProtocolEra::Modern2026, request.method(), Some(¶ms))
.map_err(|_| ClientCredentialsError::InvalidRequest)?;
let name = if matches!(request.method(), "tools/call" | "prompts/get") {
params
.get("name")
.and_then(Value::as_str)
.map(str::to_owned)
} else {
None
};
let envelope = json!({"jsonrpc":"2.0","id":id,"method":request.method(),"params":params});
let mut body = BoundedBody(Vec::new());
serde_json::to_writer(&mut body, &envelope)
.map_err(|_| ClientCredentialsError::RequestTooLarge)?;
let wire = ModernHttpRequest::new(
resource.as_str(),
body.0,
FINAL_PROTOCOL_VERSION,
request.method(),
name,
)
.map_err(|_| ClientCredentialsError::InvalidRequest)?;
Ok((wire, decoder))
}
struct BoundedBody(Vec<u8>);
impl Write for BoundedBody {
fn write(&mut self, input: &[u8]) -> io::Result<usize> {
if input.len() > MAX_REQUEST_BYTES.saturating_sub(self.0.len()) {
return Err(io::Error::other("client-credentials request limit"));
}
self.0.extend_from_slice(input);
Ok(input.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
fn authorize(
snapshot: &ClientCredentialsSnapshot,
wire: ModernHttpRequest,
) -> Result<ModernHttpRequest, ClientCredentialsError> {
check_token(&snapshot.bearer, snapshot.expires_at)?;
let wire = wire.with_authorization(&snapshot.bearer);
if !wire
.headers()
.iter()
.any(|(name, _)| name.eq_ignore_ascii_case("authorization"))
{
return Err(ClientCredentialsError::Expired);
}
Ok(wire)
}
fn decoded_result(
request: &CoreRequest,
id: &RequestId,
bytes: &[u8],
maximum: usize,
) -> Result<CoreResult, ClientCredentialsError> {
let admission = decode_strict_jsonrpc_response(bytes, maximum)
.map_err(|_| ClientCredentialsError::UnexpectedResponse)?;
let (response, source) = admission.into_parts();
if response.error.is_some()
|| !response
.id
.as_ref()
.is_some_and(|actual| actual.correlates_with(id))
{
return Err(ClientCredentialsError::UnexpectedResponse);
}
if response
.result
.as_ref()
.and_then(|result| result.get("resultType"))
.and_then(Value::as_str)
.is_some_and(|kind| !matches!(kind, "complete" | "input_required"))
{
return Err(ClientCredentialsError::UnexpectedResponse);
}
request
.decode_response_result(
&response,
&source.ok_or(ClientCredentialsError::UnexpectedResponse)?,
)
.map_err(|_| ClientCredentialsError::UnexpectedResponse)
}
fn admit_resource(
request: &CoreRequest,
id: &RequestId,
bytes: &[u8],
) -> Result<(), ClientCredentialsError> {
let params = request
.encode_params()
.map_err(|_| ClientCredentialsError::Negotiation)?
.ok_or(ClientCredentialsError::Negotiation)?;
if request.era() != ProtocolEra::Modern2026
|| request.method() != "server/discover"
|| !params
.get("_meta")
.and_then(|meta| meta.get(FINAL_CLIENT_CAPABILITIES_META_KEY))
.and_then(|caps| caps.get("extensions"))
.and_then(|ext| ext.get(CLIENT_CREDENTIALS_EXTENSION))
.is_some_and(|settings| settings.as_object().is_some_and(serde_json::Map::is_empty))
{
return Err(ClientCredentialsError::Negotiation);
}
let result = decoded_result(request, id, bytes, MAX_TOKEN_BYTES)
.map_err(|_| ClientCredentialsError::Negotiation)?;
let CoreResult::Final(FinalCoreResult::Discover(discovery)) = &result else {
return Err(ClientCredentialsError::Negotiation);
};
if !discovery
.supported_versions()
.iter()
.any(|version| version == FINAL_PROTOCOL_VERSION)
{
return Err(ClientCredentialsError::Negotiation);
}
let document: Value =
serde_json::from_slice(bytes).map_err(|_| ClientCredentialsError::Negotiation)?;
let capabilities = document
.get("result")
.and_then(|result| result.get("capabilities"))
.and_then(Value::as_object)
.ok_or(ClientCredentialsError::Negotiation)?;
if let Some(extensions) = capabilities.get("extensions") {
let extensions = extensions
.as_object()
.ok_or(ClientCredentialsError::Negotiation)?;
if extensions
.get(CLIENT_CREDENTIALS_EXTENSION)
.is_some_and(|settings| !settings.as_object().is_some_and(serde_json::Map::is_empty))
{
return Err(ClientCredentialsError::Negotiation);
}
}
Ok(())
}
fn check_token(
token: &BoundBearerCredential,
expiry: Instant,
) -> Result<(), ClientCredentialsError> {
if token.is_revoked() || Instant::now() >= expiry {
return Err(ClientCredentialsError::Expired);
}
Ok(())
}
async fn active<T>(
cx: &Cx,
deadline: Time,
owner: &McpRequestCancellation,
cancellation: &McpRequestCancellation,
token: Option<&ClientCredentialsSnapshot>,
future: impl Future<Output = Result<T, ClientCredentialsError>>,
) -> Result<T, ClientCredentialsError> {
let expiry_deadline = token
.map(|token| {
check_token(&token.bearer, token.expires_at)?;
discovery_deadline(
cx,
token.expires_at.saturating_duration_since(Instant::now()),
)
.map_err(ClientCredentialsError::from)
})
.transpose()?;
let deadline = expiry_deadline.map_or(deadline, |expiry| deadline.min(expiry));
let mut stopped = std::pin::pin!(owner.cancelled());
let mut cancelled = std::pin::pin!(cancellation.cancelled());
let mut future = std::pin::pin!(unless_revoked(
token.map(|snapshot| &snapshot.bearer.revoked),
future
));
within(cx, deadline, async {
Ok(poll_fn(|task| {
if owner.is_cancel_requested() || stopped.as_mut().poll(task).is_ready() {
return Poll::Ready(Err(ClientCredentialsError::Closed));
}
if cancellation.is_cancel_requested() || cancelled.as_mut().poll(task).is_ready() {
return Poll::Ready(Err(OAuthDiscoveryError::Cancelled.into()));
}
if let Some(token) = token {
check_token(&token.bearer, token.expires_at)?;
}
let result = future.as_mut().poll(task);
if owner.is_cancel_requested() {
return Poll::Ready(Err(ClientCredentialsError::Closed));
}
if cancellation.is_cancel_requested() {
return Poll::Ready(Err(OAuthDiscoveryError::Cancelled.into()));
}
if let Some(token) = token {
check_token(&token.bearer, token.expires_at)?;
}
result
})
.await)
})
.await?
}
pub struct ClientCredentialsResponse {
response: ModernHttpResponseStream,
snapshot: ClientCredentialsSnapshot,
owner: McpRequestCancellation,
cancellation: McpRequestCancellation,
request: CoreRequest,
request_id: RequestId,
deadline: Time,
}
impl ClientCredentialsResponse {
pub fn metadata(&self) -> &ModernHttpResponseMetadata {
self.response.metadata()
}
pub fn credential_generation(&self) -> u64 {
self.snapshot.generation
}
pub fn request(&self) -> &CoreRequest {
&self.request
}
pub fn request_id(&self) -> &RequestId {
&self.request_id
}
pub async fn read_to_end(
self,
cx: &Cx,
maximum_bytes: usize,
) -> Result<Vec<u8>, ClientCredentialsError> {
active(
cx,
self.deadline,
&self.owner,
&self.cancellation,
Some(&self.snapshot),
async {
self.response
.read_to_end_with_cancellation(cx, &self.cancellation, maximum_bytes)
.await
.map_err(|_| ClientCredentialsError::UnexpectedResponse)
},
)
.await
}
pub async fn read_json_result(
self,
cx: &Cx,
maximum_bytes: usize,
) -> Result<CoreResult, ClientCredentialsError> {
if self.metadata().status() != 200 || self.metadata().kind() != ModernHttpResponseKind::Json
{
return Err(ClientCredentialsError::UnexpectedResponse);
}
let Self {
response,
snapshot,
owner,
cancellation,
request,
request_id,
deadline,
} = self;
active(
cx,
deadline,
&owner,
&cancellation,
Some(&snapshot),
async {
let bytes = response
.read_to_end_with_cancellation(cx, &cancellation, maximum_bytes)
.await
.map_err(|_| ClientCredentialsError::UnexpectedResponse)?;
decoded_result(&request, &request_id, &bytes, maximum_bytes)
},
)
.await
}
pub fn into_sse_stream(
self,
limits: SseLimits,
) -> Result<ClientCredentialsSseStream, ClientCredentialsError> {
Ok(ClientCredentialsSseStream {
stream: Some(
self.response
.into_sse_stream(limits)
.map_err(|_| ClientCredentialsError::UnexpectedResponse)?,
),
snapshot: self.snapshot,
owner: self.owner,
cancellation: self.cancellation,
deadline: self.deadline,
finished: false,
})
}
}
pub struct ClientCredentialsSseStream {
stream: Option<ModernHttpSseResponseStream>,
snapshot: ClientCredentialsSnapshot,
owner: McpRequestCancellation,
cancellation: McpRequestCancellation,
deadline: Time,
finished: bool,
}
impl ClientCredentialsSseStream {
pub fn close(&mut self) {
self.stream = None;
}
pub async fn next_event(&mut self, cx: &Cx) -> Result<Option<String>, ClientCredentialsError> {
if self.finished {
return Ok(None);
}
let mut stream = self.stream.take().ok_or(ClientCredentialsError::Closed)?;
let result = active(
cx,
self.deadline,
&self.owner,
&self.cancellation,
Some(&self.snapshot),
async {
stream
.next_event(cx)
.await
.map_err(|_| ClientCredentialsError::UnexpectedResponse)
},
)
.await;
match &result {
Ok(Some(_)) => self.stream = Some(stream),
Ok(None) => self.finished = true,
Err(_) => {}
}
result
}
}
#[cfg(test)]
mod tests {
use super::*;
fn url(text: &str) -> CanonicalHttpUrl {
CanonicalHttpUrl::parse(text).unwrap()
}
fn plan() -> ClientCredentialsPlan {
ClientCredentialsPlan::new(
url("https://resource.example/mcp"),
TrustedOAuthIssuer::new("https://issuer.example").unwrap(),
"service-client",
"unit-secret",
vec!["read".to_owned(), "write".to_owned()],
)
.unwrap()
}
fn issuer() -> Value {
json!({"issuer":"https://issuer.example","token_endpoint":"https://issuer.example/token",
"grant_types_supported":["client_credentials"],"token_endpoint_auth_methods_supported":["client_secret_basic"]})
}
fn inner() -> ClientInner {
let plan = plan();
ClientInner {
resource: plan.discovery.resource,
token_endpoint: url("https://issuer.example/token"),
client_id: "service-client".to_owned(),
scopes: vec!["read".to_owned(), "write".to_owned()],
authentication: plan.authentication,
issuer_roots: vec![],
resource_tls: None,
timeout: Duration::from_secs(30),
maximum_lifetime: Duration::from_secs(60),
leeway: Duration::from_secs(30),
closed: McpRequestCancellation::new(),
pending: AtomicUsize::new(0),
state: Arc::new(Mutex::new(TokenState::default())),
}
}
fn core() -> CoreRequest {
CoreRequest::decode(ProtocolEra::Modern2026, "tools/list", Some(&json!({"_meta":{
"io.modelcontextprotocol/protocolVersion":"2026-07-28","io.modelcontextprotocol/clientCapabilities":{},
"com.example/tenant":"preserved"}}))).unwrap()
}
#[test]
fn machine_metadata_does_not_require_a_browser_endpoint_or_pkce() {
let plan = plan();
assert!(
admit_machine_issuer(
&plan.discovery,
&plan.discovery.issuers[0],
&serde_json::to_vec(&issuer()).unwrap(),
&plan.authentication
)
.is_ok()
);
for (key, value) in [
("issuer", json!("https://other.example")),
("token_endpoint", json!("https://other.example/token")),
("grant_types_supported", json!(["authorization_code"])),
(
"token_endpoint_auth_methods_supported",
json!(["private_key_jwt"]),
),
("protected_resources", json!(["https://other.example/mcp"])),
("signed_metadata", json!("unverified")),
] {
let mut document = issuer();
document[key] = value;
assert!(
admit_machine_issuer(
&plan.discovery,
&plan.discovery.issuers[0],
&serde_json::to_vec(&document).unwrap(),
&plan.authentication
)
.is_err()
);
}
}
#[test]
fn basic_credentials_encode_components_before_encoding_the_pair() {
assert_eq!(
basic("service", "secret").unwrap(),
"Basic c2VydmljZTpzZWNyZXQ="
);
assert_eq!(
basic("a:b +", "c/d=\u{e9}").unwrap(),
"Basic YSUzQWIrJTJCOmMlMkZkJTNEJUMzJUE5"
);
assert_eq!(component("a:b +"), "a%3Ab+%2B");
assert_eq!(component("c/d=\u{e9}"), "c%2Fd%3D%C3%A9");
assert_eq!(
form(&[("scope", ""), ("grant_type", "client_credentials")]),
"grant_type=client_credentials"
);
}
#[test]
fn token_admission_bounds_expiry_scopes_and_resource_without_refresh_authority() {
let inner = inner();
let now = Instant::now();
let token = admit_token(
&inner,
br#"{"access_token":"access","token_type":"Bearer","expires_in":600,"scope":"read"}"#,
now,
)
.unwrap();
assert_eq!(token.expires_at, now + Duration::from_secs(60));
assert_eq!(token.scopes, ["read"]);
assert!(
token
.bearer
.authorization_for_target(&inner.resource)
.is_some()
);
assert!(
token
.bearer
.authorization_for_target(&url("https://issuer.example/token"))
.is_none()
);
for raw in [
r#"{"access_token":"access","token_type":"Basic"}"#,
r#"{"access_token":"access","token_type":"Bearer","scope":"admin"}"#,
r#"{"access_token":"access","token_type":"Bearer","expires_in":null}"#,
r#"{"access_token":"access","token_type":"Bearer","expires_in":0}"#,
r#"{"access_token":"access","access_token":"duplicate","token_type":"Bearer"}"#,
r#"["access","Bearer",60]"#,
r#"{"access_token":"access","token_type":"Bearer","resource":"https://wrong.example"}"#,
] {
assert!(admit_token(&inner, raw.as_bytes(), now).is_err());
}
}
#[test]
fn explicit_profile_stamping_preserves_metadata_and_refuses_other_extensions() {
let (wire, _) = prepare(
&url("https://resource.example/mcp"),
&core(),
&RequestId::Number(7),
)
.unwrap();
let body: Value = serde_json::from_slice(wire.body()).unwrap();
assert_eq!(body["params"]["_meta"]["com.example/tenant"], "preserved");
assert_eq!(
body["params"]["_meta"][FINAL_CLIENT_CAPABILITIES_META_KEY]["extensions"],
json!({CLIENT_CREDENTIALS_EXTENSION:{}})
);
assert!(
!wire
.headers()
.iter()
.any(|(name, _)| name.eq_ignore_ascii_case("authorization"))
);
let mut params = core().encode_params().unwrap().unwrap();
params["_meta"]["io.modelcontextprotocol/clientCapabilities"]["extensions"] =
json!({"io.modelcontextprotocol/tasks":{}});
let other =
CoreRequest::decode(ProtocolEra::Modern2026, "tools/list", Some(¶ms)).unwrap();
assert!(
prepare(
&url("https://resource.example/mcp"),
&other,
&RequestId::Number(7)
)
.is_err()
);
}
#[test]
fn source_closure_revokes_existing_token_clones_without_an_http_request() {
let inner = inner();
let token = admit_token(
&inner,
br#"{"access_token":"access","token_type":"Bearer"}"#,
Instant::now(),
)
.unwrap();
let clone = token.bearer.clone();
assert!(clone.authorization_for_target(&inner.resource).is_some());
inner.closed.cancel();
assert!(clone.authorization_for_target(&inner.resource).is_none());
assert!(check_token(&token.bearer, token.expires_at).is_err());
}
#[test]
fn secrets_are_not_retained_by_debug_or_error_diagnostics() {
let policy = plan();
assert!(!format!("{policy:?}").contains("unit-secret"));
assert!(
ClientCredentialsPlan::new(
url("https://resource.example/mcp"),
TrustedOAuthIssuer::new("https://issuer.example").unwrap(),
"service",
"bad\r\nsecret",
vec![]
)
.is_err()
);
assert!(plan().with_maximum_token_lifetime(Duration::ZERO).is_err());
}
#[test]
fn shared_acquisition_bound_releases_capacity_when_work_is_dropped() {
let pending = AtomicUsize::new(0);
let permits: Vec<_> = (0..MAX_ACQUISITIONS)
.map(|_| AcquisitionPermit::new(&pending).unwrap())
.collect();
assert!(matches!(
AcquisitionPermit::new(&pending),
Err(ClientCredentialsError::Saturated)
));
drop(permits);
assert_eq!(pending.load(Ordering::Acquire), 0);
let permit = AcquisitionPermit::new(&pending).unwrap();
drop(async move {
let _permit = permit;
std::future::pending::<()>().await;
});
assert_eq!(pending.load(Ordering::Acquire), 0);
}
#[test]
fn typed_json_admission_preserves_exact_members_and_refuses_foreign_ids() {
let wire=br#"{"jsonrpc":"2.0","id":7,"result":{"resultType":"complete","tools":[],"ttlMs":0,"cacheScope":"private","x-exact":1.20e+4}}"#;
let result = decoded_result(&core(), &RequestId::Number(7), wire, 4096).unwrap();
assert!(result.encode().unwrap().contains("1.20e+4"));
assert!(decoded_result(&core(), &RequestId::Number(8), wire, 4096).is_err());
assert!(decoded_result(&core(), &RequestId::Number(7), wire, 20).is_err());
}
#[test]
fn post_selection_is_explicit_and_never_reinterprets_basic_metadata() {
let basic = plan();
let post = plan()
.with_secret_authentication(ClientSecretAuthenticationMethod::Post)
.unwrap();
assert_eq!(basic.authentication.method(), "client_secret_basic");
assert_eq!(post.authentication.method(), "client_secret_post");
for (methods, basic_ok, post_ok) in [
(json!(["client_secret_basic"]), true, false),
(json!(["client_secret_post"]), false, true),
(
json!([
"private_key_jwt",
"client_secret_post",
"client_secret_basic"
]),
true,
true,
),
(json!(["private_key_jwt"]), false, false),
(json!(["client_secret_post "]), false, false),
] {
let mut document = issuer();
document["token_endpoint_auth_methods_supported"] = methods;
let bytes = serde_json::to_vec(&document).unwrap();
for (selected, expected) in [(&basic, basic_ok), (&post, post_ok)] {
assert_eq!(
admit_machine_issuer(
&selected.discovery,
&selected.discovery.issuers[0],
&bytes,
&selected.authentication
)
.is_ok(),
expected
);
}
}
assert_eq!(basic.authentication.method(), "client_secret_basic");
assert_eq!(post.authentication.method(), "client_secret_post");
}
#[test]
fn post_grant_uses_exactly_one_authentication_channel_and_one_encoding_pass() {
let resource = url("https://resource.example/mcp");
let deadline = Time::from_nanos(99);
let scopes = vec!["read".to_owned(), "write".to_owned()];
let post = prepare_secret_grant(
ClientSecretAuthenticationMethod::Post,
"a:b +",
"c/d=%é&scope=admin",
&resource,
&scopes,
deadline,
)
.unwrap();
assert!(post.authorization.is_none());
assert_eq!(post.deadline, deadline);
assert_eq!(
std::str::from_utf8(&post.body).unwrap(),
"grant_type=client_credentials&resource=https%3A%2F%2Fresource.example%2Fmcp&scope=read+write&client_id=a%3Ab+%2B&client_secret=c%2Fd%3D%25%C3%A9%26scope%3Dadmin"
);
assert_eq!(
post.body
.iter()
.fold(0, |n, &byte| n + usize::from(byte == b'&')),
4
);
let basic = prepare_secret_grant(
ClientSecretAuthenticationMethod::Basic,
"a:b +",
"c/d=%é&scope=admin",
&resource,
&scopes,
deadline,
)
.unwrap();
assert_eq!(
basic.authorization,
Some(super::basic("a:b +", "c/d=%é&scope=admin").unwrap())
);
assert_eq!(
std::str::from_utf8(&basic.body).unwrap(),
"grant_type=client_credentials&resource=https%3A%2F%2Fresource.example%2Fmcp&scope=read+write"
);
}
#[test]
fn post_grant_omits_empty_scopes_and_enforces_the_encoded_body_bound() {
let resource = url("https://resource.example/mcp");
let deadline = Time::from_nanos(99);
let post = prepare_secret_grant(
ClientSecretAuthenticationMethod::Post,
"service-client",
"unit-secret",
&resource,
&[],
deadline,
)
.unwrap();
assert_eq!(
std::str::from_utf8(&post.body).unwrap(),
"grant_type=client_credentials&resource=https%3A%2F%2Fresource.example%2Fmcp&client_id=service-client&client_secret=unit-secret"
);
let oversized = "%".repeat(MAX_SECRET_GRANT_BYTES);
assert!(matches!(
prepare_secret_grant(
ClientSecretAuthenticationMethod::Post,
"service-client",
&oversized,
&resource,
&[],
deadline
),
Err(ClientCredentialsError::RequestTooLarge)
));
}
#[test]
fn switching_a_not_yet_discovered_secret_plan_preserves_its_authority_and_limits() {
let original = plan()
.with_timeout(Duration::from_secs(17))
.unwrap()
.with_maximum_token_lifetime(Duration::from_secs(43))
.unwrap();
let resource = original.discovery.resource.clone();
let scopes = original.discovery.scopes.clone();
let client_id = original.discovery.client_id.clone();
let post = original
.with_secret_authentication(ClientSecretAuthenticationMethod::Post)
.unwrap();
assert_eq!(post.discovery.resource, resource);
assert_eq!(post.discovery.scopes, scopes);
assert_eq!(post.discovery.client_id, client_id);
assert_eq!(post.discovery.timeout, Duration::from_secs(17));
assert_eq!(post.maximum_lifetime, Duration::from_secs(43));
assert!(format!("{post:?}").contains("client_secret_post"));
assert!(!format!("{post:?}").contains("unit-secret"));
let basic = post
.with_secret_authentication(ClientSecretAuthenticationMethod::Basic)
.unwrap();
assert_eq!(basic.authentication.method(), "client_secret_basic");
}
#[test]
fn post_metadata_keeps_issuer_scope_and_credential_destination_checks() {
let post = plan()
.with_secret_authentication(ClientSecretAuthenticationMethod::Post)
.unwrap();
let mut document = issuer();
document["token_endpoint_auth_methods_supported"] = json!(["client_secret_post"]);
let valid = document.clone();
for (key, value) in [
("issuer", json!("https://issuer.example/other")),
("token_endpoint", json!("http://issuer.example/token")),
("token_endpoint", json!("https://other.example/token")),
("scopes_supported", json!(["admin"])),
("token_endpoint_auth_methods_supported", Value::Null),
] {
document = valid.clone();
document[key] = value;
assert!(
admit_machine_issuer(
&post.discovery,
&post.discovery.issuers[0],
&serde_json::to_vec(&document).unwrap(),
&post.authentication
)
.is_err()
);
}
assert!(
admit_machine_issuer(
&post.discovery,
&post.discovery.issuers[0],
&serde_json::to_vec(&valid).unwrap(),
&post.authentication
)
.is_ok()
);
}
fn valid_machine_token() -> Value {
json!({"access_token":"access-secret-canary", "token_type":"Bearer",
"expires_in":60, "scope":"read"})
}
#[test]
fn machine_token_forbidden_members_reject_presence_not_just_values() {
let inner = inner();
let valid = valid_machine_token();
for key in [
"refresh_token",
"id_token",
"cnf",
"issued_token_type",
"error_description",
"error_uri",
] {
for value in [
Value::Null,
json!(""),
json!("peer-secret-canary"),
json!(false),
json!(7),
json!([]),
json!({}),
] {
assert!(
admit_token(&inner, &serde_json::to_vec(&valid).unwrap(), Instant::now())
.is_ok()
);
let mut changed = valid.clone();
changed[key] = value;
let error = admit_token(
&inner,
&serde_json::to_vec(&changed).unwrap(),
Instant::now(),
)
.err()
.unwrap();
assert!(matches!(error, ClientCredentialsError::InvalidToken));
let diagnostic = format!("{error:?} {error}");
assert!(!diagnostic.contains("peer-secret-canary"));
assert!(!diagnostic.contains("access-secret-canary"));
}
}
}
#[test]
fn machine_token_cannot_downgrade_sender_constraints_or_exchange_grants() {
let inner = inner();
for (key, value) in [
("token_type", json!("DPoP")),
("cnf", json!({"jkt":"proof-key"})),
("cnf", json!({"x5t#S256":"client-certificate"})),
(
"issued_token_type",
json!("urn:ietf:params:oauth:token-type:access_token"),
),
] {
let mut document = valid_machine_token();
document[key] = value;
assert!(matches!(
admit_token(
&inner,
&serde_json::to_vec(&document).unwrap(),
Instant::now()
),
Err(ClientCredentialsError::InvalidToken)
));
}
for spelling in ["Bearer", "bearer", "bEaReR"] {
let mut document = valid_machine_token();
document["token_type"] = json!(spelling);
assert!(
admit_token(
&inner,
&serde_json::to_vec(&document).unwrap(),
Instant::now()
)
.is_ok()
);
}
}
#[test]
fn machine_token_error_fields_never_coexist_with_success() {
let inner = inner();
for key in ["error", "error_description", "error_uri"] {
for value in [Value::Null, json!(""), json!("invalid_client"), json!({})] {
let mut document = valid_machine_token();
document[key] = value;
assert!(matches!(
admit_token(
&inner,
&serde_json::to_vec(&document).unwrap(),
Instant::now()
),
Err(ClientCredentialsError::InvalidToken)
));
}
}
assert!(
admit_token(
&inner,
br#"{"error":"invalid_client","error_description":"peer-secret"}"#,
Instant::now()
)
.is_err()
);
}
#[test]
fn machine_token_ignores_unknown_metadata_without_promoting_nested_authority() {
let inner = inner();
let started = Instant::now();
let mut document = valid_machine_token();
document["x-extension"] = json!({"refresh_token":"ignored-nested", "id_token":null,
"cnf":{"jkt":"nested"}, "error":"nested", "scope":"admin"});
document["registration_client_uri"] = json!("https://untrusted.example/registration");
document["token_endpoint"] = json!("https://untrusted.example/token");
let admitted =
admit_token(&inner, &serde_json::to_vec(&document).unwrap(), started).unwrap();
assert_eq!(admitted.scopes, ["read"]);
assert_eq!(admitted.expires_at, started + Duration::from_secs(60));
assert_eq!(admitted.bearer.resource(), &inner.resource);
document.as_object_mut().unwrap().remove("access_token");
document["x-extension"]["access_token"] = json!("nested-not-a-grant");
assert!(admit_token(&inner, &serde_json::to_vec(&document).unwrap(), started).is_err());
}
#[test]
fn machine_token_escaped_names_and_duplicates_cannot_hide_forbidden_members() {
let inner = inner();
for field in [
r#""refresh\u005ftoken":null"#,
r#""id\u005ftoken":"secret""#,
r#""c\u006ef":{}"#,
r#""error\u005furi":"secret""#,
r#""refresh_token":null,"refresh_token":"secret""#,
] {
let source = format!(r#"{{"access_token":"access","token_type":"Bearer",{field}}}"#);
assert!(matches!(
admit_token(&inner, source.as_bytes(), Instant::now()),
Err(ClientCredentialsError::InvalidToken)
));
}
assert!(
admit_token(
&inner,
br#"{"access_token":"one","access_token":"two","token_type":"Bearer"}"#,
Instant::now()
)
.is_err()
);
}
#[test]
fn machine_token_byte_limit_includes_ignored_members_and_trailing_space() {
let inner = inner();
let mut exact = serde_json::to_vec(&valid_machine_token()).unwrap();
exact.resize(MAX_TOKEN_BYTES, b' ');
assert!(admit_token(&inner, &exact, Instant::now()).is_ok());
exact.push(b' ');
assert!(matches!(
admit_token(&inner, &exact, Instant::now()),
Err(ClientCredentialsError::InvalidToken)
));
let mut oversized = valid_machine_token();
oversized["x-extension"] = json!("x".repeat(MAX_TOKEN_BYTES));
assert!(matches!(
admit_token(
&inner,
&serde_json::to_vec(&oversized).unwrap(),
Instant::now()
),
Err(ClientCredentialsError::InvalidToken)
));
}
#[test]
fn machine_token_opaque_jwt_text_never_supplies_scope_or_expiration() {
let mut inner = inner();
inner.maximum_lifetime = Duration::from_secs(17);
let started = Instant::now();
let document = br#"{"access_token":"e30.eyJleHAiOjB9.signature","token_type":"Bearer","exp":0,"scp":["admin"]}"#;
let admitted = admit_token(&inner, document, started).unwrap();
assert_eq!(admitted.scopes, inner.scopes);
assert_eq!(admitted.expires_at, started + Duration::from_secs(17));
assert_eq!(
admitted
.bearer
.authorization_for_target(&inner.resource)
.unwrap(),
"Bearer e30.eyJleHAiOjB9.signature"
);
}
fn declared_discovery() -> CoreRequest {
let (_, declared) = prepare(
&url("https://resource.example/mcp"),
&core(),
&RequestId::Number(7),
)
.unwrap();
let params = declared.encode_params().unwrap().unwrap();
CoreRequest::decode(
ProtocolEra::Modern2026,
"server/discover",
Some(&json!({"_meta":params["_meta"]})),
)
.unwrap()
}
fn discovery_reply(capabilities: Value) -> Vec<u8> {
serde_json::to_vec(&json!({"jsonrpc":"2.0","id":7,"result":{
"resultType":"complete","supportedVersions":[FINAL_PROTOCOL_VERSION],
"capabilities":capabilities,"ttlMs":0,"cacheScope":"private"
}}))
.unwrap()
}
#[test]
fn machine_discovery_accepts_absent_or_exact_advertisement_without_changing_the_request() {
let request = declared_discovery();
let before = request.encode_params().unwrap();
for caps in [
json!({}),
json!({"extensions":{}}),
json!({"extensions":{CLIENT_CREDENTIALS_EXTENSION:{}}}),
json!({"extensions":{"com.example/independent":{}}}),
] {
assert!(
admit_resource(&request, &RequestId::Number(7), &discovery_reply(caps)).is_ok()
);
assert_eq!(request.encode_params().unwrap(), before);
}
}
#[test]
fn machine_discovery_never_treats_present_null_or_malformed_settings_as_absence() {
let request = declared_discovery();
for settings in [
Value::Null,
json!(false),
json!(7),
json!(""),
json!([]),
json!({"enabled":true}),
] {
let bytes =
discovery_reply(json!({"extensions":{CLIENT_CREDENTIALS_EXTENSION:settings}}));
assert!(matches!(
admit_resource(&request, &RequestId::Number(7), &bytes),
Err(ClientCredentialsError::Negotiation)
));
}
for extensions in [Value::Null, json!(false), json!(7), json!(""), json!([])] {
let bytes = discovery_reply(json!({"extensions":extensions}));
assert!(matches!(
admit_resource(&request, &RequestId::Number(7), &bytes),
Err(ClientCredentialsError::Negotiation)
));
}
assert!(
admit_resource(&request, &RequestId::Number(7), &discovery_reply(json!({}))).is_ok()
);
}
#[test]
fn machine_discovery_requires_this_requests_client_declaration_even_when_server_advertises() {
let declared = declared_discovery();
let base = declared.encode_params().unwrap().unwrap();
let response = discovery_reply(json!({"extensions":{CLIENT_CREDENTIALS_EXTENSION:{}}}));
for caps in [
json!({}),
json!({"extensions":{}}),
json!({"extensions":{"com.example/independent":{}}}),
json!({"experimental":{CLIENT_CREDENTIALS_EXTENSION:{}}}),
] {
let mut params = base.clone();
params["_meta"][FINAL_CLIENT_CAPABILITIES_META_KEY] = caps;
let request =
CoreRequest::decode(ProtocolEra::Modern2026, "server/discover", Some(¶ms))
.unwrap();
assert!(matches!(
admit_resource(&request, &RequestId::Number(7), &response),
Err(ClientCredentialsError::Negotiation)
));
}
assert!(admit_resource(&declared, &RequestId::Number(7), &response).is_ok());
}
#[test]
fn optional_advertisement_does_not_bypass_version_correlation_or_document_bounds() {
let request = declared_discovery();
let valid = discovery_reply(json!({}));
assert!(admit_resource(&request, &RequestId::Number(7), &valid).is_ok());
for id in [RequestId::Number(8), RequestId::String("7".to_owned())] {
assert!(admit_resource(&request, &id, &valid).is_err());
}
let mut wrong: Value = serde_json::from_slice(&valid).unwrap();
wrong["result"]["supportedVersions"] = json!(["2024-11-05"]);
assert!(
admit_resource(
&request,
&RequestId::Number(7),
&serde_json::to_vec(&wrong).unwrap()
)
.is_err()
);
let array =
serde_json::to_vec(&json!([serde_json::from_slice::<Value>(&valid).unwrap()])).unwrap();
assert!(admit_resource(&request, &RequestId::Number(7), &array).is_err());
let mut oversized = valid;
oversized.resize(MAX_TOKEN_BYTES + 1, b' ');
assert!(admit_resource(&request, &RequestId::Number(7), &oversized).is_err());
}
#[test]
fn optional_advertisement_keeps_raw_duplicate_and_escaped_member_checks() {
let request = declared_discovery();
for members in [
r#""extensions":{"io.modelcontextprotocol/oauth-client-credentials":{},"io.modelcontextprotocol/oauth-client-credentials":null}"#,
r#""extensions":{"io.modelcontextprotocol/oauth-client-credentials":{}},"extensions":{}"#,
r#""extens\u0069ons":null"#,
r#""extensions":{"io.modelcontextprotocol/oauth-client-credent\u0069als":null}"#,
] {
let source = format!(
r#"{{"jsonrpc":"2.0","id":7,"result":{{"resultType":"complete","supportedVersions":["2026-07-28"],"capabilities":{{{members}}},"ttlMs":0,"cacheScope":"private"}}}}"#
);
assert!(matches!(
admit_resource(&request, &RequestId::Number(7), source.as_bytes()),
Err(ClientCredentialsError::Negotiation)
));
}
}
#[test]
fn machine_advertisement_is_revalidated_instead_of_cached_between_requests() {
let request = declared_discovery();
let valid = discovery_reply(json!({"extensions":{CLIENT_CREDENTIALS_EXTENSION:{}}}));
let malformed = discovery_reply(
json!({"extensions":{CLIENT_CREDENTIALS_EXTENSION:{"unexpected":true}}}),
);
let absent = discovery_reply(json!({}));
assert!(admit_resource(&request, &RequestId::Number(7), &valid).is_ok());
assert!(admit_resource(&request, &RequestId::Number(7), &malformed).is_err());
assert!(admit_resource(&request, &RequestId::Number(7), &absent).is_ok());
assert!(admit_resource(&request, &RequestId::Number(7), &malformed).is_err());
}
}