Skip to content

Commit 74451d3

Browse files
authored
Merge pull request #574 from koic/keep_the_listen_acknowledgement_first_on_transport_close
Keep the acknowledgement first when the transport closes during a listen registration
2 parents c9641fb + 77d4e4b commit 74451d3

2 files changed

Lines changed: 112 additions & 22 deletions

File tree

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

Lines changed: 38 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -929,29 +929,33 @@ def first
929929
# the legacy GET stream (`create_sse_body`).
930930
#
931931
# Registration and activation are split on purpose: the entry is inserted inactive
932-
# (reserving the cap slot atomically), the acknowledgement is written outside the lock,
932+
# (reserving the cap slot atomically), the acknowledgement is written outside the registry lock,
933933
# and only then does the entry become eligible for delivery. A concurrent notification between
934934
# the insert and the acknowledgement write skips the inactive entry,
935935
# enforcing the SEP-2575 rule that no notification precedes the acknowledgement.
936936
#
937+
# The acknowledgement write holds the stream's write mutex, which `teardown_listen_subscriptions` also takes
938+
# before it marks an entry closed. The two therefore cannot interleave: the acknowledgement either lands
939+
# before the result, or, if the transport closed first, is not written at all and the stream just closes,
940+
# so no stream ever carries a result ahead of its acknowledgement.
941+
#
937942
# The entry is keyed by an identifier minted here, not by the request id: that id is unique only among
938943
# the requesting client's own in-flight requests, and two clients that pick the same one must each get
939944
# their stream, stamped with the id they sent.
940945
def listen_sse_body(request_id, honored)
941946
ListenStreamBody.new do |stream|
942947
subscription_key = SecureRandom.uuid
943-
rejected = false
948+
subscription = nil
944949
@mutex.synchronize do
945-
if @max_listen_subscriptions && @listen_subscriptions.size >= @max_listen_subscriptions
946-
rejected = true
947-
else
948-
@listen_subscriptions[subscription_key] = {
950+
unless @max_listen_subscriptions && @listen_subscriptions.size >= @max_listen_subscriptions
951+
subscription = {
949952
request_id: request_id, stream: stream, filter: honored, active: false, write_mutex: Mutex.new, keepalive_wakeup: ConditionVariable.new
950953
}
954+
@listen_subscriptions[subscription_key] = subscription
951955
end
952956
end
953957

954-
if rejected
958+
if subscription.nil?
955959
close_stream_safely(stream)
956960
else
957961
acknowledgement = {
@@ -964,9 +968,27 @@ def listen_sse_body(request_id, honored)
964968
}
965969

966970
begin
967-
send_to_stream(stream, acknowledgement)
968-
activate_listen_subscription(subscription_key)
969-
start_listen_keepalive_thread(subscription_key, request_id)
971+
acknowledged = subscription[:write_mutex].synchronize do
972+
next false if subscription[:closed]
973+
974+
send_to_stream(stream, acknowledgement)
975+
976+
# Set on the entry itself, not through the registry: a concurrent close may already have cleared the registry
977+
# while its result write waits on this mutex, and that write must still find the stream acknowledged.
978+
# Set under the registry lock as well, since that is the lock the delivery snapshot reads the flag under.
979+
# This is the one place a write mutex is held while the registry lock is taken; it stays deadlock-free only
980+
# as long as no path takes a write mutex inside `@mutex.synchronize`, so resolve entries under `@mutex`,
981+
# release it, then write.
982+
@mutex.synchronize { subscription[:active] = true }
983+
984+
true
985+
end
986+
987+
if acknowledged
988+
start_listen_keepalive_thread(subscription_key, request_id)
989+
else
990+
close_stream_safely(stream)
991+
end
970992
rescue *STREAM_WRITE_ERRORS
971993
remove_listen_subscription(subscription_key)
972994
close_stream_safely(stream)
@@ -975,15 +997,6 @@ def listen_sse_body(request_id, honored)
975997
end
976998
end
977999

978-
# Marks a listen subscription eligible for delivery once its acknowledgement write has completed.
979-
# The entry may already be gone when the transport closed concurrently.
980-
def activate_listen_subscription(subscription_key)
981-
@mutex.synchronize do
982-
subscription = @listen_subscriptions[subscription_key]
983-
subscription[:active] = true if subscription
984-
end
985-
end
986-
9871000
# Periodically writes an SSE keepalive comment frame to a listen stream so a silently dropped
9881001
# connection is detected and its slot freed, rather than held until the next fan-out write.
9891002
# Mirrors the legacy GET stream's `start_keepalive_thread`; a comment frame (not a data frame)
@@ -1149,11 +1162,15 @@ def teardown_listen_subscriptions
11491162

11501163
removed.each_value do |subscription|
11511164
# Marking the entry closed and writing the result under the stream's write mutex orders
1152-
# this against in-flight deliveries: each one either lands before the result or observes
1153-
# `closed` and skips, keeping the graceful result the stream's final message.
1165+
# this against in-flight deliveries and against the acknowledgement write: each one either lands
1166+
# before the result or observes `closed` and skips, keeping the graceful result the stream's final message.
11541167
subscription[:write_mutex].synchronize do
11551168
subscription[:closed] = true
11561169

1170+
# A stream whose acknowledgement was never written gets no result either: SEP-2575 makes
1171+
# the acknowledgement the first message, so the stream closes abruptly and the client re-sends.
1172+
next unless subscription[:active]
1173+
11571174
begin
11581175
send_to_stream(subscription[:stream], {
11591176
jsonrpc: "2.0",

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

Lines changed: 74 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6125,7 +6125,7 @@ def string
61256125

61266126
assert_empty sse_events(io)
61276127

6128-
@transport.send(:activate_listen_subscription, "listen-1")
6128+
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"][:active] = true
61296129
@server.notify_tools_list_changed
61306130

61316131
assert_equal ["notifications/tools/list_changed"], sse_events(io).map { |event| event["method"] }
@@ -6149,6 +6149,79 @@ def string
61496149
refute(events.any? { |event| event["method"] == "notifications/tools/list_changed" })
61506150
end
61516151

6152+
test "transport close during the acknowledgement write still sends the acknowledgement first" do
6153+
# The stream blocks its first write, the acknowledgement, until released, so the close arrives while
6154+
# that write is in progress and has to queue behind it on the stream's write mutex.
6155+
reached = Queue.new
6156+
release = Queue.new
6157+
io = StringIO.new
6158+
first_write = true
6159+
6160+
io.define_singleton_method(:write) do |data|
6161+
if first_write
6162+
first_write = false
6163+
reached.push(true)
6164+
release.pop
6165+
end
6166+
super(data)
6167+
end
6168+
6169+
response = @transport.handle_request(modern_rack_request(
6170+
modern_listen_body(id: "listen-1", params: { notifications: { toolsListChanged: true } }),
6171+
))
6172+
body_thread = Thread.new { response[2].call(io) }
6173+
close_thread = nil
6174+
6175+
begin
6176+
wait_until { !reached.empty? }
6177+
close_thread = Thread.new { @transport.close }
6178+
wait_until_blocked_or_done(close_thread)
6179+
ensure
6180+
# Released whatever the waits above did, so a failed wait cannot leave the body holding
6181+
# the stream's write mutex with the close queued behind it.
6182+
release.push(true)
6183+
end
6184+
6185+
assert(body_thread.join(5), "the listen body did not finish")
6186+
assert(close_thread.join(5), "the transport close did not finish")
6187+
6188+
events = sse_events(io)
6189+
6190+
assert_equal 2, events.size
6191+
assert_equal "notifications/subscriptions/acknowledged", events[0]["method"]
6192+
assert_equal "complete", events[1].dig("result", "resultType")
6193+
assert_predicate io, :closed?
6194+
end
6195+
6196+
test "transport close before the acknowledgement closes the stream without a result" do
6197+
# An entry in the registered-but-not-yet-acknowledged state: with no acknowledgement written,
6198+
# a result would be the stream's first message, which SEP-2575 forbids.
6199+
io = StringIO.new
6200+
@transport.instance_variable_get(:@listen_subscriptions)["listen-1"] = {
6201+
request_id: "listen-1",
6202+
stream: io,
6203+
filter: { toolsListChanged: true },
6204+
active: false,
6205+
write_mutex: Mutex.new,
6206+
keepalive_wakeup: ConditionVariable.new,
6207+
}
6208+
6209+
@transport.close
6210+
6211+
assert_empty sse_events(io)
6212+
assert_predicate io, :closed?
6213+
end
6214+
6215+
# Waits until `thread` is either blocked (asleep on a lock or queue) or finished, so the step after it
6216+
# runs against a thread that has made its move.
6217+
def wait_until_blocked_or_done(thread)
6218+
deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + 5
6219+
until thread.status == "sleep" || thread.status == false
6220+
flunk("thread did not block or finish") if Process.clock_gettime(Process::CLOCK_MONOTONIC) > deadline
6221+
sleep(0.005)
6222+
end
6223+
end
6224+
61526225
test "subscriptions/listen streams for different subscriptions receive their own subscriptionId" do
61536226
first = open_listen_stream(id: "listen-1", notifications: { toolsListChanged: true })
61546227
second = open_listen_stream(id: "listen-2", notifications: { toolsListChanged: true })

0 commit comments

Comments
 (0)