fastmcp-client 0.11.0

MCP client implementation for FastMCP
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
//! Schema-bound tool execution over managed OAuth.
//!
//! A client binds one immutable, admitted tool contract to one managed login.
//! Arguments are checked before credential renewal or HTTP dispatch; successful
//! structured output is checked before publication. Protocol admission,
//! correlation, incremental notifications and finite-response EOF checks remain
//! owned by the existing managed core call. No failed call is retried.
//!
//! A definition is a contract, not permission to execute a tool. The host must
//! approve the invocation and invalidate this client when its catalog or policy
//! changes. Server annotations never grant consent, retry or header-disclosure
//! authority. Parameter headers require explicit `review_headers` or
//! `with_reviewed_headers` approval, retained by every continuation.

use std::fmt;
use std::io::{self, Write};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};

use asupersync::Cx;
use fastmcp_core::McpRequestCancellation;
use fastmcp_protocol::http_headers::AdmittedToolHeaderSchema;
use fastmcp_protocol::protocol_policy::ProtocolEra;
use fastmcp_protocol::{
    AdmittedSchema, CoreRequest, CoreResult, FinalTool, RequestId, admit_final_schema,
};
use serde_json::Value;

use super::managed::ManagedOAuthSession;
use super::rpc::tool_headers::ManagedToolHeaderError;
use super::rpc::{ManagedCoreCall, ManagedCoreError, ManagedCoreEvent, ManagedCoreLimits};
use crate::http_executor::parameter_headers::{ReviewedToolHeaders, ToolHeaderDispatchError};

/// Caller-driven catalog watches publishing invalidation-bound tool clients.
pub mod catalog;
/// Disclosure review bound to this client's exact schema and invalidation.
pub mod headers;
/// Explicit multi-round tool operations retaining this same schema contract.
pub mod interaction;

mod validity;
use validity::await_validity;

/// Combined encoded-byte ceiling for one retained input/output schema pair.
pub const MAX_MANAGED_TOOL_SCHEMA_BYTES: usize = 512 * 1024;
/// Maximum UTF-8 bytes retained for a tool's exact, case-sensitive name.
pub const MAX_MANAGED_TOOL_NAME_BYTES: usize = 1024;

/// Fixed diagnostics never retain argument values, schema paths or peer output.
#[derive(Debug)]
pub enum ManagedToolError {
    InvalidDefinition,
    SchemaTooLarge,
    InvalidInputSchema,
    InvalidOutputSchema,
    RequestMismatch,
    InvalidArguments,
    InvalidResult,
    MissingStructuredOutput,
    InvalidStructuredOutput,
    HeaderBindingMismatch,
    Invalidated,
    Closed,
    Headers(ToolHeaderDispatchError),
    Core(ManagedCoreError),
}

impl fmt::Display for ManagedToolError {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.write_str(match self {
            Self::InvalidDefinition => "invalid managed tool definition",
            Self::SchemaTooLarge => "managed tool schemas exceed the retained-byte limit",
            Self::InvalidInputSchema => "managed tool input schema failed admission",
            Self::InvalidOutputSchema => "managed tool output schema failed admission",
            Self::RequestMismatch => "request does not match the bound modern tool",
            Self::InvalidArguments => "tool arguments do not satisfy the admitted input schema",
            Self::InvalidResult => "managed tool result failed protocol admission",
            Self::MissingStructuredOutput => {
                "successful tool result omitted required structured output"
            }
            Self::InvalidStructuredOutput => {
                "tool structured output does not satisfy its admitted schema"
            }
            Self::HeaderBindingMismatch => "header review does not match the bound tool contract",
            Self::Invalidated => "managed tool contract has been invalidated",
            Self::Closed => "managed tool call is closed",
            Self::Headers(error) => return fmt::Display::fmt(error, f),
            Self::Core(error) => return fmt::Display::fmt(error, f),
        })
    }
}

impl std::error::Error for ManagedToolError {}

impl From<ManagedCoreError> for ManagedToolError {
    fn from(error: ManagedCoreError) -> Self {
        Self::Core(error)
    }
}

