Skip to content

Commit 6da3009

Browse files
authored
Merge pull request #567 from koic/key_listen_streams_by_a_server_minted_identifier
Serve `subscriptions/listen` streams that carry the same request id
2 parents 08f1acf + 0c89b43 commit 6da3009

2 files changed

Lines changed: 93 additions & 42 deletions

File tree

‎lib/mcp/server/transports/streamable_http_transport.rb‎

Lines changed: 40 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -169,8 +169,10 @@ def initialize(
169169
@allowed_origins = Array(allowed_origins).map(&:downcase).freeze
170170
@pending_responses = {}
171171

172-
# Maps a `subscriptions/listen` request id to
173-
# `{ stream: stream_object, filter: honored_subscription_filter, active: boolean, write_mutex: Mutex }` (SEP-2575).
172+
# Maps a key the transport mints for each `subscriptions/listen` stream to
173+
# `{ request_id: listen_request_id, stream: stream_object, filter: honored_subscription_filter, active: boolean,
174+
# write_mutex: Mutex }` (SEP-2575). The request id is the client's, unique only among that client's own
175+
# in-flight requests, so it stamps `subscriptionId` but cannot serve as the key: two clients may pick the same one.
174176
# In-process only; a multi-worker deployment needs an external event bus to fan notifications out across processes,
175177
# which is a follow-up.
176178
@listen_subscriptions = {}
@@ -925,19 +927,25 @@ def first
925927
# the legacy GET stream (`create_sse_body`).
926928
#
927929
# Registration and activation are split on purpose: the entry is inserted inactive
928-
# (reserving the id and the cap slot atomically), the acknowledgement is written outside the lock,
930+
# (reserving the cap slot atomically), the acknowledgement is written outside the lock,
929931
# and only then does the entry become eligible for delivery. A concurrent notification between
930932
# the insert and the acknowledgement write skips the inactive entry,
931933
# enforcing the SEP-2575 rule that no notification precedes the acknowledgement.
934+
#
935+
# The entry is keyed by an identifier minted here, not by the request id: that id is unique only among
936+
# the requesting client's own in-flight requests, and two clients that pick the same one must each get
937+
# their stream, stamped with the id they sent.
932938
def listen_sse_body(request_id, honored)
933939
ListenStreamBody.new do |stream|
940+
subscription_key = SecureRandom.uuid
934941
rejected = false
935942
@mutex.synchronize do
936-
if @listen_subscriptions.key?(request_id) ||
937-
(@max_listen_subscriptions && @listen_subscriptions.size >= @max_listen_subscriptions)
943+
if @max_listen_subscriptions && @listen_subscriptions.size >= @max_listen_subscriptions
938944
rejected = true
939945
else
940-
@listen_subscriptions[request_id] = { stream: stream, filter: honored, active: false, write_mutex: Mutex.new }
946+
@listen_subscriptions[subscription_key] = {
947+
request_id: request_id, stream: stream, filter: honored, active: false, write_mutex: Mutex.new,
948+
}
941949
end
942950
end
943951

@@ -955,10 +963,10 @@ def listen_sse_body(request_id, honored)
955963

956964
begin
957965
send_to_stream(stream, acknowledgement)
958-
activate_listen_subscription(request_id)
959-
start_listen_keepalive_thread(request_id)
966+
activate_listen_subscription(subscription_key)
967+
start_listen_keepalive_thread(subscription_key, request_id)
960968
rescue *STREAM_WRITE_ERRORS
961-
remove_listen_subscription(request_id)
969+
remove_listen_subscription(subscription_key)
962970
close_stream_safely(stream)
963971
end
964972
end
@@ -967,9 +975,9 @@ def listen_sse_body(request_id, honored)
967975

968976
# Marks a listen subscription eligible for delivery once its acknowledgement write has completed.
969977
# The entry may already be gone when the transport closed concurrently.
970-
def activate_listen_subscription(request_id)
978+
def activate_listen_subscription(subscription_key)
971979
@mutex.synchronize do
972-
subscription = @listen_subscriptions[request_id]
980+
subscription = @listen_subscriptions[subscription_key]
973981
subscription[:active] = true if subscription
974982
end
975983
end
@@ -978,37 +986,39 @@ def activate_listen_subscription(request_id)
978986
# connection is detected and its slot freed, rather than held until the next fan-out write.
979987
# Mirrors the legacy GET stream's `start_keepalive_thread`; a comment frame (not a data frame)
980988
# cannot corrupt an interleaved notification's JSON.
981-
def start_listen_keepalive_thread(request_id)
989+
def start_listen_keepalive_thread(subscription_key, request_id)
982990
return unless @listen_keepalive_interval
983991

984992
Thread.new do
985-
while listen_subscription_active?(request_id)
993+
while listen_subscription_active?(subscription_key)
986994
sleep(@listen_keepalive_interval)
987-
send_listen_keepalive_ping(request_id)
995+
send_listen_keepalive_ping(subscription_key)
988996
end
989997
rescue *STREAM_WRITE_ERRORS
990998
# The peer went away; the ensure frees the slot. A dropped listen stream is the normal
991999
# way this loop ends, so it is not reported.
9921000
rescue StandardError => e
1001+
# The request id is taken from the caller rather than the registry: a delivery failure may have
1002+
# removed the entry already, and the report should still name the stream.
9931003
MCP.configuration.exception_reporter.call(e, { subscription_id: request_id })
9941004
ensure
9951005
stream = @mutex.synchronize do
996-
subscription = @listen_subscriptions.delete(request_id)
1006+
subscription = @listen_subscriptions.delete(subscription_key)
9971007
subscription && subscription[:stream]
9981008
end
9991009
close_stream_safely(stream) if stream
10001010
end
10011011
end
10021012

1003-
def listen_subscription_active?(request_id)
1004-
@mutex.synchronize { @listen_subscriptions.key?(request_id) }
1013+
def listen_subscription_active?(subscription_key)
1014+
@mutex.synchronize { @listen_subscriptions.key?(subscription_key) }
10051015
end
10061016

10071017
# Resolves the stream under the lock, then writes outside it so a stalled reader cannot block
10081018
# every other subscription on `@mutex`. A write error propagates to end the keepalive loop.
1009-
def send_listen_keepalive_ping(request_id)
1019+
def send_listen_keepalive_ping(subscription_key)
10101020
stream = @mutex.synchronize do
1011-
subscription = @listen_subscriptions[request_id]
1021+
subscription = @listen_subscriptions[subscription_key]
10121022
subscription && subscription[:stream]
10131023
end
10141024
return unless stream
@@ -1054,7 +1064,7 @@ def deliver_to_listen_subscriptions(method, params)
10541064
# The matching snapshot is taken under `@mutex`, but stream writes happen outside it:
10551065
# a slow or stalled subscriber must not block the transport, matching the legacy delivery paths.
10561066
matched = @mutex.synchronize do
1057-
@listen_subscriptions.filter_map do |request_id, subscription|
1067+
@listen_subscriptions.filter_map do |subscription_key, subscription|
10581068
# An inactive entry has not finished writing its acknowledgement yet;
10591069
# delivering to it would put a notification ahead of the acknowledgement.
10601070
next unless subscription[:active]
@@ -1067,12 +1077,12 @@ def deliver_to_listen_subscriptions(method, params)
10671077
uris.is_a?(Array) && uris.include?(uri)
10681078
end
10691079

1070-
[request_id, subscription] if hit
1080+
[subscription_key, subscription] if hit
10711081
end
10721082
end
10731083

1074-
matched.each do |request_id, subscription|
1075-
meta = { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => request_id }
1084+
matched.each do |subscription_key, subscription|
1085+
meta = { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => subscription[:request_id] }
10761086
notification_params = (params || {}).merge(_meta: meta)
10771087
notification = { jsonrpc: "2.0", method: method, params: notification_params }
10781088

@@ -1089,16 +1099,16 @@ def deliver_to_listen_subscriptions(method, params)
10891099
rescue *STREAM_WRITE_ERRORS => e
10901100
MCP.configuration.exception_reporter.call(
10911101
e,
1092-
{ subscription_id: request_id, error: "Failed to send notification" },
1102+
{ subscription_id: subscription[:request_id], error: "Failed to send notification" },
10931103
)
1094-
remove_listen_subscription(request_id)
1104+
remove_listen_subscription(subscription_key)
10951105
close_stream_safely(subscription[:stream])
10961106
end
10971107
end
10981108
end
10991109

1100-
def remove_listen_subscription(request_id)
1101-
@mutex.synchronize { @listen_subscriptions.delete(request_id) }
1110+
def remove_listen_subscription(subscription_key)
1111+
@mutex.synchronize { @listen_subscriptions.delete(subscription_key) }
11021112
end
11031113

11041114
# Graceful teardown (SEP-2575): each open listen stream receives its `SubscriptionsListenResult` response
@@ -1110,7 +1120,7 @@ def teardown_listen_subscriptions
11101120
subscriptions
11111121
end
11121122

1113-
removed.each do |request_id, subscription|
1123+
removed.each_value do |subscription|
11141124
# Marking the entry closed and writing the result under the stream's write mutex orders
11151125
# this against in-flight deliveries: each one either lands before the result or observes
11161126
# `closed` and skips, keeping the graceful result the stream's final message.
@@ -1120,13 +1130,13 @@ def teardown_listen_subscriptions
11201130
begin
11211131
send_to_stream(subscription[:stream], {
11221132
jsonrpc: "2.0",
1123-
id: request_id,
1133+
id: subscription[:request_id],
11241134
result: {
11251135
# `SubscriptionsListenResult` is served at the transport layer and never
11261136
# passes through the dispatch path, so the REQUIRED 2026-07-28 `resultType` is
11271137
# stamped at its construction site.
11281138
resultType: ResultType::COMPLETE,
1129-
_meta: { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => request_id },
1139+
_meta: { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => subscription[:request_id] },
11301140
},
11311141
})
11321142
rescue *STREAM_WRITE_ERRORS

‎test/mcp/server/transports/streamable_http_transport_test.rb‎

Lines changed: 53 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -6113,7 +6113,7 @@ def string
61136113
# the registry insert and the acknowledgement write, which happens outside the lock.
61146114
io = StringIO.new
61156115
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = {
6116-
stream: io, filter: { toolsListChanged: true }, active: false, write_mutex: Mutex.new
6116+
request_id: "listen-1", stream: io, filter: { toolsListChanged: true }, active: false, write_mutex: Mutex.new
61176117
}
61186118

61196119
@server.notify_tools_list_changed
@@ -6128,13 +6128,14 @@ def string
61286128

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

61336134
@transport.close
61346135

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

61406141
events = sse_events(io)
@@ -6257,16 +6258,54 @@ def string
62576258
transport.close
62586259
end
62596260

6260-
test "subscriptions/listen rejects a duplicate subscription id by closing the new stream" do
6261-
open_listen_stream(id: "listen-1", notifications: { toolsListChanged: true })
6261+
test "subscriptions/listen serves two streams that carry the same request id" do
6262+
# A request id is unique only among one client's own in-flight requests, and clients that number
6263+
# their requests with a counter reach the same small integers, so two clients' listen requests may
6264+
# carry the same id. Each gets its own stream, stamped with the id it sent.
6265+
first = open_listen_stream(id: 1, notifications: { toolsListChanged: true })
6266+
second = open_listen_stream(id: 1, notifications: { toolsListChanged: true })
62626267

6263-
duplicate = StringIO.new
6264-
response = @transport.handle_request(modern_rack_request(
6265-
modern_listen_body(id: "listen-1", params: { notifications: { toolsListChanged: true } }),
6266-
))
6267-
response[2].call(duplicate)
6268+
@server.notify_tools_list_changed
6269+
6270+
[first, second].each do |io|
6271+
refute_predicate io, :closed?
6272+
events = sse_events(io)
6273+
assert_equal ["notifications/subscriptions/acknowledged", "notifications/tools/list_changed"], events.map { |event| event["method"] }
6274+
assert_equal [1, 1], events.map { |event| event.dig("params", "_meta", "io.modelcontextprotocol/subscriptionId") }
6275+
end
6276+
end
6277+
6278+
test "transport close sends each of two streams sharing a request id its own result" do
6279+
first = open_listen_stream(id: 1, notifications: { toolsListChanged: true })
6280+
second = open_listen_stream(id: 1, notifications: { toolsListChanged: true })
6281+
6282+
@transport.close
62686283

6269-
assert_predicate duplicate, :closed?
6284+
[first, second].each do |io|
6285+
result = sse_events(io).last
6286+
assert_equal 1, result["id"]
6287+
assert_equal "complete", result.dig("result", "resultType")
6288+
assert_predicate io, :closed?
6289+
end
6290+
end
6291+
6292+
test "a failed write on one of two streams sharing a request id frees only that stream" do
6293+
first = open_listen_stream(id: 1, notifications: { toolsListChanged: true })
6294+
second = open_listen_stream(id: 1, notifications: { toolsListChanged: true })
6295+
first.define_singleton_method(:write) { |_data| raise Errno::EPIPE }
6296+
6297+
@server.notify_tools_list_changed
6298+
6299+
assert_predicate first, :closed?
6300+
refute_predicate second, :closed?
6301+
registry = @transport.instance_variable_get(:@listen_subscriptions)
6302+
assert_equal 1, registry.size
6303+
assert_same second, registry.values.first[:stream]
6304+
6305+
# The entry that survived is the second stream's, so a further notification still reaches it.
6306+
@server.notify_tools_list_changed
6307+
6308+
assert_equal 2, sse_events(second).count { |event| event["method"] == "notifications/tools/list_changed" }
62706309
end
62716310

62726311
test "listen keepalive writes a comment frame outside the mutex" do
@@ -6281,7 +6320,9 @@ def string
62816320
ping = data
62826321
end
62836322
stream.define_singleton_method(:flush) {}
6284-
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = { stream: stream, filter: {}, write_mutex: Mutex.new }
6323+
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = {
6324+
request_id: "listen-1", stream: stream, filter: {}, write_mutex: Mutex.new,
6325+
}
62856326

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

0 commit comments

Comments
 (0)