Skip to content

Commit 618ff85

Browse files
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
1 parent c7fef64 commit 618ff85

2 files changed

Lines changed: 99 additions & 2 deletions

File tree

‎mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -645,16 +645,27 @@ else if (contentType.contains(TEXT_EVENT_STREAM)) {
645645
});
646646
}
647647
else if (contentType.contains(APPLICATION_JSON)) {
648-
deliveredSink.success();
649648
String data = ((ResponseSubscribers.AggregateResponseEvent) responseEvent).data();
650649
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
651650
logger.warn("Notification: {} received non-compliant response: {}", sentMessage,
652651
Utils.hasText(data) ? data : "[empty]");
652+
deliveredSink.success();
653653
return Mono.empty();
654654
}
655655

656656
try {
657-
return Mono.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
657+
McpSchema.JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper, data);
658+
// Signal delivery only after the payload has been parsed
659+
// successfully.
660+
// Completing the sink before deserialization would swallow a
661+
// parse
662+
// failure: McpClientSession relies on the error signal to
663+
// remove the
664+
// pending response, and without it the caller waits for the
665+
// full
666+
// request timeout and only sees a TimeoutException.
667+
deliveredSink.success();
668+
return Mono.just(message);
658669
}
659670
catch (IOException e) {
660671
return Mono.error(new McpTransportException(
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
/*
2+
* Copyright 2024-2026 the original author or authors.
3+
*/
4+
5+
package io.modelcontextprotocol.client.transport;
6+
7+
import java.io.IOException;
8+
import java.net.InetSocketAddress;
9+
import java.time.Duration;
10+
11+
import org.junit.jupiter.api.AfterAll;
12+
import org.junit.jupiter.api.BeforeAll;
13+
import org.junit.jupiter.api.Test;
14+
import org.junit.jupiter.api.Timeout;
15+
16+
import com.sun.net.httpserver.HttpServer;
17+
18+
import io.modelcontextprotocol.spec.McpSchema;
19+
import io.modelcontextprotocol.spec.McpTransportException;
20+
import io.modelcontextprotocol.spec.ProtocolVersions;
21+
import io.modelcontextprotocol.server.transport.TomcatTestUtil;
22+
import reactor.test.StepVerifier;
23+
24+
import static org.assertj.core.api.Assertions.assertThat;
25+
26+
/**
27+
* Verifies that an {@code application/json} response whose body is not valid JSON fails
28+
* the {@link HttpClientStreamableHttpTransport#sendMessage} mono with the parsing error
29+
* instead of completing it successfully.
30+
*
31+
* <p>
32+
* Completing the delivery sink before deserialization used to swallow the parse failure:
33+
* the {@code McpClientSession} then never received the error, kept the pending response
34+
* entry and the caller only saw a {@code TimeoutException} once the request timeout
35+
* elapsed.
36+
*
37+
* @see <a href="https://github.com/modelcontextprotocol/java-sdk/issues/1147">#1147</a>
38+
*/
39+
public class HttpClientStreamableHttpTransportInvalidJsonResponseTest {
40+
41+
static int PORT = TomcatTestUtil.findAvailablePort();
42+
43+
static String host = "http://localhost:" + PORT;
44+
45+
static HttpServer server;
46+
47+
@BeforeAll
48+
static void startServer() throws IOException {
49+
server = HttpServer.create(new InetSocketAddress(PORT), 0);
50+
51+
// 200 OK with an invalid JSON body for the /mcp endpoint
52+
server.createContext("/mcp", exchange -> {
53+
byte[] body = "{broken".getBytes();
54+
exchange.getResponseHeaders().set("Content-Type", "application/json");
55+
exchange.sendResponseHeaders(200, body.length);
56+
exchange.getResponseBody().write(body);
57+
exchange.close();
58+
});
59+
60+
server.setExecutor(null);
61+
server.start();
62+
}
63+
64+
@AfterAll
65+
static void stopServer() {
66+
server.stop(1);
67+
}
68+
69+
@Test
70+
@Timeout(10)
71+
void testInvalidJsonResponseFailsWithParseError() {
72+
var transport = HttpClientStreamableHttpTransport.builder(host).build();
73+
74+
var initializeRequest = McpSchema.InitializeRequest
75+
.builder(ProtocolVersions.MCP_2025_03_26, McpSchema.ClientCapabilities.builder().roots(true).build(),
76+
McpSchema.Implementation.builder("MCP Client", "0.3.1").build())
77+
.build();
78+
var testMessage = new McpSchema.JSONRPCRequest(McpSchema.METHOD_INITIALIZE, "test-id", initializeRequest);
79+
80+
StepVerifier.create(transport.sendMessage(testMessage)).expectErrorSatisfies(error -> {
81+
// The parse failure must surface as the delivery error, not a timeout
82+
assertThat(error).isInstanceOf(McpTransportException.class);
83+
}).verify(Duration.ofSeconds(5));
84+
}
85+
86+
}

0 commit comments

Comments
 (0)