From 6601fd48ee5fbd8dd9ac64a550546bfc04f02ec0 Mon Sep 17 00:00:00 2001 From: Koichi ITO Date: Sat, 10 Oct 2026 00:40:56 +0900 Subject: [PATCH] Close a refused SSE stream outside the registry lock ## Motivation and Context `StreamableHTTPTransport` closes every stream outside its registry lock, so a close that blocks on the peer cannot hold every other session on `@mutex`. One path broke that rule: when a `GET` stream could not be registered, because the session had gone or another stream had taken its place, `store_stream_for_session` closed the refused stream inside the lock, and an error raised by that close escaped the Rack body. The refused stream is now closed after the lock is released and through `close_stream_safely`, like the streams the transport drops everywhere else. ## How Has This Been Tested? Three new tests in `test/mcp/server/transports/streamable_http_transport_test.rb`: two refuse a second stream for a session, one checking from inside the stream's `close` that the lock is free, the other that a `close` raising `IOError` is swallowed and the session keeps serving, and the third refuses a stream for a session that is gone, checking the lock the same way. Against the previous library, the first and the third find the lock held and the second raises. ## Breaking Changes None. --- .../transports/streamable_http_transport.rb | 30 +++++------ .../streamable_http_transport_test.rb | 50 +++++++++++++++++++ 2 files changed, 66 insertions(+), 14 deletions(-) diff --git a/lib/mcp/server/transports/streamable_http_transport.rb b/lib/mcp/server/transports/streamable_http_transport.rb index 7d02841f..1ccadaf4 100644 --- a/lib/mcp/server/transports/streamable_http_transport.rb +++ b/lib/mcp/server/transports/streamable_http_transport.rb @@ -2444,22 +2444,24 @@ def create_sse_body(session_id, auth_info = nil) # The expiry of the token that authenticated the GET is kept with the stream so the keepalive loop and # the delivery paths can close the stream once that token expires. def store_stream_for_session(session_id, stream, auth_info = nil) - @mutex.synchronize do + stored = @mutex.synchronize do session = @sessions[session_id] - if session && !session[:get_sse_stream] - session[:get_sse_stream] = stream - session[:get_sse_stream_expires_at] = stream_token_expiry(auth_info) - # The expiry is nil for a stream opened without a token, and callers read a falsy return as "not stored", - # so the stream itself is the return value. - stream - else - # Either session was removed, or another request already established a stream. - stream.close - # `stream.close` may return a truthy value depending on the stream class. - # Explicitly return nil to guarantee a falsy return for callers. - nil - end + next false unless session && !session[:get_sse_stream] + + session[:get_sse_stream] = stream + session[:get_sse_stream_expires_at] = stream_token_expiry(auth_info) + true end + + # The expiry is nil for a stream opened without a token, and callers read a falsy return as "not stored", + # so the stream itself is the return value. + return stream if stored + + # Either the session was removed, or another request already established a stream. The refused stream is + # closed outside the lock like every other stream this transport closes, so a close that blocks on the peer + # cannot hold every other session on `@mutex`. + close_stream_safely(stream) + nil end # The thread acts on the stream it was started for and on no other: the session may have detached that diff --git a/test/mcp/server/transports/streamable_http_transport_test.rb b/test/mcp/server/transports/streamable_http_transport_test.rb index 30be27fe..06fe41ab 100644 --- a/test/mcp/server/transports/streamable_http_transport_test.rb +++ b/test/mcp/server/transports/streamable_http_transport_test.rb @@ -743,6 +743,56 @@ def string assert stream_b.closed? end + test "store_stream_for_session closes a refused stream outside the mutex" do + session_id = initialize_test_session + @transport.send(:store_stream_for_session, session_id, StringIO.new) + + mutex = @transport.instance_variable_get(:@mutex) + closed_outside_mutex = false + refused = Object.new + refused.define_singleton_method(:close) do + if mutex.try_lock + closed_outside_mutex = true + mutex.unlock + end + end + + assert_nil @transport.send(:store_stream_for_session, session_id, refused) + assert closed_outside_mutex, "the refused stream was closed while the mutex was held" + end + + test "store_stream_for_session swallows an error from closing a refused stream" do + session_id = initialize_test_session + @transport.send(:store_stream_for_session, session_id, StringIO.new) + + failing = Object.new + failing.define_singleton_method(:close) { raise IOError, "already closed by the peer" } + + assert_nil @transport.send(:store_stream_for_session, session_id, failing) + assert_equal 200, @transport.handle_request(create_rack_request( + "POST", + "/", + { "CONTENT_TYPE" => "application/json", "HTTP_MCP_SESSION_ID" => session_id }, + { jsonrpc: "2.0", method: "ping", id: "ping-1" }.to_json, + ))[0] + end + + test "store_stream_for_session closes the stream outside the mutex when the session is gone" do + # The session was removed between the GET's checks and its body running, so there is nothing to attach to. + mutex = @transport.instance_variable_get(:@mutex) + closed_outside_mutex = false + orphan = Object.new + orphan.define_singleton_method(:close) do + if mutex.try_lock + closed_outside_mutex = true + mutex.unlock + end + end + + assert_nil @transport.send(:store_stream_for_session, "no-such-session", orphan) + assert closed_outside_mutex, "the orphaned stream was closed while the mutex was held" + end + test "handles GET request with invalid session ID" do request = create_rack_request( "GET",