impl From<ManagedToolHeaderError> for ManagedToolError {
    fn from(error: ManagedToolHeaderError) -> Self {
        match error {
            ManagedToolHeaderError::Headers(error) => Self::Headers(error),
            ManagedToolHeaderError::Core(error) => Self::Core(error),
        }
    }
}

struct ToolContract {
    name: String,
    input: AdmittedToolHeaderSchema,
    output: Option<AdmittedSchema>,
    invalidated: AtomicBool,
    invalidation: McpRequestCancellation,
    catalog_invalidated: Option<Arc<AtomicBool>>,
    // A repair-only contract keeps its original validity owner alive. It is
    // never published as a reusable client, and inheritance is one hop only.
    source_contract: Option<Arc<ToolContract>>,
}

impl ToolContract {
    fn admit(tool: FinalTool) -> Result<Self, ManagedToolError> {
        if tool.name.is_empty()
            || tool.name.len() > MAX_MANAGED_TOOL_NAME_BYTES
            || tool.name.bytes().any(|byte| byte < 0x20 || byte == 0x7f)
        {
            return Err(ManagedToolError::InvalidDefinition);
        }
        // FinalTool has public fields: constructing it in Rust must not bypass
        // the root-shape checks otherwise performed by its wire deserializer.
        if tool.input_schema.get("type").and_then(Value::as_str) != Some("object") {
            return Err(ManagedToolError::InvalidInputSchema);
        }
        if tool
            .output_schema
            .as_ref()
            .is_some_and(|schema| !schema.is_object())
        {
            return Err(ManagedToolError::InvalidOutputSchema);
        }
        // Standard header annotations are part of the tool definition, not
        // permission to disclose arguments. Admit their syntax while retaining
        // the exact source and validating through the ordinary schema engine.
        // Admission itself never executes this plan or approves disclosure.
        let input = AdmittedToolHeaderSchema::admit(tool.input_schema)
            .map_err(|_| ManagedToolError::InvalidInputSchema)?;
        let output = tool
            .output_schema
            .map(admit_final_schema)
            .transpose()
            .map_err(|_| ManagedToolError::InvalidOutputSchema)?;
        // Shared admission bounds nesting/nodes before serialization. Counting
        // does not allocate another copy of potentially large schema strings.
        let mut bytes = SchemaBytes(0);
        serde_json::to_writer(&mut bytes, input.schema())
            .map_err(|_| ManagedToolError::SchemaTooLarge)?;
        if let Some(output) = &output {
            serde_json::to_writer(&mut bytes, output.schema())
                .map_err(|_| ManagedToolError::SchemaTooLarge)?;
        }
        Ok(Self {
            name: tool.name,
            input,
            output,
            invalidated: AtomicBool::new(false),
            invalidation: McpRequestCancellation::new(),
            catalog_invalidated: None,
            source_contract: None,
        })
    }

    fn invalidate(&self) {
        self.invalidated.store(true, Ordering::Release);
        self.invalidation.cancel();
    }

    fn is_invalidated(&self) -> bool {
        self.invalidated.load(Ordering::Acquire)
            || self
                .catalog_invalidated
                .as_ref()
                .is_some_and(|flag| flag.load(Ordering::Acquire))
            || self
                .source_contract
                .as_ref()
                .is_some_and(|source| source.is_invalidated())
    }

    fn check(&self) -> Result<(), ManagedToolError> {
        if self.is_invalidated() {
            Err(ManagedToolError::Invalidated)
        } else {
            Ok(())
        }
    }

    fn validate_request(&self, request: &CoreRequest) -> Result<(), ManagedToolError> {
        self.check()?;
        if request.era() != ProtocolEra::Modern2026 || request.method() != "tools/call" {
            return Err(ManagedToolError::RequestMismatch);
        }
        let params = request
            .encode_params()
            .map_err(|_| ManagedToolError::InvalidArguments)?
            .ok_or(ManagedToolError::InvalidArguments)?;
        if params.get("name").and_then(Value::as_str) != Some(self.name.as_str()) {
            return Err(ManagedToolError::RequestMismatch);
        }
        // Omission means no arguments, not permission to skip required fields.
        // Do not rewrite the request: absent and explicit-empty retain their
        // original wire representation, including continuation metadata.
        let empty = Value::Object(serde_json::Map::new());
        let arguments = params.get("arguments").unwrap_or(&empty);
        if !arguments.is_object() {
            return Err(ManagedToolError::InvalidArguments);
        }
        self.input
            .validate(arguments)
            .map_err(|_| ManagedToolError::InvalidArguments)?;
        self.check()
    }

