Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
Next Next commit
Fix Streamable HTTP: surface invalid JSON response as transport error
In HttpClientStreamableHttpTransport.sendMessage, the application/json
branch completed the delivery sink before deserializing the payload.
When the server returned a body that is not valid JSON, the parse
failure only travelled through the response Flux; the delivery sink had
already completed, so McpClientSession.sendRequest never received the
error, never removed its pending response entry, and the caller waited
for the full request timeout and only saw a TimeoutException. The
original parsing exception was missing from the terminal chain.

Complete the sink only after the payload has been parsed successfully.
Notifications keep completing before their early return since they have
no response to parse.

Fixes #1147
  • Loading branch information
NewPeople-star committed Sep 30, 2026
commit 618ff855ad4fd93a476b3504e8fe689ad2f50e60
Original file line number Diff line number Diff line change
Expand Up @@ -645,16 +645,27 @@ else if (contentType.contains(TEXT_EVENT_STREAM)) {
});
}
else if (contentType.contains(APPLICATION_JSON)) {
deliveredSink.success();
String data = ((ResponseSubscribers.AggregateResponseEvent) responseEvent).data();
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
logger.warn("Notification: {} received non-compliant response: {}", sentMessage,
Utils.hasText(data) ? data : "[empty]");
deliveredSink.success();
return Mono.empty();
}

try {
return Mono.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
McpSchema.JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper, data);
// Signal delivery only after the payload has been parsed
// successfully.
// Completing the sink before deserialization would swallow a
// parse
// failure: McpClientSession relies on the error signal to
// remove the
// pending response, and without it the caller waits for the
// full
// request timeout and only sees a TimeoutException.
deliveredSink.success();
return Mono.just(message);
}
catch (IOException e) {
return Mono.error(new McpTransportException(
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
/*
* Copyright 2024-2026 the original author or authors.
*/

package io.modelcontextprotocol.client.transport;

import java.io.IOException;
import java.net.InetSocketAddress;
import java.time.Duration;

import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;

import com.sun.net.httpserver.HttpServer;

import io.modelcontextprotocol.spec.McpSchema;
import io.modelcontextprotocol.spec.McpTransportException;
import io.modelcontextprotocol.spec.ProtocolVersions;
import io.modelcontextprotocol.server.transport.TomcatTestUtil;
import reactor.test.StepVerifier;

import static org.assertj.core.api.Assertions.assertThat;

/**
* Verifies that an {@code application/json} response whose body is not valid JSON fails
* the {@link HttpClientStreamableHttpTransport#sendMessage} mono with the parsing error
* instead of completing it successfully.
*
* <p>
* Completing the delivery sink before deserialization used to swallow the parse failure:
* the {@code McpClientSession} then never received the error, kept the pending response
* entry and the caller only saw a {@code TimeoutException} once the request timeout
* elapsed.
*
* @see <a href="https://github.com/modelcontextprotocol/java-sdk/issues/1147">#1147</a>
*/
public class HttpClientStreamableHttpTransportInvalidJsonResponseTest {

static int PORT = TomcatTestUtil.findAvailablePort();

static String host = "http://localhost:" + PORT;

static HttpServer server;

@BeforeAll
static void startServer() throws IOException {
server = HttpServer.create(new InetSocketAddress(PORT), 0);

// 200 OK with an invalid JSON body for the /mcp endpoint
server.createContext("/mcp", exchange -> {
byte[] body = "{broken".getBytes();
exchange.getResponseHeaders().set("Content-Type", "application/json");
exchange.sendResponseHeaders(200, body.length);
exchange.getResponseBody().write(body);
exchange.close();
});

server.setExecutor(null);
server.start();
}

@AfterAll
static void stopServer() {
server.stop(1);
}

@Test
@Timeout(10)
void testInvalidJsonResponseFailsWithParseError() {
var transport = HttpClientStreamableHttpTransport.builder(host).build();

var initializeRequest = McpSchema.InitializeRequest
.builder(ProtocolVersions.MCP_2025_03_26, McpSchema.ClientCapabilities.builder().roots(true).build(),
McpSchema.Implementation.builder("MCP Client", "0.3.1").build())
.build();
var testMessage = new McpSchema.JSONRPCRequest(McpSchema.METHOD_INITIALIZE, "test-id", initializeRequest);

StepVerifier.create(transport.sendMessage(testMessage)).expectErrorSatisfies(error -> {
// The parse failure must surface as the delivery error, not a timeout
assertThat(error).isInstanceOf(McpTransportException.class);
}).verify(Duration.ofSeconds(5));
}

}