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
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
//! LEG-NEG-01 B — the HTTP `Auto` fallback coordinator.
//!
//! This module owns exactly one decision: given a typed observation of the one
//! modern HTTP discovery probe, may the client open a single legacy SSE `GET`
//! against its configured endpoint, and does the first event on that stream
//! actually select the exact-2024 era?
//!
//! It deliberately does **not** classify probe responses. `ClientHttpNegotiation`
//! (HTTP-03 / CLT-02) remains the sole observation producer, and this
//! coordinator consumes its [`ClientHttpNegotiationDecision`] rather than
//! restating the eligibility table. The dependency runs one way, `leg_neg` →
//! `negotiation`, so the classifier never learns that a coordinator exists.
//!
//! The eligibility matrix this coordinator enforces, under
//! [`ProtocolPolicy::Auto`] only:
//!
//! | status | `Empty` | `Unrecognized` | `RecognizedModernJsonRpc` |
//! |--------|---------|----------------|---------------------------|
//! | 400    | one GET | one GET        | no GET                    |
//! | 404    | one GET | one GET        | no GET                    |
//! | 405    | one GET | one GET        | no GET                    |
//!
//! Six rows authorize exactly one `GET`; the three recognized-modern rows
//! forbid it at every status, because a recognized modern JSON-RPC body — result
//! or error — is proof the peer speaks the modern era and must never be read as
//! a downgrade signal. Every other status/body combination is ineligible.
//!
//! Authorization is not selection. A permitted `GET` still selects nothing: only
//! the first valid `endpoint` event admitted from that stream moves the era to
//! [`ProtocolEra::Legacy2024`]. The event names the configured message POST
//! endpoint, not the SSE GET endpoint. Missing, malformed, duplicate, late, or
//! wrong-target events are refused with a typed error and leave every observable
//! field untouched.

use fastmcp_core::CanonicalHttpUrl;
use fastmcp_protocol::protocol_policy::{
    HttpEndpointBundleKey, HttpModernProbe, HttpProbeBody, ProtocolEra, ProtocolPolicy,
};
use std::sync::Arc;

use crate::negotiation::{
    ClientHttpNegotiation, ClientHttpNegotiationDecision, ClientHttpNegotiationError,
};
use crate::session::ClientProtocolPlan;

/// Typed refusal from the HTTP fallback coordinator.
///
/// Every variant is reached before any effect: no legacy `GET` is opened, no
/// credential is acquired or mutated, no era or discovery cache entry is
/// written, and no era is selected.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum HttpFallbackError {
    /// Only `Auto` may coordinate a fallback. `ModernOnly` has nothing to fall
    /// back to and `LegacyOnly` never probes the modern era first.
    PolicyForbidsFallback {
        /// The immutable policy that refused.
        policy: ProtocolPolicy,
    },
    /// The plan carried no configured HTTP endpoint bundle.
    MissingHttpEndpointBundle,
    /// The plan carried no configured modern POST target.
    MissingModernPostTarget,
    /// The plan carried no configured legacy SSE GET target.
    MissingLegacySseTarget,
    /// The plan carried no configured legacy message POST target.
    MissingLegacyMessagePostTarget,
    /// The classifier refused to start an attempt for this plan.
    Negotiation(ClientHttpNegotiationError),
    /// An observation must name a nonempty modern target and a nonzero attempt.
    InvalidObservation,
    /// The observation or GET permit belongs to a different endpoint bundle or
    /// coordinator instance.
    ///
    /// Bundle identity covers the complete canonical targets, the credential and
    /// security partitions, the transport profile, and the policy/configuration
    /// generations. Permits additionally bind their exact issuing coordinator:
    /// two attempts with identical configuration and IDs cannot exchange them.
    CrossBundleObservation,
    /// The observation names a different modern POST target than the plan.
    ModernTargetMismatch,
    /// This coordinator already settled a different attempt's observation.
    ObservationAlreadyAdmitted {
        /// The attempt that was admitted first.
        admitted_attempt: u64,
    },
    /// The same attempt's observation was replayed after it settled.
    ReplayedObservation {
        /// The attempt identity that was replayed.
        attempt_id: u64,
    },
    /// The status/body row cannot authorize a legacy `GET`.
    IneligibleObservation {
        /// The observed HTTP status.
        status: u16,
        /// The observed body classification.
        body: HttpProbeBody,
    },
    /// This coordinator has not authorized a legacy GET. In particular, a
    /// recognized modern response must never be bypassed by a foreign permit.
    LegacyGetNotAuthorized,
    /// The single authorized `GET` was already opened.
    LegacyGetAlreadyOpened,
    /// An `endpoint` event arrived without an authorized, opened `GET`.
    EndpointEventWithoutAuthorization,
    /// The `endpoint` event did not carry a bounded canonical HTTP URL without
    /// userinfo or a fragment.
    EndpointEventMalformed,
    /// The advertised endpoint is not the configured legacy message POST target.
    EndpointEventTargetMismatch,
    /// A second `endpoint` event arrived after the era was already selected.
    DuplicateEndpointEvent,
}