    fn validate_result(&self, result: &CoreResult) -> Result<(), ManagedToolError> {
        self.check()?;
        let Some(output) = &self.output else {
            return Ok(());
        };
        // This is called only on the method-owned result produced by the
        // bound call's protocol decoder, never on an arbitrary result supplied
        // by the application. Keep and return that original lossless result.
        let encoded = result
            .encode()
            .map_err(|_| ManagedToolError::InvalidResult)?;
        let value: Value =
            serde_json::from_str(&encoded).map_err(|_| ManagedToolError::InvalidResult)?;
        match value.get("resultType").and_then(Value::as_str) {
            // A suspended invocation has not produced the tool's output yet.
            Some("input_required") => return self.check(),
            Some("complete") => {}
            _ => return Err(ManagedToolError::InvalidResult),
        }
        // A tool-level execution error is not successful structured output.
        // Preserve it as a typed tool result, not a client validation failure.
        if value.get("isError").and_then(Value::as_bool) == Some(true) {
            return self.check();
        }
        let structured = value
            .get("structuredContent")
            .ok_or(ManagedToolError::MissingStructuredOutput)?;
        output
            .validate(structured)
            .map_err(|_| ManagedToolError::InvalidStructuredOutput)?;
        self.check()
    }
}

struct SchemaBytes(usize);

impl Write for SchemaBytes {
    fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
        if buffer.len() > MAX_MANAGED_TOOL_SCHEMA_BYTES.saturating_sub(self.0) {
            return Err(io::Error::other("managed tool schema byte limit"));
        }
        self.0 += buffer.len();
        Ok(buffer.len())
    }

    fn flush(&mut self) -> io::Result<()> {
        Ok(())
    }
}

/// One tool definition bound to one managed login. Clones share the same
/// immutable schemas and irreversible invalidation state, not active requests.
/// Definitions/annotations are not treated as authorization or retry policy.
#[derive(Clone)]
pub struct ManagedToolClient {
    session: ManagedOAuthSession,
    contract: Arc<ToolContract>,
    header_review: Option<Arc<ReviewedToolHeaders>>,
}

impl fmt::Debug for ManagedToolClient {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.debug_struct("ManagedToolClient")
            .field("has_output_schema", &self.contract.output.is_some())
            .field("parameter_headers", &self.header_review.is_some())
            .field("invalidated", &self.is_invalidated())
            .finish_non_exhaustive()
    }
}

impl ManagedToolClient {
    /// Admits both schemas before retaining a tool binding. There is no
    /// discovery, credential acquisition, network schema resolution or I/O.
    /// The host is responsible for obtaining this definition from the intended
    /// server and approving its use with this exact managed login.
    pub fn new(session: ManagedOAuthSession, tool: FinalTool) -> Result<Self, ManagedToolError> {
        let contract = Arc::new(ToolContract::admit(tool)?);
        Ok(Self {
            session,
            contract,
            header_review: None,
        })
    }

    pub fn tool_name(&self) -> &str {
        &self.contract.name
    }

    /// Refuses newly started calls and later publication through every clone.
    /// Calls admitted before invalidation may already be dispatching.
    /// Already-delivered events and server side effects cannot be recalled.
    /// Wakes pending requests, reads and continuations through every clone.
    /// Their owners must poll or drop them to release resources; no background
    /// worker is created. The caller's cancellation handle and shared session
    /// are not cancelled. Abandoning an in-flight OAuth renewal still retains
    /// the session's existing fail-closed refresh-lineage policy and may require
    /// a new login. A new definition requires a new client.
    pub fn invalidate(&self) {
        self.contract.invalidate();
    }

    pub fn is_invalidated(&self) -> bool {
        self.contract.is_invalidated()
    }

