Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 16 additions & 14 deletions lib/mcp/server/transports/streamable_http_transport.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
50 changes: 50 additions & 0 deletions test/mcp/server/transports/streamable_http_transport_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Loading