From 06b54aaf78fb00a50a9e42c1d859b308bda8d858 Mon Sep 17 00:00:00 2001 From: Serhii Lizunov <34252204+LizunovSergey@users.noreply.github.com> Date: Fri, 11 Sep 2026 00:02:24 +0200 Subject: [PATCH] fix(transport): stop keeping a session and GET stream at 2026-07-28 SEP-2567 removes sessions and the standalone GET endpoint at 2026-07-28, so an Mcp-Session-Id and a GET stream are both artifacts of a pre-2026-07-28 server shape. The modern startup path already gets this right: a server/discover bootstrap sets session_id to None unconditionally. The legacy-shaped paths did not, and all three spawn_common_stream call sites were guarded only on a session id being present, never on the negotiated version: - legacy startup, where a server can answer legacy initialize with both a session id and protocolVersion 2026-07-28; - the fallback initialize taken after server/discover fails; - session re-establishment after an expired-session 404, which respawned on the new id without reconsulting the version. Gating happens where a session id is adopted rather than at each spawn: session_id_for_version drops the id once the negotiated version has no sessions, which leaves all three call sites with nothing to open and stops the id being echoed on requests. PeerRequestAssociation Unassociated then becomes unreachable over streamable HTTP instead of merely rejected, and the receive-side enforcement stays as the defence for anything that still slips through. A volunteered id is dropped with a warning rather than being treated as fatal, and the id is still recorded for the shutdown DELETE, so a session the server really did create is torn down. --- crates/rmcp/Cargo.toml | 10 + .../src/transport/streamable_http_client.rs | 34 +++ ...est_streamable_http_sessionless_version.rs | 214 ++++++++++++++++++ 3 files changed, 258 insertions(+) create mode 100644 crates/rmcp/tests/test_streamable_http_sessionless_version.rs diff --git a/crates/rmcp/Cargo.toml b/crates/rmcp/Cargo.toml index ac51ebfcd..2ccdba34e 100644 --- a/crates/rmcp/Cargo.toml +++ b/crates/rmcp/Cargo.toml @@ -338,6 +338,16 @@ required-features = [ ] path = "tests/test_discover_http_client_startup.rs" +[[test]] +name = "test_streamable_http_sessionless_version" +required-features = [ + "client", + "reqwest", + "transport-streamable-http-client-reqwest", + "transport-streamable-http-server", +] +path = "tests/test_streamable_http_sessionless_version.rs" + [[test]] name = "test_streamable_http_standard_headers" required-features = ["server", "client", "transport-streamable-http-server", "reqwest"] diff --git a/crates/rmcp/src/transport/streamable_http_client.rs b/crates/rmcp/src/transport/streamable_http_client.rs index 6f27abf90..276661580 100644 --- a/crates/rmcp/src/transport/streamable_http_client.rs +++ b/crates/rmcp/src/transport/streamable_http_client.rs @@ -90,6 +90,37 @@ fn request_version_headers( (version, headers) } +/// Decides whether a session id survives version negotiation. +/// +/// SEP-2567 removes sessions and the standalone GET endpoint at +/// [`ProtocolVersion::STANDARD_HEADERS`], so at that version an `Mcp-Session-Id` and a GET +/// stream are both artifacts of a pre-`2026-07-28` server shape. A legacy-shaped handshake +/// can still answer with a session id while negotiating that version; the id is dropped +/// rather than echoed, which also leaves every `spawn_common_stream` call site — all three +/// of which are guarded on a session id being present — with no stream to open. +/// +/// Dropping is deliberate rather than fatal: refusing to start would break clients against +/// servers that work today, and the receive-side enforcement added for SEP-2260 still +/// rejects anything that reaches the client over a stream it should not have. The caller +/// keeps the original id for the shutdown `DELETE`, so a session the server really did +/// create is still torn down. +fn session_id_for_version( + session_id: Option>, + negotiated_version: &ProtocolVersion, +) -> Option> { + if negotiated_version < &ProtocolVersion::STANDARD_HEADERS { + return session_id; + } + if session_id.is_some() { + tracing::warn!( + version = negotiated_version.as_str(), + "server returned an Mcp-Session-Id while negotiating a version that has no sessions; \ + the id will not be sent on requests and no standalone GET stream will be opened" + ); + } + None +} + fn cache_tools_from_response( cache: &mut HashMap>, message: &mut ServerJsonRpcMessage, @@ -1133,6 +1164,7 @@ impl Worker for StreamableHttpClientWorker { auth_header: config.auth_header.clone(), protocol_headers: protocol_headers.clone(), }); + session_id = session_id_for_version(session_id, &negotiated_version); context.send_to_handler(message).await?; if is_legacy_startup { @@ -1261,6 +1293,7 @@ impl Worker for StreamableHttpClientWorker { auth_header: config.auth_header.clone(), protocol_headers: protocol_headers.clone(), }); + session_id = session_id_for_version(session_id, &negotiated_version); // Do not send controls queued during recovery to the new session. context.advance_control_generation(); session_cancellation = CancellationToken::new(); @@ -1517,6 +1550,7 @@ impl Worker for StreamableHttpClientWorker { auth_header: config.auth_header.clone(), protocol_headers: protocol_headers.clone(), }); + session_id = session_id_for_version(session_id, &negotiated_version); context.send_to_handler(initialize_response).await?; awaiting_fallback_initialized = true; continue; diff --git a/crates/rmcp/tests/test_streamable_http_sessionless_version.rs b/crates/rmcp/tests/test_streamable_http_sessionless_version.rs new file mode 100644 index 000000000..ce95df369 --- /dev/null +++ b/crates/rmcp/tests/test_streamable_http_sessionless_version.rs @@ -0,0 +1,214 @@ +#![cfg(all( + not(feature = "local"), + feature = "client", + feature = "reqwest", + feature = "transport-streamable-http-server" +))] +//! SEP-2567 removes sessions and the standalone GET stream at 2026-07-28. A legacy-shaped +//! handshake can still answer with an `Mcp-Session-Id` while negotiating that version; the +//! client must not then echo the id or open the stream. + +use std::sync::{Arc, Mutex}; + +use axum::{ + Router, + body::{Body, Bytes}, + extract::State, + http::{HeaderMap, Response, StatusCode}, + routing::any, +}; +use rmcp::{ + ClientLifecycleMode, ClientServiceExt, + model::ClientInfo, + transport::{ + StreamableHttpClientTransport, streamable_http_client::StreamableHttpClientTransportConfig, + }, +}; +use serde_json::json; +use tokio_util::sync::CancellationToken; + +const SESSION_ID: &str = "session-the-server-should-not-have-issued"; + +/// One request the client made: HTTP method, JSON-RPC method, `Mcp-Session-Id` header. +type Call = (String, String, Option); + +#[derive(Clone, Default)] +struct Recorder { + seen: Arc>>, + negotiated_version: Arc>, +} + +impl Recorder { + fn calls(&self) -> Vec { + self.seen.lock().expect("recorder poisoned").clone() + } + + fn get_requests(&self) -> Vec { + self.calls() + .into_iter() + .filter(|(http_method, ..)| http_method == "GET") + .collect() + } + + fn session_headers_after_handshake(&self) -> Vec> { + self.calls() + .into_iter() + .filter(|(_, jsonrpc_method, _)| jsonrpc_method != "initialize") + .map(|(.., session)| session) + .collect() + } +} + +async fn handler( + State(state): State, + method: axum::http::Method, + headers: HeaderMap, + body: Bytes, +) -> Response { + let session = headers + .get("mcp-session-id") + .and_then(|value| value.to_str().ok()) + .map(str::to_owned); + + if method == axum::http::Method::GET { + state.seen.lock().expect("recorder poisoned").push(( + "GET".to_owned(), + "-".to_owned(), + session, + )); + // Hang so the client keeps the stream if it opens one; the test cancels it. + return Response::builder() + .status(StatusCode::METHOD_NOT_ALLOWED) + .body(Body::empty()) + .expect("build GET rejection"); + } + + let request: serde_json::Value = serde_json::from_slice(&body).expect("valid JSON-RPC body"); + let jsonrpc_method = request["method"].as_str().unwrap_or("-").to_owned(); + state.seen.lock().expect("recorder poisoned").push(( + "POST".to_owned(), + jsonrpc_method.clone(), + session, + )); + + if jsonrpc_method == "server/discover" { + // Force the legacy initialize path. + return Response::builder() + .status(StatusCode::NOT_FOUND) + .body(Body::empty()) + .expect("build discover rejection"); + } + + if jsonrpc_method == "initialize" { + let version = state + .negotiated_version + .lock() + .expect("recorder poisoned") + .clone(); + // The shape this issue is about: a session id alongside a negotiated version + // that, from 2026-07-28 on, has no sessions at all. + return Response::builder() + .status(StatusCode::OK) + .header("content-type", "application/json") + .header("mcp-session-id", SESSION_ID) + .body(Body::from( + json!({ + "jsonrpc": "2.0", + "id": request["id"], + "result": { + "protocolVersion": version, + "capabilities": {}, + "serverInfo": {"name": "dual-era", "version": "1.0"} + } + }) + .to_string(), + )) + .expect("build initialize response"); + } + + Response::builder() + .status(StatusCode::ACCEPTED) + .body(Body::empty()) + .expect("build notification response") +} + +async fn connect_and_record(negotiated_version: &str) -> Recorder { + let recorder = Recorder::default(); + *recorder + .negotiated_version + .lock() + .expect("recorder poisoned") = negotiated_version.to_owned(); + + let ct = CancellationToken::new(); + let router = Router::new() + .route("/mcp", any(handler)) + .with_state(recorder.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("listener should bind"); + let address = listener.local_addr().expect("listener address"); + let server = tokio::spawn({ + let ct = ct.clone(); + async move { + let _ = axum::serve(listener, router) + .with_graceful_shutdown(async move { ct.cancelled_owned().await }) + .await; + } + }); + + let transport = StreamableHttpClientTransport::from_config( + StreamableHttpClientTransportConfig::with_uri(format!("http://{address}/mcp")), + ); + let client = ClientInfo::default() + .serve_with_lifecycle(transport, ClientLifecycleMode::Initialize) + .await + .expect("client should start against a legacy handshake"); + + // Give a standalone GET stream, if one were opened, time to reach the server. + tokio::time::sleep(std::time::Duration::from_millis(300)).await; + client.cancel().await.expect("cancel client"); + ct.cancel(); + let _ = server.await; + recorder +} + +#[tokio::test] +async fn modern_version_drops_the_session_and_opens_no_stream() { + let recorder = connect_and_record("2026-07-28").await; + + // No standalone GET stream: SEP-2567 removed the endpoint at this version. + assert_eq!( + recorder.get_requests(), + Vec::new(), + "client opened a standalone GET stream at a version that has none" + ); + // And the id the server volunteered is not echoed back on anything. + let sessions = recorder.session_headers_after_handshake(); + assert!( + !sessions.is_empty(), + "expected at least one post-handshake request to inspect" + ); + assert!( + sessions.iter().all(Option::is_none), + "client echoed Mcp-Session-Id at a version with no sessions: {sessions:?}" + ); +} + +#[tokio::test] +async fn legacy_version_keeps_the_session_and_opens_the_stream() { + let recorder = connect_and_record("2025-11-25").await; + + // The legacy shape is untouched: the stream is opened and carries the session id. + let gets = recorder.get_requests(); + assert_eq!( + gets.len(), + 1, + "expected exactly one standalone GET stream on a legacy session, got {gets:?}" + ); + assert_eq!(gets[0].2.as_deref(), Some(SESSION_ID)); + let sessions = recorder.session_headers_after_handshake(); + assert!( + sessions.iter().any(|s| s.as_deref() == Some(SESSION_ID)), + "legacy session id was not echoed after the handshake: {sessions:?}" + ); +}