impl std::fmt::Display for HttpFallbackError {
    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            Self::PolicyForbidsFallback { policy } => {
                write!(formatter, "{policy:?} cannot coordinate an HTTP fallback")
            }
            Self::MissingHttpEndpointBundle => {
                formatter.write_str("the plan has no configured HTTP endpoint bundle")
            }
            Self::MissingModernPostTarget => {
                formatter.write_str("the plan has no configured modern POST target")
            }
            Self::MissingLegacySseTarget => {
                formatter.write_str("the plan has no configured legacy SSE GET target")
            }
            Self::MissingLegacyMessagePostTarget => {
                formatter.write_str("the plan has no configured legacy message POST target")
            }
            Self::Negotiation(error) => write!(formatter, "probe classification refused: {error}"),
            Self::InvalidObservation => {
                formatter.write_str("an observation needs a nonempty target and nonzero attempt")
            }
            Self::CrossBundleObservation => {
                formatter.write_str("the observation or permit belongs to another coordinator")
            }
            Self::ModernTargetMismatch => {
                formatter.write_str("the observation names another modern POST target")
            }
            Self::ObservationAlreadyAdmitted { admitted_attempt } => write!(
                formatter,
                "attempt {admitted_attempt} was already admitted by this coordinator"
            ),
            Self::ReplayedObservation { attempt_id } => {
                write!(
                    formatter,
                    "attempt {attempt_id} was replayed after settling"
                )
            }
            Self::IneligibleObservation { status, body } => write!(
                formatter,
                "status {status} with {body:?} cannot authorize a legacy GET"
            ),
            Self::LegacyGetNotAuthorized => {
                formatter.write_str("this coordinator has not authorized a legacy GET")
            }
            Self::LegacyGetAlreadyOpened => {
                formatter.write_str("the one authorized legacy GET was already opened")
            }
            Self::EndpointEventWithoutAuthorization => {
                formatter.write_str("an endpoint event requires an opened authorized GET")
            }
            Self::EndpointEventMalformed => {
                formatter.write_str("the endpoint event carried no usable target")
            }
            Self::EndpointEventTargetMismatch => formatter
                .write_str("the advertised endpoint is not the configured message POST target"),
            Self::DuplicateEndpointEvent => {
                formatter.write_str("the era was already selected by an earlier endpoint event")
            }
        }
    }
}

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

/// One typed observation of the single modern HTTP discovery probe.
///
/// The observation is bound to the exact attempt that produced it: the complete
/// endpoint bundle key, the exact configured modern target, and a caller-owned
/// nonzero attempt identity. A coordinator refuses any observation whose binding
/// does not match, which is what makes a stale, cross-bundle, or replayed
/// observation inert rather than merely unlikely.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ModernProbeObservation {
    bundle_key: HttpEndpointBundleKey,
    modern_target: String,
    attempt_id: u64,
    probe: HttpModernProbe,
}

impl ModernProbeObservation {
    /// Binds one probe observation to the attempt that produced it.
    pub fn new(
        bundle_key: HttpEndpointBundleKey,
        modern_target: impl Into<String>,
        attempt_id: u64,
        probe: HttpModernProbe,
    ) -> Result<Self, HttpFallbackError> {
        let modern_target = modern_target.into();
        if modern_target.is_empty() || attempt_id == 0 {
            return Err(HttpFallbackError::InvalidObservation);
        }
        Ok(Self {
            bundle_key,
            modern_target,
            attempt_id,
            probe,
        })
    }