    /// Validates an invocation without acquiring credentials or dispatching it.
    /// The exact name, protocol era and full argument schema must match.
    pub fn validate_request(&self, request: &CoreRequest) -> Result<(), ManagedToolError> {
        self.contract.validate_request(request)
    }

    /// Opens one schema-checked tool call. `input_required` is returned intact;
    /// it is not treated as a successful output or as an automatic retry.
    pub async fn request(
        &self,
        cx: &Cx,
        request: CoreRequest,
        request_id: RequestId,
        limits: ManagedCoreLimits,
    ) -> Result<ManagedToolCall, ManagedToolError> {
        Box::pin(self.request_with_cancellation(
            cx,
            &McpRequestCancellation::new(),
            request,
            request_id,
            limits,
        ))
        .await
    }

    pub async fn request_with_cancellation(
        &self,
        cx: &Cx,
        cancellation: &McpRequestCancellation,
        request: CoreRequest,
        request_id: RequestId,
        limits: ManagedCoreLimits,
    ) -> Result<ManagedToolCall, ManagedToolError> {
        check_tool_call(cx, cancellation, &self.contract)?;
        self.contract.validate_request(&request)?;
        check_tool_call(cx, cancellation, &self.contract)?;
        let call = match self.header_review.as_deref() {
            Some(reviewed) => {
                Box::pin(await_validity(
                    cx,
                    cancellation,
                    &self.contract,
                    self.session.request_tool_with_headers_and_cancellation(
                        cx,
                        cancellation,
                        request,
                        request_id,
                        reviewed,
                        limits,
                    ),
                ))
                .await??
            }
            None => {
                Box::pin(await_validity(
                    cx,
                    cancellation,
                    &self.contract,
                    self.session.request_core_with_cancellation(
                        cx,
                        cancellation,
                        request,
                        request_id,
                        limits,
                    ),
                ))
                .await??
            }
        };
        check_tool_call(cx, cancellation, &self.contract)?;
        Ok(ManagedToolCall {
            call: Some(call),
            contract: self.contract.clone(),
            cancellation: cancellation.clone(),
            finished: false,
        })
    }
}

/// Incremental schema-checked call. Failed validation closes only this call;
/// malformed output never becomes a published success. Dropping an in-progress
/// read drops its owned response rather than making partial framing reusable.
pub struct ManagedToolCall {
    call: Option<ManagedCoreCall>,
    contract: Arc<ToolContract>,
    cancellation: McpRequestCancellation,
    finished: bool,
}

impl ManagedToolCall {
    pub fn close(&mut self) {
        self.call = None;
    }

    pub async fn next_event(
        &mut self,
        cx: &Cx,
    ) -> Result<Option<ManagedCoreEvent>, ManagedToolError> {
        if self.finished {
            return Ok(None);
        }
        let mut call = self.call.take().ok_or(ManagedToolError::Closed)?;
        check_tool_call(cx, &self.cancellation, &self.contract)?;
        let event = Box::pin(await_validity(
            cx,
            &self.cancellation,
            &self.contract,
            call.next_event(cx),
        ))
        .await??
        .ok_or(ManagedCoreError::MissingTerminal)?;
        check_tool_call(cx, &self.cancellation, &self.contract)?;
        match &event {
            ManagedCoreEvent::Result(result) => {
                self.contract.validate_result(result)?;
                check_tool_call(cx, &self.cancellation, &self.contract)?;
                self.finished = true;
            }
            ManagedCoreEvent::Notification(_) => self.call = Some(call),
        }
        Ok(Some(event))
    }
}

// Recheck the retained request-local domain around synchronous schema work,
// not just the caller Cx. The transport still owns its original time budgets.
fn check_tool_call(
    cx: &Cx,
    cancellation: &McpRequestCancellation,
    contract: &ToolContract,
) -> Result<(), ManagedToolError> {
    if cancellation.is_cancel_requested() || cx.checkpoint().is_err() {
        return Err(ManagedCoreError::Cancelled.into());
    }
    contract.check()
}

#[cfg(test)]
mod tests;

#[cfg(test)]
mod header_schema_tests;