use std::fmt;
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Instant;
use asupersync::Cx;
use fastmcp_core::McpRequestCancellation;
use fastmcp_protocol::{
CoreRequest, CoreResult, FINAL_PROTOCOL_VERSION, FinalCoreRequest, FinalCoreResult, RequestId,
ServerNotification,
};
use super::{
ManagedCoreError, ManagedCoreEvent, ManagedCoreLimits, ManagedOAuthSession, bounded_wait,
call_deadline, check_call, prepare,
};
use crate::cache::{
CachePartitionKey, FinalCacheGeneration, FinalCacheInsert, FinalCacheKey, FinalCacheLookup,
FinalCacheResultSet, FinalCacheStats, FinalResultCache, MAX_FINAL_CACHE_CAPACITY,
MAX_FINAL_CACHE_MAX_BYTES,
};
use crate::http_auth::managed::OAuthCredentialSnapshot;
pub mod watch;
#[derive(Clone, Copy, Debug)]
pub struct ManagedResourceLimits {
core: ManagedCoreLimits,
maximum_contents: usize,
}
impl Default for ManagedResourceLimits {
fn default() -> Self {
Self {
core: ManagedCoreLimits::default(),
maximum_contents: 1024,
}
}
}
impl ManagedResourceLimits {
pub fn new(
core: ManagedCoreLimits,
maximum_contents: usize,
) -> Result<Self, ManagedResourceError> {
if maximum_contents > 100_000 {
return Err(ManagedResourceError::InvalidLimits);
}
Ok(Self {
core,
maximum_contents,
})
}
}
#[derive(Debug)]
pub enum ManagedResourceError {
InvalidLimits,
NotResourceRead,
InvalidResult,
ContentsLimit,
Invalidated,
CredentialChanged,
CredentialUnavailable,
CacheUnavailable,
AbortedByHost,
Core(ManagedCoreError),
}
impl fmt::Display for ManagedResourceError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Core(error) => fmt::Display::fmt(error, f),
Self::InvalidLimits => f.write_str("invalid managed resource limits"),
Self::NotResourceRead => f.write_str("request is not a modern core resource read"),
Self::InvalidResult => f.write_str("result is not an admitted resource-read outcome"),
Self::ContentsLimit => f.write_str("resource contents count exceeds its limit"),
Self::Invalidated => f.write_str("resource read was invalidated before delivery"),
Self::CredentialChanged => {
f.write_str("resource request used a different credential generation")
}
Self::CredentialUnavailable => {
f.write_str("resource credential is expired or locally revoked")
}
Self::CacheUnavailable => f.write_str("resource cache state is unavailable"),
Self::AbortedByHost => f.write_str("resource read stopped by its host"),
}
}
}
impl std::error::Error for ManagedResourceError {}
impl From<ManagedCoreError> for ManagedResourceError {
fn from(error: ManagedCoreError) -> Self {
Self::Core(error)
}
}
pub struct ManagedResourceRead {
uri: String,
result: CoreResult,
credential_generation: u64,
cache_hit: bool,
}
impl ManagedResourceRead {
pub fn uri(&self) -> &str {
&self.uri
}
pub fn result(&self) -> &CoreResult {
&self.result
}
pub fn into_result(self) -> CoreResult {
self.result
}
pub fn credential_generation(&self) -> u64 {
self.credential_generation
}
pub fn is_cache_hit(&self) -> bool {
self.cache_hit
}
pub fn is_complete(&self) -> bool {
matches!(
&self.result,
CoreResult::Final(FinalCoreResult::ResourcesRead { .. })
)
}
}
#[derive(Clone)]
pub struct ManagedResourceClient {
session: ManagedOAuthSession,
limits: ManagedResourceLimits,
cache: Arc<Mutex<FinalResultCache>>,
}
impl ManagedResourceClient {
pub fn new(session: ManagedOAuthSession, limits: ManagedResourceLimits) -> Self {
let mut cache = FinalResultCache::default();
cache.set_enabled(false);
Self {
session,
limits,
cache: Arc::new(Mutex::new(cache)),
}
}
pub fn with_cache_limits(
mut self,
entries: usize,
bytes: usize,
) -> Result<Self, ManagedResourceError> {
if !(1..=MAX_FINAL_CACHE_CAPACITY).contains(&entries)
|| !(1..=MAX_FINAL_CACHE_MAX_BYTES).contains(&bytes)
{
return Err(ManagedResourceError::InvalidLimits);
}
self.cache = Arc::new(Mutex::new(FinalResultCache::with_limits(entries, bytes)));
Ok(self)
}
pub fn clear(&self) -> Result<(), ManagedResourceError> {
self.cache()?.clear();
Ok(())
}
pub fn invalidate_notification(
&self,
notification: &ServerNotification,
) -> Result<(), ManagedResourceError> {
self.cache()?.invalidate_notification(notification);
Ok(())
}
pub fn cache_stats(&self) -> Result<FinalCacheStats, ManagedResourceError> {
Ok(self.cache()?.stats())
}
pub async fn read<I, O>(
&self,
cx: &Cx,
request: CoreRequest,
next_id: I,
observe: O,
) -> Result<ManagedResourceRead, ManagedResourceError>
where
I: FnOnce() -> Result<RequestId, ManagedResourceError>,
O: FnMut(Box<ServerNotification>) -> Result<(), ManagedResourceError>,
{
self.read_with_cancellation(
cx,
&McpRequestCancellation::new(),
request,
next_id,
observe,
)
.await
}
pub async fn read_with_cancellation<I, O>(
&self,
cx: &Cx,
cancellation: &McpRequestCancellation,
request: CoreRequest,
next_id: I,
mut observe: O,
) -> Result<ManagedResourceRead, ManagedResourceError>
where
I: FnOnce() -> Result<RequestId, ManagedResourceError>,
O: FnMut(Box<ServerNotification>) -> Result<(), ManagedResourceError>,
{
let deadline = call_deadline(cx, cancellation, self.limits.core.timeout)?;
let (uri, reusable) = read_identity(&request)?;
let uri = uri.to_owned();
let _ = prepare(
self.session.resource().as_str(),
request.clone(),
RequestId::Number(0),
self.limits.core,
)?;
let result_set = FinalCacheResultSet::Resource(uri.clone());
let captured = self.cache()?.begin_fetch(&result_set);
bounded_wait(cx, cancellation, deadline, async {
Ok(async {
let credential = self
.session
.credential_with_cancellation(cx, cancellation)
.await
.map_err(ManagedCoreError::from)?;
require_credential(&credential)?;
self.require_generation(&result_set, captured)?;
let key = if reusable {
Some(cache_key(
self.session.resource().as_str(),
&request,
credential.generation(),
)?)
} else {
None
};
let cached = match &key {
Some(key) => match self.cache()?.lookup(key) {
FinalCacheLookup::Fresh(result) => Some(result),
FinalCacheLookup::Miss(_) => None,
},
None => None,
};
let cache_hit = cached.is_some();
let (result, receipt) = if let Some(result) = cached {
let bytes = result
.encode()
.map_err(|_| ManagedResourceError::InvalidResult)?
.len();
if bytes > self.limits.core.frame_bytes || bytes > self.limits.core.total_bytes
{
return Err(ManagedCoreError::ResponseByteLimit.into());
}
(result, Instant::now())
} else {
let id = next_id()?;
check_call(cx, cancellation, deadline)?;
require_credential(&credential)?;
self.require_generation(&result_set, captured)?;
let mut call = Box::pin(self.session.request_core_with_cancellation(
cx,
cancellation,
request,
id,
self.limits.core,
))
.await?;
if call.credential_generation() != credential.generation() {
return Err(ManagedResourceError::CredentialChanged);
}
call.deadline = deadline;
let result = loop {
let event = call
.next_event(cx)
.await?
.ok_or(ManagedResourceError::InvalidResult)?;
check_call(cx, cancellation, deadline)?;
require_credential(&credential)?;
match event {
ManagedCoreEvent::Notification(notification) => {
self.invalidate_notification(¬ification)?;
observe(notification)?;
check_call(cx, cancellation, deadline)?;
require_credential(&credential)?;
self.require_generation(&result_set, captured)?;
}
ManagedCoreEvent::Result(result) => break *result,
}
};
(result, Instant::now())
};
let complete = admit_result(&result, self.limits.maximum_contents)?;
check_call(cx, cancellation, deadline)?;
{
let mut cache = self.cache()?;
if cache.begin_fetch(&result_set) != captured {
return Err(ManagedResourceError::Invalidated);
}
require_credential(&credential)?;
if complete && !cache_hit && cache.is_enabled() {
if let Some(key) = key {
if cache.insert_if_current_at(key, captured, result.clone(), receipt)
== FinalCacheInsert::InvalidatedDuringFetch
{
return Err(ManagedResourceError::Invalidated);
}
}
}
}
check_call(cx, cancellation, deadline)?;
require_credential(&credential)?;
self.require_generation(&result_set, captured)?;
Ok(ManagedResourceRead {
uri,
result,
credential_generation: credential.generation(),
cache_hit,
})
}
.await)
})
.await?
}
fn cache(&self) -> Result<MutexGuard<'_, FinalResultCache>, ManagedResourceError> {
self.cache
.lock()
.map_err(|_| ManagedResourceError::CacheUnavailable)
}
fn require_generation(
&self,
result_set: &FinalCacheResultSet,
captured: FinalCacheGeneration,
) -> Result<(), ManagedResourceError> {
if self.cache()?.begin_fetch(result_set) != captured {
return Err(ManagedResourceError::Invalidated);
}
Ok(())
}
}
fn require_credential(credential: &OAuthCredentialSnapshot) -> Result<(), ManagedResourceError> {
if credential.credential().is_revoked() || Instant::now() >= credential.expires_at() {
return Err(ManagedResourceError::CredentialUnavailable);
}
Ok(())
}
fn read_identity(request: &CoreRequest) -> Result<(&str, bool), ManagedResourceError> {
let CoreRequest::Final(FinalCoreRequest::ResourcesRead(params)) = request else {
return Err(ManagedResourceError::NotResourceRead);
};
Ok((
params.uri.as_str(),
params.input_responses.is_none() && params.request_state.is_none(),
))
}
fn cache_key(
target: &str,
request: &CoreRequest,
generation: u64,
) -> Result<FinalCacheKey, ManagedResourceError> {
let (uri, reusable) = read_identity(request)?;
if !reusable {
return Err(ManagedCoreError::InvalidRequest.into());
}
let params = request
.encode_params()
.map_err(|_| ManagedCoreError::InvalidRequest)?
.ok_or(ManagedCoreError::InvalidRequest)?;
let projection =
serde_json::to_string(¶ms).map_err(|_| ManagedCoreError::InvalidRequest)?;
Ok(FinalCacheKey::new(
target,
FINAL_PROTOCOL_VERSION,
"included-in-exact-params",
"core-only",
"resources/read",
projection,
None,
0,
0,
0,
0,
CachePartitionKey::new(format!("managed-resource-generation-{generation}")),
FinalCacheResultSet::Resource(uri.to_owned()),
))
}
fn admit_result(
result: &CoreResult,
maximum_contents: usize,
) -> Result<bool, ManagedResourceError> {
match result {
CoreResult::Final(FinalCoreResult::ResourcesRead { result, .. }) => {
if result.payload.contents.len() > maximum_contents {
return Err(ManagedResourceError::ContentsLimit);
}
Ok(true)
}
CoreResult::Final(FinalCoreResult::ResourcesReadInputRequired { .. }) => Ok(false),
_ => Err(ManagedResourceError::InvalidResult),
}
}
#[cfg(test)]
mod tests {
use super::*;
use fastmcp_protocol::protocol_policy::ProtocolEra;
use fastmcp_protocol::{ClientCapabilities, FinalRequestMeta, JsonRpcRequest};
use serde_json::json;
fn request(uri: &str) -> CoreRequest {
decode(json!({"uri":uri}))
}
fn decode(mut params: serde_json::Value) -> CoreRequest {
params["_meta"] =
serde_json::to_value(FinalRequestMeta::new(ClientCapabilities::default())).unwrap();
CoreRequest::decode(ProtocolEra::Modern2026, "resources/read", Some(¶ms)).unwrap()
}
fn complete(request: &CoreRequest, ttl: u64, scope: &str) -> CoreResult {
request.decode_result(&format!(r#"{{"resultType":"complete","contents":[{{"uri":"file:///one","text":"first"}},{{"uri":"file:///two","blob":"AAEC"}}],"ttlMs":{ttl},"cacheScope":"{scope}","x-exact":{{"z":900719925474099312345,"a":1.20e+4}}}}"#)).unwrap()
}
fn notification(method: &str, params: Option<serde_json::Value>) -> ServerNotification {
ServerNotification::decode(&JsonRpcRequest::notification(method, params)).unwrap()
}
#[test]
fn complete_read_retains_every_content_item_and_exact_unknown_payload() {
let request = request("file:///one");
let result = complete(&request, 60000, "private");
assert!(admit_result(&result, 2).unwrap());
assert!(matches!(
admit_result(&result, 1),
Err(ManagedResourceError::ContentsLimit)
));
let encoded = result.encode().unwrap();
assert!(encoded.contains("file:///two") && encoded.contains("AAEC"));
assert!(encoded.contains("900719925474099312345") && encoded.contains("1.20e+4"));
assert!(encoded.find("first").unwrap() < encoded.find("AAEC").unwrap());
}
#[test]
fn input_required_cache_lookalikes_remain_noncacheable() {
let request = request("file:///one");
let result = request.decode_result(r#"{"resultType":"input_required","requestState":"","ttlMs":60000,"cacheScope":"public"}"#).unwrap();
assert!(!admit_result(&result, 0).unwrap());
assert!(crate::cache::final_cache_hints(&result).is_none());
}
#[test]
fn either_present_continuation_field_bypasses_cache_even_when_empty() {
assert!(read_identity(&request("file:///one")).unwrap().1);
for params in [
json!({"uri":"file:///one","requestState":""}),
json!({"uri":"file:///one","inputResponses":{}}),
json!({"uri":"file:///one","requestState":"opaque","inputResponses":{}}),
] {
let request = decode(params);
assert!(!read_identity(&request).unwrap().1);
assert!(cache_key("https://mcp.example/mcp", &request, 1).is_err());
}
}
#[test]
fn cache_identity_includes_uri_metadata_target_and_credential_generation() {
let original = request("file:///one");
let baseline = cache_key("https://mcp.example/mcp", &original, 1).unwrap();
assert_ne!(
baseline,
cache_key("https://mcp.example/mcp", &original, 2).unwrap()
);
assert_ne!(
baseline,
cache_key("https://other.example/mcp", &original, 1).unwrap()
);
assert_ne!(
baseline,
cache_key("https://mcp.example/mcp", &request("file:///two"), 1).unwrap()
);
let mut params = original.encode_params().unwrap().unwrap();
params["_meta"]["com.example/tenant"] = json!("other");
let changed =
CoreRequest::decode(ProtocolEra::Modern2026, "resources/read", Some(¶ms)).unwrap();
assert_ne!(
baseline,
cache_key("https://mcp.example/mcp", &changed, 1).unwrap()
);
}
#[test]
fn public_hints_cannot_cross_credential_partitions() {
let request = request("file:///one");
let key = cache_key("https://mcp.example/mcp", &request, 1).unwrap();
let mut cache = FinalResultCache::default();
let generation = cache.begin_fetch(key.result_set());
assert_eq!(
cache.insert_if_current(key.clone(), generation, complete(&request, 60000, "public")),
FinalCacheInsert::Stored
);
assert!(matches!(cache.lookup(&key), FinalCacheLookup::Fresh(_)));
let other = cache_key("https://mcp.example/mcp", &request, 2).unwrap();
assert!(matches!(cache.lookup(&other), FinalCacheLookup::Miss(_)));
}
#[test]
fn resource_and_catalog_notifications_fence_fills_but_tools_changes_do_not() {
let request = request("file:///one");
let key = cache_key("https://mcp.example/mcp", &request, 1).unwrap();
for (method, params) in [
(
"notifications/resources/updated",
Some(json!({"uri":"file:///one"})),
),
("notifications/resources/list_changed", None),
] {
let mut cache = FinalResultCache::default();
let before = cache.begin_fetch(key.result_set());
cache.invalidate_notification(¬ification(method, params));
assert_ne!(cache.begin_fetch(key.result_set()), before);
assert_eq!(
cache.insert_if_current(key.clone(), before, complete(&request, 60000, "private")),
FinalCacheInsert::InvalidatedDuringFetch
);
}
let mut cache = FinalResultCache::default();
let before = cache.begin_fetch(key.result_set());
cache.invalidate_notification(¬ification("notifications/tools/list_changed", None));
assert_eq!(cache.begin_fetch(key.result_set()), before);
}
#[test]
fn zero_ttl_and_oversized_results_never_become_hits() {
let request = request("file:///one");
let key = cache_key("https://mcp.example/mcp", &request, 1).unwrap();
let mut cache = FinalResultCache::with_limits(1, 1);
let generation = cache.begin_fetch(key.result_set());
assert_eq!(
cache.insert_if_current(key.clone(), generation, complete(&request, 0, "private")),
FinalCacheInsert::ImmediatelyStale
);
assert_eq!(
cache.insert_if_current(
key.clone(),
generation,
complete(&request, 60000, "private")
),
FinalCacheInsert::Oversized
);
assert!(matches!(cache.lookup(&key), FinalCacheLookup::Miss(_)));
}
#[test]
fn non_read_requests_and_non_read_results_are_not_reinterpreted() {
let params = json!({"_meta":FinalRequestMeta::new(ClientCapabilities::default())});
let catalog =
CoreRequest::decode(ProtocolEra::Modern2026, "resources/list", Some(¶ms)).unwrap();
assert!(matches!(
read_identity(&catalog),
Err(ManagedResourceError::NotResourceRead)
));
let result = catalog
.decode_result(
r#"{"resultType":"complete","resources":[],"ttlMs":0,"cacheScope":"private"}"#,
)
.unwrap();
assert!(matches!(
admit_result(&result, 10),
Err(ManagedResourceError::InvalidResult)
));
assert!(ManagedResourceLimits::new(ManagedCoreLimits::default(), 100_001).is_err());
}
}