    /// Returns the endpoint bundle identity this observation is bound to.
    #[must_use]
    pub const fn bundle_key(&self) -> &HttpEndpointBundleKey {
        &self.bundle_key
    }

    /// Returns the exact configured modern POST target.
    #[must_use]
    pub fn modern_target(&self) -> &str {
        &self.modern_target
    }

    /// Returns the caller-owned attempt identity.
    #[must_use]
    pub const fn attempt_id(&self) -> u64 {
        self.attempt_id
    }

    /// Returns the observed probe status and body classification.
    #[must_use]
    pub const fn probe(&self) -> HttpModernProbe {
        self.probe
    }
}

/// Single-use authorization for exactly one legacy SSE `GET`.
///
/// The type is deliberately neither `Clone` nor `Copy`, and
/// [`HttpFallbackCoordinator::open_legacy_get`] consumes it by value, so a
/// second `GET` cannot be opened from one authorization even by mistake.
/// Holding a permit is not an era selection and never mutates coordinator state.
/// Its private allocation identity binds it to its exact issuing coordinator,
/// even when another coordinator uses the same endpoint bundle and attempt ID.
#[derive(Debug)]
pub struct LegacyGetPermit {
    target: String,
    attempt_id: u64,
    owner: Arc<()>,
}

impl PartialEq for LegacyGetPermit {
    fn eq(&self, other: &Self) -> bool {
        Arc::ptr_eq(&self.owner, &other.owner)
            && self.target == other.target
            && self.attempt_id == other.attempt_id
    }
}

impl Eq for LegacyGetPermit {}

impl LegacyGetPermit {
    /// Returns the configured legacy SSE target this permit authorizes.
    #[must_use]
    pub fn target(&self) -> &str {
        &self.target
    }

    /// Returns the attempt identity that earned this authorization.
    #[must_use]
    pub const fn attempt_id(&self) -> u64 {
        self.attempt_id
    }
}

/// The coordinator's decision for one admitted observation.
#[derive(Debug, PartialEq, Eq)]
pub enum FallbackDecision {
    /// No legacy `GET` is permitted; the modern observation stands.
    ModernRetained,
    /// Exactly one legacy SSE `GET` is authorized.
    LegacyGetAuthorized(LegacyGetPermit),
}

/// Observable progress counters of one coordinator.
///
/// Tests compare this value and the separately exposed admitted POST target
/// before and after a refusal to prove an invalid input changed nothing.
/// `credential_mutations` and `era_cache_mutations` are always zero: this
/// coordinator has no authority to acquire credentials or write an era cache.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct FallbackState {
    /// Observations admitted by this coordinator, at most one.
    pub observations_admitted: usize,
    /// Legacy `GET` authorizations issued, at most one.
    pub legacy_gets_authorized: usize,
    /// Legacy `GET`s actually opened, at most one.
    pub legacy_gets_opened: usize,
    /// `endpoint` events admitted, at most one.
    pub endpoint_events_admitted: usize,
    /// The era selected, only ever by a valid first `endpoint` event.
    pub selected_era: Option<ProtocolEra>,
    /// Credential acquisitions or mutations performed. Always zero.
    pub credential_mutations: usize,
    /// Era or discovery cache entries written. Always zero.
    pub era_cache_mutations: usize,
}

/// The sole HTTP `Auto` fallback coordinator.
///
/// One coordinator serves one connection attempt: it admits at most one
/// observation, issues at most one `GET` authorization, and admits at most one
/// `endpoint` event. The admitted message POST URL is retained only after its
/// syntax and immutable route binding pass; callers need not reuse peer bytes.
/// Debug output omits URLs, including server-generated session query values.
pub struct HttpFallbackCoordinator {
    plan: ClientProtocolPlan,
    bundle_key: HttpEndpointBundleKey,
    modern_target: String,
    legacy_sse_target: String,
    legacy_message_post_target: String,
    admitted_message_post_target: Option<CanonicalHttpUrl>,
    state: FallbackState,
    settled_attempt: Option<u64>,
    permit_owner: Arc<()>,
}

impl std::fmt::Debug for HttpFallbackCoordinator {
    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        formatter
            .debug_struct("HttpFallbackCoordinator")
            .field("state", &self.state)
            .field("settled_attempt", &self.settled_attempt)
            .finish_non_exhaustive()
    }
}

