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
57 changes: 46 additions & 11 deletions lib/mcp/server/transports/streamable_http_transport.rb
Original file line number Diff line number Diff line change
Expand Up @@ -904,6 +904,14 @@ def handle_subscriptions_listen(body)
return invalid_request_response("Invalid Request: subscriptions/listen requires an id")
end

# The id is echoed in every message of the stream, so one that cannot be written back as JSON
# (a String holding bytes that are not valid UTF-8 parses, but does not generate) is refused here,
# while the refusal can still answer with a null id, rather than surfacing as a failed write after
# the stream has been registered.
unless json_encodable?(request_id)
return invalid_request_response("Invalid Request: subscriptions/listen id must be valid UTF-8")
end

begin
if RequestEnvelope.modern?(params)
RequestEnvelope.parse!(params, request: params)
Expand Down Expand Up @@ -946,7 +954,43 @@ def handle_subscriptions_listen(body)
return too_many_listen_subscriptions_response(request_id)
end

[200, SSE_HEADERS.dup, listen_sse_body(request_id, honored_filter(filter))]
honored = honored_filter(filter)

# The acknowledgement echoes the honored filter, whose resource URIs are client-supplied strings,
# so it is encoded before the stream is committed to: a URI that cannot be written as JSON is
# refused with a 400 instead of failing the first write after the stream was registered.
unless (acknowledgement = listen_acknowledgement_json(request_id, honored))
return json_rpc_error_response(
status: 400,
code: JsonRpcHandler::ErrorCode::INVALID_PARAMS,
message: "Invalid params: subscriptions/listen `notifications` must be valid UTF-8",
id: request_id,
)
end

[200, SSE_HEADERS.dup, listen_sse_body(request_id, honored, acknowledgement)]
end

# The `notifications/subscriptions/acknowledged` frame for a listen stream as JSON, or `nil` when
# the request id or the honored filter holds a string JSON cannot encode.
def listen_acknowledgement_json(request_id, honored)
{
jsonrpc: "2.0",
method: Methods::NOTIFICATIONS_SUBSCRIPTIONS_ACKNOWLEDGED,
params: {
notifications: honored,
_meta: { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => request_id },
},
}.to_json
rescue JSON::GeneratorError
nil
end

def json_encodable?(value)
value.to_json
true
rescue JSON::GeneratorError
false
end

def listen_subscriptions_full?
Expand Down Expand Up @@ -1005,7 +1049,7 @@ def first
# The entry is keyed by an identifier minted here, not by the request id: that id is unique only among
# the requesting client's own in-flight requests, and two clients that pick the same one must each get
# their stream, stamped with the id they sent.
def listen_sse_body(request_id, honored)
def listen_sse_body(request_id, honored, acknowledgement)
ListenStreamBody.new do |stream|
subscription_key = SecureRandom.uuid
subscription = nil
Expand All @@ -1029,15 +1073,6 @@ def listen_sse_body(request_id, honored)
if subscription.nil?
close_stream_safely(stream)
else
acknowledgement = {
jsonrpc: "2.0",
method: Methods::NOTIFICATIONS_SUBSCRIPTIONS_ACKNOWLEDGED,
params: {
notifications: honored,
_meta: { RequestEnvelope::SUBSCRIPTION_ID_META_KEY.to_sym => request_id },
},
}

begin
acknowledged = subscription[:write_mutex].synchronize do
next false if subscription[:closed]
Expand Down
74 changes: 74 additions & 0 deletions test/mcp/server/transports/streamable_http_transport_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -6564,6 +6564,80 @@ def string
refute(events.any? { |event| event["method"] == "notifications/tools/list_changed" })
end

test "subscriptions/listen refuses an id that is not valid UTF-8 without registering a stream" do
body = modern_listen_body(id: "listen-BAD", params: { notifications: { toolsListChanged: true } }).b.sub("BAD", "\xFF".b)

response = @transport.handle_request(modern_rack_request(body))

assert_equal 400, response[0]

parsed = JSON.parse(response[2][0])

# The member is present and null, as JSON-RPC 2.0 requires when the id cannot be determined.
assert parsed.key?("id")
assert_nil parsed["id"]
assert_equal JsonRpcHandler::ErrorCode::INVALID_REQUEST, parsed.dig("error", "code")
assert_empty @transport.instance_variable_get(:@listen_subscriptions)
end

test "subscriptions/listen refuses a resource URI that is not valid UTF-8 without registering a stream" do
server = Server.new(name: "listen_test", capabilities: { resources: { subscribe: true } })
transport = StreamableHTTPTransport.new(server, listen_keepalive_interval: nil)
body = modern_listen_body(
id: "listen-1",
params: { notifications: { resourceSubscriptions: ["file:///BAD"] } }
).b.sub("BAD", "\xFF".b)

response = transport.handle_request(modern_rack_request(body))

assert_equal 400, response[0]

parsed = JSON.parse(response[2][0])

assert_equal "listen-1", parsed["id"]
assert_equal(-32602, parsed.dig("error", "code"))
assert_empty transport.instance_variable_get(:@listen_subscriptions)
ensure
transport.close
end

test "subscriptions/listen acknowledges without a resource URI that is not valid UTF-8 when subscribe is not declared" do
# Without the `subscribe` capability the URIs are not honored, so the acknowledgement never echoes
# the unencodable one and the request is served.
body = modern_listen_body(
id: "listen-1",
params: { notifications: { toolsListChanged: true, resourceSubscriptions: ["file:///BAD"] } },
).b.sub("BAD", "\xFF".b)

response = @transport.handle_request(modern_rack_request(body))

assert_equal 200, response[0]

io = StringIO.new
response[2].call(io)
acknowledgement = sse_events(io).first

assert_equal "notifications/subscriptions/acknowledged", acknowledgement["method"]
refute acknowledgement.dig("params", "notifications").key?("resourceSubscriptions")
end

test "subscriptions/listen refuses a request past the stream cap before checking its resource URIs" do
server = Server.new(name: "listen_test", capabilities: { resources: { subscribe: true } })
transport = StreamableHTTPTransport.new(server, max_listen_subscriptions: 1, listen_keepalive_interval: nil)
open_listen_stream(id: "listen-1", notifications: { toolsListChanged: true }, transport: transport)
body = modern_listen_body(
id: "listen-2",
params: { notifications: { resourceSubscriptions: ["file:///BAD"] } }
).b.sub("BAD", "\xFF".b)

response = transport.handle_request(modern_rack_request(body))

assert_equal 503, response[0]
assert_equal "listen-2", JSON.parse(response[2][0])["id"]
ensure
transport.close
end

test "transport close during the acknowledgement write still sends the acknowledgement first" do
# The stream blocks its first write, the acknowledgement, until released, so the close arrives while
# that write is in progress and has to queue behind it on the stream's write mutex.
Expand Down
Loading