From 2200cbd2bbf39499e8dec64703b4f2d0a31cbbbb Mon Sep 17 00:00:00 2001 From: Koichi ITO Date: Wed, 26 Aug 2026 00:24:10 +0900 Subject: [PATCH] Serialize `subscriptions/listen` writes so the graceful result is the final message ## Motivation and Context On graceful teardown (`transport.close`), each open `subscriptions/listen` stream receives its `SubscriptionsListenResult` response before closing, signaling a clean end the client can distinguish from an abrupt disconnect. But a delivery that snapshotted the stream from the registry before teardown cleared it performs its write outside the lock, so a change notification could land after the final result. Each subscription entry now carries a per-stream write mutex. Teardown marks the entry closed and writes the result under that mutex; a racing delivery either lands before the result or observes `closed` and skips, keeping the graceful result the stream's final message. The mutex is per stream and held only around the single write call, so deliveries to different streams stay as parallel as before, and no path nests it inside the transport mutex (delivery snapshots release `@mutex` first, and the write-error cleanup takes `@mutex` only after the write mutex is released), ruling out lock-order inversions. Keepalive pings stay outside the mutex: they are SSE comment frames, which clients ignore by specification, so one landing after the result is harmless. ## How Has This Been Tested? With a new regression test pinning the race (an entry captured before teardown and re-presented to delivery afterwards writes nothing after the result), the full suite, RuboCop, and the conformance suite, all green. ## Breaking Changes None. The write mutex is internal to the listen registry; the wire only gains the ordering guarantee. --- .../transports/streamable_http_transport.rb | 56 ++++++++++++------- .../streamable_http_transport_test.rb | 21 ++++++- 2 files changed, 55 insertions(+), 22 deletions(-) diff --git a/lib/mcp/server/transports/streamable_http_transport.rb b/lib/mcp/server/transports/streamable_http_transport.rb index 8f329b5c..b406d6df 100644 --- a/lib/mcp/server/transports/streamable_http_transport.rb +++ b/lib/mcp/server/transports/streamable_http_transport.rb @@ -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 = {} @@ -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 @@ -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 @@ -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 diff --git a/test/mcp/server/transports/streamable_http_transport_test.rb b/test/mcp/server/transports/streamable_http_transport_test.rb index a7d192bf..36fe9ac2 100644 --- a/test/mcp/server/transports/streamable_http_transport_test.rb +++ b/test/mcp/server/transports/streamable_http_transport_test.rb @@ -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 @@ -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 }) @@ -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")