impl HttpFallbackCoordinator {
    /// Starts one coordinator from an immutable `Auto` HTTP plan.
    ///
    /// This is a side-effect-free admission boundary: it opens no socket and
    /// acquires no credential.
    pub fn new(plan: ClientProtocolPlan) -> Result<Self, HttpFallbackError> {
        let policy = plan.policy();
        if !matches!(policy, ProtocolPolicy::Auto) {
            return Err(HttpFallbackError::PolicyForbidsFallback { policy });
        }
        let bundle_key = plan
            .http_endpoints()
            .ok_or(HttpFallbackError::MissingHttpEndpointBundle)?
            .key();
        let modern_target = plan
            .modern_post_target()
            .ok_or(HttpFallbackError::MissingModernPostTarget)?
            .to_owned();
        let legacy_sse_target = plan
            .legacy_sse_target()
            .ok_or(HttpFallbackError::MissingLegacySseTarget)?
            .to_owned();
        let legacy_message_post_target = plan
            .legacy_message_post_target()
            .ok_or(HttpFallbackError::MissingLegacyMessagePostTarget)?
            .to_owned();
        // Prove the classifier will accept this plan now, so a later refusal
        // cannot be mistaken for an eligibility outcome.
        ClientHttpNegotiation::from_protocol_plan(&plan).map_err(HttpFallbackError::Negotiation)?;
        Ok(Self {
            plan,
            bundle_key,
            modern_target,
            legacy_sse_target,
            legacy_message_post_target,
            admitted_message_post_target: None,
            state: FallbackState::default(),
            settled_attempt: None,
            permit_owner: Arc::new(()),
        })
    }

    /// Returns the observable progress counters.
    #[must_use]
    pub const fn state(&self) -> FallbackState {
        self.state
    }

    /// Returns the era selected so far, if any.
    #[must_use]
    pub const fn selected_era(&self) -> Option<ProtocolEra> {
        self.state.selected_era
    }

    /// Returns the exact configured legacy SSE target.
    #[must_use]
    pub fn legacy_sse_target(&self) -> &str {
        &self.legacy_sse_target
    }

    /// Returns the configured legacy message POST target, distinct from GET.
    #[must_use]
    pub fn legacy_message_post_target(&self) -> &str {
        &self.legacy_message_post_target
    }

    /// Returns the admitted POST target, including any permitted session query.
    ///
    /// This is absent until a valid endpoint event selects the legacy era.
    /// Refused and duplicate events cannot replace it. Treat a session query
    /// as sensitive routing data rather than a diagnostic or cache identity.
    #[must_use]
    pub fn advertised_message_post_target(&self) -> Option<&str> {
        self.admitted_message_post_target
            .as_ref()
            .map(CanonicalHttpUrl::as_str)
    }

    /// Returns the endpoint bundle identity this coordinator is bound to.
    #[must_use]
    pub const fn bundle_key(&self) -> &HttpEndpointBundleKey {
        &self.bundle_key
    }

    /// Classifies one probe through the shipped observation producer.
    ///
    /// A fresh attempt is used for every call so that a refused classification
    /// leaves no retained state anywhere — neither here nor in the classifier.
    /// This is also why the eligibility table is not restated in this module.
    fn classify(
        &self,
        probe: HttpModernProbe,
    ) -> Result<ClientHttpNegotiationDecision, ClientHttpNegotiationError> {
        let mut negotiation = ClientHttpNegotiation::from_protocol_plan(&self.plan)?;
        negotiation.observe_modern_probe(probe)
    }

