Skip to content
Merged
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
56 changes: 36 additions & 20 deletions lib/mcp/server/transports/streamable_http_transport.rb
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,8 @@ def initialize(
@allowed_origins = Array(allowed_origins).map(&:downcase).freeze
@pending_responses = {}

# Maps a `subscriptions/listen` request id to `{ stream: stream_object, filter: honored_subscription_filter, active: boolean }` (SEP-2575).
# Maps a `subscriptions/listen` request id to
# `{ stream: stream_object, filter: honored_subscription_filter, active: boolean, write_mutex: Mutex }` (SEP-2575).
# In-process only; a multi-worker deployment needs an external event bus to fan notifications out across processes,
# which is a follow-up.
@listen_subscriptions = {}
Expand Down Expand Up @@ -936,7 +937,7 @@ def listen_sse_body(request_id, honored)
(@max_listen_subscriptions && @listen_subscriptions.size >= @max_listen_subscriptions)
rejected = true
else
@listen_subscriptions[request_id] = { stream: stream, filter: honored, active: false }
@listen_subscriptions[request_id] = { stream: stream, filter: honored, active: false, write_mutex: Mutex.new }
end
end

Expand Down Expand Up @@ -1066,24 +1067,32 @@ def deliver_to_listen_subscriptions(method, params)
uris.is_a?(Array) && uris.include?(uri)
end

[request_id, subscription[:stream]] if hit
[request_id, subscription] if hit
end
end

matched.each do |request_id, stream|
matched.each do |request_id, subscription|
meta = { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => request_id }
notification_params = (params || {}).merge(_meta: meta)
notification = { jsonrpc: "2.0", method: method, params: notification_params }

begin
send_to_stream(stream, notification)
# The per-stream write mutex orders this write against a concurrent graceful teardown:
# once teardown has marked the entry closed and written its `SubscriptionsListenResult`,
# a delivery that snapshotted the entry before the registry was cleared skips it instead
# of writing after the final message.
subscription[:write_mutex].synchronize do
next if subscription[:closed]

send_to_stream(subscription[:stream], notification)
end
rescue *STREAM_WRITE_ERRORS => e
MCP.configuration.exception_reporter.call(
e,
{ subscription_id: request_id, error: "Failed to send notification" },
)
remove_listen_subscription(request_id)
close_stream_safely(stream)
close_stream_safely(subscription[:stream])
end
end
end
Expand All @@ -1102,20 +1111,27 @@ def teardown_listen_subscriptions
end

removed.each do |request_id, subscription|
begin
send_to_stream(subscription[:stream], {
jsonrpc: "2.0",
id: request_id,
result: {
# `SubscriptionsListenResult` is served at the transport layer and never
# passes through the dispatch path, so the REQUIRED 2026-07-28 `resultType` is
# stamped at its construction site.
resultType: ResultType::COMPLETE,
_meta: { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => request_id },
},
})
rescue *STREAM_WRITE_ERRORS
nil
# Marking the entry closed and writing the result under the stream's write mutex orders
# this against in-flight deliveries: each one either lands before the result or observes
# `closed` and skips, keeping the graceful result the stream's final message.
subscription[:write_mutex].synchronize do
subscription[:closed] = true

begin
send_to_stream(subscription[:stream], {
jsonrpc: "2.0",
id: request_id,
result: {
# `SubscriptionsListenResult` is served at the transport layer and never
# passes through the dispatch path, so the REQUIRED 2026-07-28 `resultType` is
# stamped at its construction site.
resultType: ResultType::COMPLETE,
_meta: { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => request_id },
},
})
rescue *STREAM_WRITE_ERRORS
nil
end
end
close_stream_safely(subscription[:stream])
end
Expand Down
21 changes: 19 additions & 2 deletions test/mcp/server/transports/streamable_http_transport_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -6113,7 +6113,7 @@ def string
# the registry insert and the acknowledgement write, which happens outside the lock.
io = StringIO.new
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = {
stream: io, filter: { toolsListChanged: true }, active: false
stream: io, filter: { toolsListChanged: true }, active: false, write_mutex: Mutex.new
}

@server.notify_tools_list_changed
Expand All @@ -6126,6 +6126,23 @@ def string
assert_equal ["notifications/tools/list_changed"], sse_events(io).map { |event| event["method"] }
end

test "a delivery racing the graceful teardown cannot write after the final result" do
io = open_listen_stream(id: "listen-1", notifications: { toolsListChanged: true })
entry = @transport.instance_variable_get(:@listen_subscriptions)["listen-1"]

@transport.close

# Simulate an in-flight delivery that snapshotted the entry before teardown cleared
# the registry: the closed flag set under the write mutex makes it a no-op.
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = entry
@server.notify_tools_list_changed

events = sse_events(io)

assert_equal "complete", events.last.dig("result", "resultType")
refute(events.any? { |event| event["method"] == "notifications/tools/list_changed" })
end

test "subscriptions/listen streams for different subscriptions receive their own subscriptionId" do
first = open_listen_stream(id: "listen-1", notifications: { toolsListChanged: true })
second = open_listen_stream(id: "listen-2", notifications: { toolsListChanged: true })
Expand Down Expand Up @@ -6264,7 +6281,7 @@ def string
ping = data
end
stream.define_singleton_method(:flush) {}
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = { stream: stream, filter: {} }
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = { stream: stream, filter: {}, write_mutex: Mutex.new }

@transport.send(:send_listen_keepalive_ping, "listen-1")

Expand Down