    /// Admits one typed probe observation and applies the eligibility matrix.
    ///
    /// Binding is checked before classification, so an observation from another
    /// bundle, another modern target, or a settled attempt is inert.
    pub fn observe(
        &mut self,
        observation: &ModernProbeObservation,
    ) -> Result<FallbackDecision, HttpFallbackError> {
        if observation.bundle_key != self.bundle_key {
            return Err(HttpFallbackError::CrossBundleObservation);
        }
        if observation.modern_target != self.modern_target {
            return Err(HttpFallbackError::ModernTargetMismatch);
        }
        if let Some(settled) = self.settled_attempt {
            return Err(if settled == observation.attempt_id {
                HttpFallbackError::ReplayedObservation {
                    attempt_id: settled,
                }
            } else {
                HttpFallbackError::ObservationAlreadyAdmitted {
                    admitted_attempt: settled,
                }
            });
        }

        let probe = observation.probe;
        let decision =
            self.classify(probe)
                .map_err(|_| HttpFallbackError::IneligibleObservation {
                    status: probe.status,
                    body: probe.body,
                })?;

        // Only a settled decision mutates this coordinator.
        self.settled_attempt = Some(observation.attempt_id);
        self.state.observations_admitted += 1;
        match decision {
            // A recognized modern JSON-RPC body — result or error — is proof of
            // the modern era at every eligible status, so it forbids the GET.
            ClientHttpNegotiationDecision::ModernSelected => Ok(FallbackDecision::ModernRetained),
            ClientHttpNegotiationDecision::LegacySseFallbackAuthorized => {
                self.state.legacy_gets_authorized += 1;
                Ok(FallbackDecision::LegacyGetAuthorized(LegacyGetPermit {
                    target: self.legacy_sse_target.clone(),
                    attempt_id: observation.attempt_id,
                    owner: Arc::clone(&self.permit_owner),
                }))
            }
        }
    }

    /// Consumes the single authorization and opens the one permitted `GET`.
    ///
    /// Returns the exact configured target the caller must request. Opening a
    /// `GET` is still not an era selection. A matching attempt number alone is
    /// not authority: the permit must have been issued by this coordinator.
    pub fn open_legacy_get(
        &mut self,
        permit: LegacyGetPermit,
    ) -> Result<String, HttpFallbackError> {
        if self.state.legacy_gets_opened > 0 {
            return Err(HttpFallbackError::LegacyGetAlreadyOpened);
        }
        if self.state.legacy_gets_authorized != 1 {
            return Err(HttpFallbackError::LegacyGetNotAuthorized);
        }
        if !Arc::ptr_eq(&self.permit_owner, &permit.owner)
            || Some(permit.attempt_id) != self.settled_attempt
            || permit.target != self.legacy_sse_target
        {
            return Err(HttpFallbackError::CrossBundleObservation);
        }
        self.state.legacy_gets_opened += 1;
        Ok(self.legacy_sse_target.clone())
    }

    /// Admits the first valid `endpoint` event from the opened `GET`.
    ///
    /// Only this call may select [`ProtocolEra::Legacy2024`]. The advertised
    /// target must be the configured message POST target, byte for byte, or the
    /// configured query-free target extended with a server-generated query.
    /// Scheme, authority and path never change. Admission never trims or
    /// silently repairs peer bytes, and never accepts userinfo or fragments.
    pub fn admit_endpoint_event(
        &mut self,
        advertised: &str,
    ) -> Result<ProtocolEra, HttpFallbackError> {
        if self.state.legacy_gets_opened == 0 {
            return Err(HttpFallbackError::EndpointEventWithoutAuthorization);
        }
        if self.state.endpoint_events_admitted > 0 {
            return Err(HttpFallbackError::DuplicateEndpointEvent);
        }
        let target = CanonicalHttpUrl::parse(advertised)
            .map_err(|_| HttpFallbackError::EndpointEventMalformed)?;
        if target.as_str() != advertised || target.has_userinfo() || target.fragment().is_some() {
            return Err(HttpFallbackError::EndpointEventMalformed);
        }
        if !advertised_target_is_admissible(&self.legacy_message_post_target, target.as_str()) {
            return Err(HttpFallbackError::EndpointEventTargetMismatch);
        }
        self.admitted_message_post_target = Some(target);
        self.state.endpoint_events_admitted += 1;
        self.state.selected_era = Some(ProtocolEra::Legacy2024);
        Ok(ProtocolEra::Legacy2024)
    }
}

/// Binds an already syntax-admitted message URL to the configured POST route.
///
/// Byte equality admits. A configured target with no query component may
/// additionally be extended by a nonempty server-generated query, because the
/// exact 2024-11-05 lane advertises a session query no client can preconfigure.
/// A configured query is immutable; it cannot be replaced or extended by the
/// peer. Syntax admission precedes this comparison, including the byte bound.
fn advertised_target_is_admissible(configured: &str, advertised: &str) -> bool {
    if advertised == configured {
        return true;
    }
    match advertised.split_once('?') {
        Some((base, query)) => !query.is_empty() && base == configured && !configured.contains('?'),
        None => false,
    }
}