From 618ff855ad4fd93a476b3504e8fe689ad2f50e60 Mon Sep 17 00:00:00 2001
From: NewPeople-star <232190515+NewPeople-star@users.noreply.github.com>
Date: Wed, 30 Sep 2026 11:11:32 +0800
Subject: [PATCH 1/3] 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
---
.../HttpClientStreamableHttpTransport.java | 15 +++-
...eHttpTransportInvalidJsonResponseTest.java | 86 +++++++++++++++++++
2 files changed, 99 insertions(+), 2 deletions(-)
create mode 100644 mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportInvalidJsonResponseTest.java
diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java
index 5517823b6..31bd4c3a6 100644
--- a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java
+++ b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java
@@ -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(
diff --git a/mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportInvalidJsonResponseTest.java b/mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportInvalidJsonResponseTest.java
new file mode 100644
index 000000000..c1b988624
--- /dev/null
+++ b/mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportInvalidJsonResponseTest.java
@@ -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.
+ *
+ *
+ * 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 #1147
+ */
+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));
+ }
+
+}
From 1cf7903935ac6a99ea8920b4ce85a43e4f1d0689 Mon Sep 17 00:00:00 2001
From: Daniel Garnier-Moiroux
Date: Wed, 30 Sep 2026 15:05:36 +0200
Subject: [PATCH 2/3] Refactor HttpClient-based transports to use Publisher
instead of Subscriber (#1079)
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
HttpClient-based transports used to capture the enclosing sseSink in the subscriber, leading to HttpClient leaks when the client was closed. This PR addresses this, and adds many other improvements to HttpClient-based transports.
This has no public API change.
Improved transport reliability:
- Closed transports no longer keep their `HttpClient` alive, so selector threads and memory stop piling up
- `closeGracefully()` now releases open connections even when the session DELETE fails, e.g. when the server is down
- `connect()` on the legacy SSE transport no longer hangs when it gets the stream ends (error, stream closed,`closeGracefully()`, ...) before the first event or
- `sendMessage()` on Streamable HTTP no longer hangs when the SSE stream is closed without response or before the response arrives
- Responses the client never reads are always released (e.g. `DELETE`), so connections go back to the pool.
Errors surface immediately instead of as timeouts:
- On Streamable HTTP, a JSON response that can't be read (malformed, or over maxResponseSize) now fails the request immediately with the real cause, instead of a TimeoutException after requestTimeout`
- A server that answers a request with an empty JSON body now makes that request fail instead of silently timing out. An empty body in reply to a notification is still tolerated.
- Server-caused errors are now McpTransportException instead of a plain RuntimeException, and the message includes the response body the server sent.
- Errors that happen after connect() or sendMessage() has already completed now reach the transport's exception handler instead of Reactor "onErrorDropped" logs.
- A 404 or 400 invalidates the session only if the request that got it carried a session id. Fixes a race condition where are reconnect got a new session while another request was already in flight (with the old session). The old request could have ended up invalidating the new session.
Performance
- Large SSE responses, such as multi-MB tool results, are no longer slow to receive.
SSE parsing spec compliance
- Unknown fields such as retry: are ignored instead of failing the stream with "Invalid SSE response".
- The event type resets after each event, so a message that follows a named event is no longer misclassified and dropped.
- A data: line containing U+2028, U+2029 or U+0085 is no longer truncated.
- The legacy SSE transport skips empty "primer" events and unknown event types instead of failing.
- An empty id: clears the last event id.
Fixes #547
Fixes #620
Fixes #1042
Fixes #1047
Fixes #1147
Signed-off-by: Daniel Garnier-Moiroux
Signed-off-by: Dariusz Jędrzejczyk
---
.../HttpClientSseClientTransport.java | 150 +++--
.../HttpClientStreamableHttpTransport.java | 519 ++++++---------
.../transport/ResponseBodyHandlers.java | 576 ++++++++++++++++
.../client/transport/ResponseSubscribers.java | 619 ------------------
.../spec/DefaultMcpTransportSession.java | 1 +
.../transport/BoundedBodySubscriberTests.java | 334 ----------
.../HttpClientHttpTransportLeakTests.java | 107 +++
...pClientSseClientTransportConnectTests.java | 109 +++
...reamableHttpTransportSendMessageTests.java | 121 ++++
.../transport/LargeSseEventDecodingTests.java | 219 +++++++
.../transport/LoopbackMcpHttpServer.java | 136 ++++
.../ResponseBodyHandlersSendAsyncTests.java | 78 +++
.../client/transport/SseEventParserTests.java | 185 ++++++
.../transport/Utf8LineDecoderBoundTests.java | 131 ++++
.../transport/Utf8LineDecoderTests.java | 256 ++++++++
.../spec/DefaultMcpTransportSessionTests.java | 26 +
.../HttpClientBoundedReadTestSupport.java | 98 ++-
...entSseClientTransportBoundedReadTests.java | 28 +-
...reamableHttpTransportBoundedReadTests.java | 36 +-
...bleHttpTransportEmptyJsonResponseTest.java | 94 ---
...amableHttpTransportEmptyResponseTests.java | 147 +++++
...amableHttpTransportLargeResponseTests.java | 295 +++++++++
22 files changed, 2809 insertions(+), 1456 deletions(-)
create mode 100644 mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseBodyHandlers.java
delete mode 100644 mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseSubscribers.java
delete mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/BoundedBodySubscriberTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/HttpClientHttpTransportLeakTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/HttpClientSseClientTransportConnectTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportSendMessageTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/LargeSseEventDecodingTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/LoopbackMcpHttpServer.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/ResponseBodyHandlersSendAsyncTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/SseEventParserTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/Utf8LineDecoderBoundTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/client/transport/Utf8LineDecoderTests.java
create mode 100644 mcp-core/src/test/java/io/modelcontextprotocol/spec/DefaultMcpTransportSessionTests.java
delete mode 100644 mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportEmptyJsonResponseTest.java
create mode 100644 mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportEmptyResponseTests.java
create mode 100644 mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportLargeResponseTests.java
diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientSseClientTransport.java b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientSseClientTransport.java
index 3e4b613ff..9ed5c5cd4 100644
--- a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientSseClientTransport.java
+++ b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientSseClientTransport.java
@@ -11,12 +11,11 @@
import java.net.http.HttpResponse;
import java.time.Duration;
import java.util.List;
-import java.util.concurrent.CompletableFuture;
+import java.util.Optional;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
import java.util.function.Function;
-import io.modelcontextprotocol.client.transport.ResponseSubscribers.ResponseEvent;
import io.modelcontextprotocol.client.transport.customizer.McpAsyncHttpClientRequestCustomizer;
import io.modelcontextprotocol.client.transport.customizer.McpSyncHttpClientRequestCustomizer;
import io.modelcontextprotocol.common.McpTransportContext;
@@ -390,70 +389,93 @@ public Mono connect(Function, Mono> h
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
return Mono.from(this.httpRequestCustomizer.customize(builder, "GET", uri, null, transportContext));
}).flatMap(requestBuilder -> Mono.create(sink -> {
- Disposable connection = Flux.create(
- sseSink -> this.httpClient
- .sendAsync(requestBuilder.build(),
- responseInfo -> ResponseSubscribers.sseToBodySubscriber(responseInfo, sseSink,
- this.maxResponseSize))
- .exceptionallyCompose(e -> {
- sseSink.error(e);
- return CompletableFuture.failedFuture(e);
- }))
- .map(responseEvent -> (ResponseSubscribers.SseResponseEvent) responseEvent)
- .flatMap(responseEvent -> {
+ Disposable connection = ResponseBodyHandlers.sendAsync(this.httpClient, requestBuilder.build())
+ .flatMapMany(response -> {
if (isClosing) {
- return Mono.empty();
+ // The body is handed over as a publisher and the connection is
+ // only released once it is subscribed to. It is an SSE stream
+ // that may never end, so it is cancelled rather than drained.
+ return ResponseBodyHandlers.cancel(response.body());
}
- int statusCode = responseEvent.responseInfo().statusCode();
+ int statusCode = response.statusCode();
if (statusCode >= 200 && statusCode < 300) {
- try {
- if (ENDPOINT_EVENT_TYPE.equals(responseEvent.sseEvent().event())) {
- String messageEndpointUri = responseEvent.sseEvent().data();
- try {
- messageEndpointValidator.validate(uri, messageEndpointUri);
- }
- catch (InvalidSseMessageEndpointException e) {
- sink.error(e);
- this.messageEndpointSink.tryEmitError(e);
- return Flux.error(e);
- }
- if (this.messageEndpointSink.tryEmitValue(messageEndpointUri).isSuccess()) {
- sink.success();
- return Flux.empty(); // No further processing needed
- }
- else {
- sink.error(new RuntimeException("Failed to handle SSE endpoint event"));
- }
+ Flux lines = ResponseBodyHandlers.decodeLines(response.body(), this.maxResponseSize);
+ return ResponseBodyHandlers.decodeSseResponse(lines, this.maxResponseSize);
+ }
+ else {
+ return ResponseBodyHandlers.readThenError(response.body(), this.maxResponseSize,
+ "Failed to connect to SSE stream: " + statusCode);
+ }
+ })
+ // Every successfully processed event yields exactly one element, empty
+ // when it carries no message, so that the first one can mark the
+ // connection as established.
+ .>handle((sseEvent, events) -> {
+ try {
+ if (ENDPOINT_EVENT_TYPE.equals(sseEvent.event())) {
+ String messageEndpointUri = sseEvent.data();
+ try {
+ messageEndpointValidator.validate(uri, messageEndpointUri);
+ }
+ catch (InvalidSseMessageEndpointException e) {
+ this.messageEndpointSink.tryEmitError(e);
+ events.error(e);
+ return;
}
- else if (MESSAGE_EVENT_TYPE.equals(responseEvent.sseEvent().event())) {
- JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(jsonMapper,
- responseEvent.sseEvent().data());
- sink.success();
- return Flux.just(message);
+ if (this.messageEndpointSink.tryEmitValue(messageEndpointUri).isSuccess()) {
+ events.next(Optional.empty());
}
else {
- logger.debug("Received unrecognized SSE event type: {}", responseEvent.sseEvent());
- sink.success();
+ events.error(new McpTransportException("Failed to handle SSE endpoint event"));
}
}
- catch (IOException e) {
- sink.error(new McpTransportException("Error processing SSE event", e));
+ else if (MESSAGE_EVENT_TYPE.equals(sseEvent.event())) {
+ String data = sseEvent.data();
+ if (data == null || data.isBlank()) {
+ logger.debug("Skipping SSE event with empty data (stream primer)");
+ events.next(Optional.empty());
+ }
+ else {
+ events.next(Optional.of(McpSchema.deserializeJsonRpcMessage(jsonMapper, data)));
+ }
+ }
+ else {
+ logger.debug("Received unrecognized SSE event type: {}", sseEvent);
+ events.next(Optional.empty());
}
}
- return Flux.error(
- new RuntimeException("Failed to send message: " + responseEvent));
-
+ catch (IOException e) {
+ events.error(new McpTransportException("Error processing SSE event", e));
+ }
})
- .flatMap(jsonRpcMessage -> handler.apply(Mono.just(jsonRpcMessage)))
+ // connect() is resolved by the first signal only: any later failure is
+ // merely logged below, as connect() has already completed by then.
+ .switchOnFirst((first, events) -> {
+ if (first.hasValue()) {
+ sink.success();
+ }
+ else if (first.isOnError()) {
+ sink.error(first.getThrowable());
+ }
+ else if (first.isOnComplete()) {
+ sink.error(new McpTransportException("SSE stream closed before any event was received"));
+ }
+ return events;
+ })
+ .handle((message, messages) -> message.ifPresent(messages::next))
+ .flatMap(message -> handler.apply(Mono.just(message)))
.onErrorComplete(t -> {
if (!isClosing) {
logger.warn("SSE stream observed an error", t);
- sink.error(t);
}
return true;
})
+ // A closeGracefully() before the first signal cancels the stream:
+ // complete
+ // connect() instead of leaving it pending. A no-op once it has resolved.
+ .doOnCancel(sink::success)
.doFinally(s -> {
Disposable ref = this.sseSubscription.getAndSet(null);
if (ref != null && !ref.isDisposed()) {
@@ -486,17 +508,7 @@ public Mono sendMessage(JSONRPCMessage message) {
}
return this.serializeMessage(message)
- .flatMap(body -> sendHttpPost(messageEndpointUri, body).handle((response, sink) -> {
- if (response.statusCode() != 200 && response.statusCode() != 201 && response.statusCode() != 202
- && response.statusCode() != 206) {
- sink.error(new RuntimeException("Sending message failed with a non-OK HTTP code: "
- + response.statusCode() + " - " + response.body()));
- }
- else {
- sink.next(response);
- sink.complete();
- }
- }))
+ .flatMap(body -> sendHttpPost(messageEndpointUri, body))
.doOnError(error -> {
if (!isClosing) {
logger.error("Error sending message: {}", error.getMessage());
@@ -517,7 +529,16 @@ private Mono serializeMessage(final JSONRPCMessage message) {
});
}
- private Mono> sendHttpPost(final String endpoint, final String body) {
+ /**
+ * POSTs {@code body} to {@code endpoint} and consumes the response, failing if the
+ * server did not accept the message.
+ *
+ *
+ * The response body is streamed rather than aggregated: it is only read as text when
+ * a non-OK status makes it part of the failure message, and discarded otherwise.
+ * Either way it has to be consumed, or the connection is never released.
+ */
+ private Mono sendHttpPost(final String endpoint, final String body) {
final URI requestUri = Utils.resolveUri(baseUri, endpoint);
return Mono.deferContextual(ctx -> {
var builder = this.requestBuilder.copy()
@@ -529,8 +550,15 @@ private Mono> sendHttpPost(final String endpoint, final Str
return Mono.from(this.httpRequestCustomizer.customize(builder, "POST", requestUri, body, transportContext));
}).flatMap(customizedBuilder -> {
var request = customizedBuilder.build();
- return Mono.fromFuture(
- httpClient.sendAsync(request, ResponseSubscribers.boundedStringBodyHandler(this.maxResponseSize)));
+ return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMap(response -> {
+ int statusCode = response.statusCode();
+ if (statusCode == 200 || statusCode == 201 || statusCode == 202 || statusCode == 206) {
+ return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize).then();
+ }
+ return ResponseBodyHandlers.decodeAggregateResponse(response.body(), this.maxResponseSize)
+ .flatMap(text -> Mono.error(new McpTransportException(
+ "Sending message failed with a non-OK HTTP code: " + statusCode + " - " + text)));
+ });
});
}
diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java
index 5517823b6..9d55e816c 100644
--- a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java
+++ b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransport.java
@@ -9,19 +9,19 @@
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
-import java.net.http.HttpResponse.BodyHandler;
+import java.nio.ByteBuffer;
import java.time.Duration;
import java.util.Collections;
import java.util.Comparator;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.CompletionException;
+import java.util.concurrent.Flow;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
import java.util.function.Function;
import io.modelcontextprotocol.client.McpAsyncClient;
-import io.modelcontextprotocol.client.transport.ResponseSubscribers.ResponseEvent;
import io.modelcontextprotocol.client.transport.customizer.McpAsyncHttpClientRequestCustomizer;
import io.modelcontextprotocol.client.transport.customizer.McpHttpClientAuthorizationErrorHandler;
import io.modelcontextprotocol.client.transport.customizer.McpHttpClientTransportAuthorizationErrorHandler;
@@ -49,7 +49,6 @@
import org.slf4j.LoggerFactory;
import reactor.core.Disposable;
import reactor.core.publisher.Flux;
-import reactor.core.publisher.FluxSink;
import reactor.core.publisher.Mono;
import reactor.util.function.Tuple2;
import reactor.util.function.Tuples;
@@ -207,7 +206,7 @@ public static Builder builder(String baseUri) {
@Override
public Mono connect(Function, Mono> handler) {
- return Mono.deferContextual(ctx -> {
+ return Mono.defer(() -> {
this.handler.set(handler);
if (this.openConnectionOnStartup) {
logger.debug("Eagerly opening connection on startup");
@@ -240,11 +239,13 @@ private Publisher createDelete(String sessionId) {
.DELETE();
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
return Mono.from(this.httpRequestCustomizer.customize(builder, "DELETE", uri, null, transportContext));
- }).flatMap(requestBuilder -> {
- var request = requestBuilder.build();
- return Mono.fromFuture(() -> this.httpClient.sendAsync(request,
- ResponseSubscribers.boundedStringBodyHandler(this.maxResponseSize)));
- }).then();
+ })
+ .flatMap(requestBuilder -> ResponseBodyHandlers.sendAsync(this.httpClient, requestBuilder.build())
+ // The response is not inspected, but the body still has to be consumed
+ // to release the connection.
+ .flatMapMany(response -> ResponseBodyHandlers.drain(response.body(), this.maxResponseSize))
+ .then())
+ .then();
}
@Override
@@ -254,7 +255,8 @@ public void setExceptionHandler(Consumer handler) {
}
private void handleException(Throwable t) {
- logger.debug("Handling exception for session {}", sessionIdOrPlaceholder(this.activeSession.get()), t);
+ logger.debug("Handling exception for session {}", sessionIdOrPlaceholder(
+ activeSession.get() != null ? activeSession.get().sessionId() : Optional.empty()), t);
if (t instanceof McpTransportSessionNotFoundException) {
McpTransportSession> invalidSession = this.activeSession.getAndSet(createTransportSession());
logger.warn("Server does not recognize session {}. Invalidating.", invalidSession.sessionId());
@@ -266,6 +268,15 @@ private void handleException(Throwable t) {
}
}
+ private void handleExceptionSafely(Throwable t) {
+ try {
+ handleException(t);
+ }
+ catch (Exception e) {
+ logger.error("Error handling exception {}", t.getMessage(), e);
+ }
+ }
+
@Override
public Mono closeGracefully() {
return Mono.defer(() -> {
@@ -279,6 +290,38 @@ public Mono closeGracefully() {
});
}
+ /**
+ * Every successfully processed event yields exactly one element, empty when it
+ * carries no message, so that callers can tell when the first one has arrived.
+ */
+ private Flux> consumeSseStream(Flow.Publisher> body,
+ McpTransportStream existingStream) {
+ Flux lines = ResponseBodyHandlers.decodeLines(body, this.maxResponseSize);
+ return ResponseBodyHandlers.decodeSseResponse(lines, this.maxResponseSize).flatMap(sseEvent -> {
+ if (!isMessageEvent(sseEvent.event())) {
+ logger.debug("Received SSE event with type: {}", sseEvent);
+ return Flux.just(Optional.empty());
+ }
+ String data = sseEvent.data();
+ if (data == null || data.isBlank()) {
+ logger.debug("Skipping SSE event with empty data (stream primer)");
+ return Flux.just(Optional.empty());
+ }
+ try {
+ McpSchema.JSONRPCMessage message = McpSchema.deserializeJsonRpcMessage(this.jsonMapper, data);
+ Tuple2, Iterable> idWithMessages = Tuples
+ .of(Optional.ofNullable(sseEvent.id()), List.of(message));
+ McpTransportStream sessionStream = existingStream != null ? existingStream
+ : new DefaultMcpTransportStream<>(this.resumableStreams, this::reconnect);
+ return Flux.from(sessionStream.consumeSseStream(Flux.just(idWithMessages))).map(Optional::of);
+ }
+ catch (IOException e) {
+ return Flux.>error(
+ new McpTransportException("Error parsing JSON-RPC message: " + data, e));
+ }
+ });
+ }
+
private Mono reconnect(McpTransportStream stream) {
return Mono.deferContextual(ctx -> {
var rh = this.handler.get();
@@ -304,9 +347,8 @@ private Mono reconnect(McpTransportStream stream) {
final AtomicReference disposableRef = new AtomicReference<>();
- var uri = Utils.resolveUri(this.baseUri, this.endpoint);
-
Disposable connection = Mono.deferContextual(connectionCtx -> {
+ var uri = Utils.resolveUri(this.baseUri, this.endpoint);
HttpRequest.Builder requestBuilder = this.requestBuilder.copy();
if (transportSession != null && transportSession.sessionId().isPresent()) {
@@ -327,124 +369,54 @@ private Mono reconnect(McpTransportStream stream) {
.GET();
var transportContext = connectionCtx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
return Mono.from(this.httpRequestCustomizer.customize(builder, "GET", uri, null, transportContext));
+ }).flatMapMany(requestBuilder -> {
+ var request = requestBuilder.build();
+ return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMapMany(httpResponse -> {
+ int statusCode = httpResponse.statusCode();
+ if (statusCode == 401 || statusCode == 403) {
+ logger.debug("Authorization error in reconnect with code {}", statusCode);
+ var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
+ request.headers());
+ return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
+ new McpHttpClientTransportAuthorizationException(
+ "Authorization error connecting to SSE stream", requestSnapshot,
+ toResponseInfo(httpResponse)));
+ }
+ if (statusCode == METHOD_NOT_ALLOWED) {
+ logger.debug("The server does not support SSE streams, using request-response mode.");
+ return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
+ }
+ if (statusCode < 200 || statusCode >= 300) {
+ return statusError(request, httpResponse);
+ }
+ String contentType = httpResponse.headers()
+ .firstValue(HttpHeaders.CONTENT_TYPE)
+ .orElse("")
+ .toLowerCase();
+ if (!contentType.contains(TEXT_EVENT_STREAM)) {
+ return ResponseBodyHandlers.readThenError(httpResponse.body(), this.maxResponseSize,
+ "Unrecognized server error when connecting to SSE stream, status code: " + statusCode);
+ }
+ logger.debug("SSE connection established successfully");
+ return consumeSseStream(httpResponse.body(), stream);
+ });
})
- .flatMapMany(requestBuilder -> Flux.create(sseSink -> this.httpClient
- .sendAsync(requestBuilder.build(), this.toSendMessageBodySubscriber(sseSink))
- .whenComplete((response, throwable) -> {
- if (throwable != null) {
- sseSink.error(throwable);
- }
- else {
- logger.debug("SSE connection established successfully");
- }
- })).flatMap(responseEvent -> {
- int statusCode = responseEvent.responseInfo().statusCode();
- if (statusCode == 401 || statusCode == 403) {
- logger.debug("Authorization error in reconnect with code {}", statusCode);
- var request = requestBuilder.build();
- var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
- request.headers());
- return Mono.error(
- new McpHttpClientTransportAuthorizationException(
- "Authorization error connecting to SSE stream", requestSnapshot,
- responseEvent.responseInfo()));
- }
- else if (statusCode == METHOD_NOT_ALLOWED) {
- logger.debug("The server does not support SSE streams, using request-response mode.");
- return Flux.empty();
- }
-
- if (!(responseEvent instanceof ResponseSubscribers.SseResponseEvent sseResponseEvent)) {
- return Flux.error(new McpTransportException(
- "Unrecognized server error when connecting to SSE stream, status code: "
- + statusCode));
- }
- else if (statusCode >= 200 && statusCode < 300) {
- if (isMessageEvent(sseResponseEvent.sseEvent().event())) {
- String data = sseResponseEvent.sseEvent().data();
- // Per 2025-11-25 spec (SEP-1699), servers may
- // send SSE events
- // with empty data to prime the client for
- // reconnection.
- // Skip these events as they contain no JSON-RPC
- // message.
- if (data == null || data.isBlank()) {
- logger.debug("Skipping SSE event with empty data (stream primer)");
- return Flux.empty();
- }
- try {
- // We don't support batching ATM and probably
- // won't since the next version considers
- // removing it.
- McpSchema.JSONRPCMessage message = McpSchema
- .deserializeJsonRpcMessage(this.jsonMapper, data);
-
- Tuple2, Iterable> idWithMessages = Tuples
- .of(Optional.ofNullable(sseResponseEvent.sseEvent().id()), List.of(message));
-
- McpTransportStream sessionStream = stream != null ? stream
- : new DefaultMcpTransportStream<>(this.resumableStreams, this::reconnect);
- logger.debug("Connected stream {}", sessionStream.streamId());
-
- return Flux.from(sessionStream.consumeSseStream(Flux.just(idWithMessages)));
-
- }
- catch (IOException ioException) {
- return Flux.error(new McpTransportException(
- "Error parsing JSON-RPC message: " + responseEvent, ioException));
- }
- }
- else {
- logger.debug("Received SSE event with type: {}", sseResponseEvent.sseEvent());
- return Flux.empty();
- }
- }
- else if (statusCode == NOT_FOUND) {
-
- if (transportSession != null && transportSession.sessionId().isPresent()) {
- // only if the request was sent with a session id
- // and the response is 404, we consider it a
- // session not found error.
- logger.debug("Session not found for session ID: {}",
- transportSession.sessionId().get());
- String sessionIdRepresentation = sessionIdOrPlaceholder(transportSession);
- McpTransportSessionNotFoundException exception = new McpTransportSessionNotFoundException(
- "Session not found for session ID: " + sessionIdRepresentation);
- return Flux.error(exception);
- }
- return Flux.error(
- new McpTransportException("Server Not Found. Status code:" + statusCode
- + ", response-event:" + responseEvent));
- }
- else if (statusCode == BAD_REQUEST) {
- if (transportSession != null && transportSession.sessionId().isPresent()) {
- // only if the request was sent with a session id
- // and thre response is 404, we consider it a
- // session not found error.
- String sessionIdRepresentation = sessionIdOrPlaceholder(transportSession);
- McpTransportSessionNotFoundException exception = new McpTransportSessionNotFoundException(
- "Session not found for session ID: " + sessionIdRepresentation);
- return Flux.error(exception);
- }
- return Flux.error(new McpTransportException(
- "Bad Request. Status code:" + statusCode + ", response-event:" + responseEvent));
- }
- return Flux.error(new McpTransportException(
- "Received unrecognized SSE event type: " + sseResponseEvent.sseEvent().event()));
- })
- .retryWhen(authorizationErrorRetrySpec())
- .flatMap(jsonrpcMessage -> requestHandler.apply(Mono.just(jsonrpcMessage)))
- .onErrorMap(CompletionException.class, t -> t.getCause())
- .doFinally(s -> {
- Disposable ref = disposableRef.getAndSet(null);
- if (ref != null) {
- transportSession.removeConnection(ref);
- }
- }))
+ .retryWhen(authorizationErrorRetrySpec()).handle((message, messages) -> message.ifPresent(messages::next))
+ .flatMap(jsonrpcMessage -> requestHandler.apply(Mono.just(jsonrpcMessage)))
.onErrorComplete(t -> {
+ if (t instanceof CompletionException) {
+ t = t.getCause();
+ }
this.handleException(t);
return true;
})
+ .doFinally(s -> {
+ Disposable ref = disposableRef.getAndSet(null);
+ if (ref != null) {
+ transportSession.removeConnection(ref);
+ }
+ })
.contextWrite(ctx)
.subscribe();
@@ -455,6 +427,14 @@ else if (statusCode == BAD_REQUEST) {
}
+ private static HttpResponse.ResponseInfo toResponseInfo(HttpResponse>> response) {
+ return new HttpClientResponseInfo(response.statusCode(), response.headers(), response.version());
+ }
+
+ private record HttpClientResponseInfo(int statusCode, java.net.http.HttpHeaders headers,
+ HttpClient.Version version) implements HttpResponse.ResponseInfo {
+ }
+
private Retry authorizationErrorRetrySpec() {
return Retry.from(companion -> companion.flatMap(retrySignal -> {
if (!(retrySignal.failure() instanceof McpHttpClientTransportAuthorizationException authException)) {
@@ -475,31 +455,6 @@ private Retry authorizationErrorRetrySpec() {
}));
}
- private BodyHandler toSendMessageBodySubscriber(FluxSink sink) {
-
- BodyHandler responseBodyHandler = responseInfo -> {
-
- String contentType = responseInfo.headers().firstValue(HttpHeaders.CONTENT_TYPE).orElse("").toLowerCase();
-
- if (contentType.contains(TEXT_EVENT_STREAM)) {
- // For SSE streams, use line subscriber that returns Void
- logger.debug("Received SSE stream response, using line subscriber");
- return ResponseSubscribers.sseToBodySubscriber(responseInfo, sink, this.maxResponseSize);
- }
- else if (contentType.contains(APPLICATION_JSON)) {
- // For JSON responses and others, use string subscriber
- logger.debug("Received response, using string subscriber");
- return ResponseSubscribers.aggregateBodySubscriber(responseInfo, sink, this.maxResponseSize);
- }
-
- logger.debug("Received Bodyless response, using discarding subscriber");
- return ResponseSubscribers.bodilessBodySubscriber(responseInfo, sink, this.maxResponseSize);
- };
-
- return responseBodyHandler;
-
- }
-
public String toString(McpSchema.JSONRPCMessage message) {
try {
return this.jsonMapper.writeValueAsString(message);
@@ -529,9 +484,6 @@ public Mono sendMessage(McpSchema.JSONRPCMessage sentMessage) {
final AtomicReference disposableRef = new AtomicReference<>();
- var uri = Utils.resolveUri(this.baseUri, this.endpoint);
- String jsonBody = this.toString(sentMessage);
-
Disposable connection = Mono.deferContextual(ctx -> {
HttpRequest.Builder requestBuilder = this.requestBuilder.copy();
@@ -540,6 +492,8 @@ public Mono sendMessage(McpSchema.JSONRPCMessage sentMessage) {
transportSession.sessionId().get());
}
+ String jsonBody = this.toString(sentMessage);
+ var uri = Utils.resolveUri(this.baseUri, this.endpoint);
var builder = requestBuilder.uri(uri)
.header(HttpHeaders.ACCEPT, APPLICATION_JSON + ", " + TEXT_EVENT_STREAM)
.header(HttpHeaders.CONTENT_TYPE, APPLICATION_JSON_UTF8)
@@ -551,179 +505,114 @@ public Mono sendMessage(McpSchema.JSONRPCMessage sentMessage) {
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
return Mono
.from(this.httpRequestCustomizer.customize(builder, "POST", uri, jsonBody, transportContext));
- }).flatMapMany(requestBuilder -> Flux.create(responseEventSink -> {
- // Create the async request with proper body subscriber selection
- Mono.fromFuture(this.httpClient
- .sendAsync(requestBuilder.build(), this.toSendMessageBodySubscriber(responseEventSink))
- .whenComplete((response, throwable) -> {
- if (throwable != null) {
- responseEventSink.error(throwable);
- }
- else {
- logger.debug("SSE connection established successfully");
- }
- })).onErrorMap(CompletionException.class, t -> t.getCause()).onErrorComplete().subscribe();
-
- }).flatMap(responseEvent -> {
- int statusCode = responseEvent.responseInfo().statusCode();
- if (statusCode == 401 || statusCode == 403) {
- var request = requestBuilder.build();
- var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(), request.headers());
- logger.debug("Authorization error in sendMessage with code {}", statusCode);
- return Mono.error(new McpHttpClientTransportAuthorizationException(
- "Authorization error when sending message", requestSnapshot, responseEvent.responseInfo()));
- }
-
- if (transportSession.markInitialized(
- responseEvent.responseInfo().headers().firstValue("mcp-session-id").orElseGet(() -> null))) {
- // Once we have a session, we try to open an async stream for
- // the server to send notifications and requests out-of-band.
-
- reconnect(null).contextWrite(deliveredSink.contextView()).subscribe();
- }
+ }).flatMapMany(requestBuilder -> {
+ var request = requestBuilder.build();
+ return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMapMany(httpResponse -> {
+ int statusCode = httpResponse.statusCode();
+ if (statusCode == 401 || statusCode == 403) {
+ logger.debug("Authorization error in sendMessage with code {}", statusCode);
+ var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
+ request.headers());
+ return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
+ new McpHttpClientTransportAuthorizationException(
+ "Authorization error when sending message", requestSnapshot,
+ toResponseInfo(httpResponse)));
+ }
- String sessionRepresentation = sessionIdOrPlaceholder(transportSession);
+ if (transportSession
+ .markInitialized(httpResponse.headers().firstValue("mcp-session-id").orElse(null))) {
+ // Fails only when the transport has been closed in the meantime,
+ // in which case there is no stream left to open.
+ reconnect(null).contextWrite(deliveredSink.contextView()).subscribe(ignored -> {
+ }, t -> logger.debug("Not opening the SSE stream: {}", t.getMessage()));
+ }
- if (statusCode >= 200 && statusCode < 300) {
+ if (statusCode < 200 || statusCode >= 300) {
+ return statusError(request, httpResponse);
+ }
- String contentType = responseEvent.responseInfo()
- .headers()
+ String sessionRepresentation = sessionIdOrPlaceholder(
+ request.headers().firstValue(HttpHeaders.MCP_SESSION_ID));
+ String contentType = httpResponse.headers()
.firstValue(HttpHeaders.CONTENT_TYPE)
.orElse("")
.toLowerCase();
+ String contentLength = httpResponse.headers().firstValue(HttpHeaders.CONTENT_LENGTH).orElse(null);
- String contentLength = responseEvent.responseInfo()
- .headers()
- .firstValue(HttpHeaders.CONTENT_LENGTH)
- .orElse(null);
-
- // For empty content or HTTP code 202 (ACCEPTED), assume success
if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
- // if (contentType.isBlank() || "0".equals(contentLength)) {
logger.debug("No body returned for POST in session {}", sessionRepresentation);
- // No content type means no response body, so we can just
- // return an empty stream
- deliveredSink.success();
- return Flux.empty();
+ return ResponseBodyHandlers.>drain(httpResponse.body(),
+ this.maxResponseSize)
+ .startWith(Optional.empty());
}
else if (contentType.contains(TEXT_EVENT_STREAM)) {
- return Flux.just(((ResponseSubscribers.SseResponseEvent) responseEvent).sseEvent())
- .flatMap(sseEvent -> {
- String data = sseEvent.data();
- // Per 2025-11-25 spec (SEP-1699), servers may send SSE
- // events
- // with empty data to prime the client for reconnection.
- // Skip these events as they contain no JSON-RPC message.
- if (data == null || data.isBlank()) {
- logger.debug("Skipping SSE event with empty data (stream primer)");
- return Flux.empty();
- }
- try {
- // We don't support batching ATM and probably
- // won't
- // since the
- // next version considers removing it.
- McpSchema.JSONRPCMessage message = McpSchema
- .deserializeJsonRpcMessage(this.jsonMapper, data);
-
- Tuple2, Iterable> idWithMessages = Tuples
- .of(Optional.ofNullable(sseEvent.id()), List.of(message));
-
- McpTransportStream sessionStream = new DefaultMcpTransportStream<>(
- this.resumableStreams, this::reconnect);
-
- logger.debug("Connected stream {}", sessionStream.streamId());
-
- deliveredSink.success();
-
- return Flux.from(sessionStream.consumeSseStream(Flux.just(idWithMessages)));
- }
- catch (IOException ioException) {
- return Flux.error(new McpTransportException(
- "Error parsing JSON-RPC message: " + responseEvent, ioException));
- }
- });
+ return consumeSseStream(httpResponse.body(), null);
}
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]");
- return Mono.empty();
- }
-
- try {
- return Mono.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
- }
- catch (IOException e) {
- return Mono.error(new McpTransportException(
- "Error deserializing JSON-RPC message: " + responseEvent, e));
- }
+ return ResponseBodyHandlers.decodeAggregateResponse(httpResponse.body(),
+ this.maxResponseSize).>handle((data, messages) -> {
+ if (sentMessage instanceof McpSchema.JSONRPCNotification) {
+ logger.warn("Notification: {} received non-compliant response: {}", sentMessage,
+ Utils.hasText(data) ? data : "[empty]");
+ messages.next(Optional.empty());
+ return;
+ }
+ try {
+ messages
+ .next(Optional.of(McpSchema.deserializeJsonRpcMessage(jsonMapper, data)));
+ }
+ catch (IOException e) {
+ messages.error(new McpTransportException(
+ "Error deserializing JSON-RPC message: " + data, e));
+ }
+ })
+ .flux();
}
+
logger.warn("Unknown media type {} returned for POST in session {}", contentType,
sessionRepresentation);
-
- return Flux.error(
- new RuntimeException("Unknown media type returned: " + contentType));
- }
- else if (statusCode == NOT_FOUND) {
- if (transportSession != null && transportSession.sessionId().isPresent()) {
- // only if the request was sent with a session id and the
- // response is 404, we consider it a session not found error.
- logger.debug("Session not found for session ID: {}", transportSession.sessionId().get());
- McpTransportSessionNotFoundException exception = new McpTransportSessionNotFoundException(
- "Session not found for session ID: " + sessionRepresentation);
- return Flux.error(exception);
- }
- return Flux.error(new McpTransportException(
- "Server Not Found. Status code:" + statusCode + ", response-event:" + responseEvent));
- }
- else if (statusCode == BAD_REQUEST) {
- // Some implementations can return 400 when presented with a
- // session id that it doesn't know about, so we will
- // invalidate the session
- // https://github.com/modelcontextprotocol/typescript-sdk/issues/389
-
- if (transportSession != null && transportSession.sessionId().isPresent()) {
- // only if the request was sent with a session id and the
- // response is 404, we consider it a session not found error.
- McpTransportSessionNotFoundException exception = new McpTransportSessionNotFoundException(
- "Session not found for session ID: " + sessionRepresentation);
- return Flux.error(exception);
- }
- return Flux.error(new McpTransportException(
- "Bad Request. Status code:" + statusCode + ", response-event:" + responseEvent));
- }
- else if (statusCode >= 400 && statusCode < 500) {
- return Flux.error(
- new McpTransportException("Invalid request. Status code: " + statusCode));
- }
-
- return Flux.error(
- new RuntimeException("Failed to send message: " + responseEvent));
+ return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
+ new McpTransportException("Unknown media type returned: " + contentType));
+ });
})
.retryWhen(authorizationErrorRetrySpec())
- .flatMap(jsonRpcMessage -> requestHandler.apply(Mono.just(jsonRpcMessage)))
.onErrorMap(CompletionException.class, t -> t.getCause())
+ // sendMessage() is resolved by the first signal only: any later failure
+ // is
+ // merely handled below, as sendMessage() has already completed by then.
+ // An exchange ending without any event still means the server accepted
+ // the message, so completion resolves it successfully too.
+ .switchOnFirst((first, messages) -> {
+ if (first.isOnError()) {
+ // Handled before failing sendMessage(), so that a session the
+ // server does not recognise is already invalidated by the time
+ // the caller learns about it. Consumed here so that it is not
+ // handled a second time below.
+ handleExceptionSafely(first.getThrowable());
+ deliveredSink.error(first.getThrowable());
+ return Flux.empty();
+ }
+ deliveredSink.success();
+ return messages;
+ }).handle((message, messages) -> message.ifPresent(messages::next))
+ .flatMap(jsonRpcMessage -> requestHandler.apply(Mono.just(jsonRpcMessage)))
.doFinally(s -> {
- logger.debug("SendMessage finally: {}", s);
Disposable ref = disposableRef.getAndSet(null);
if (ref != null) {
transportSession.removeConnection(ref);
}
- })).onErrorComplete(t -> {
- // handle the error first
- try {
- this.handleException(t);
- }
- catch (Exception e) {
- logger.error("Error handling exception {}", t.getMessage(), e);
- }
- // inform the caller of sendMessage
- deliveredSink.error(t);
+ })
+ .onErrorComplete(t -> {
+ handleExceptionSafely(t);
return true;
- }).contextWrite(deliveredSink.contextView()).subscribe();
+ })
+ // Closing the session before the first signal cancels the exchange:
+ // complete sendMessage() instead of leaving it pending. A no-op once it
+ // has resolved.
+ .doOnCancel(deliveredSink::success)
+ .contextWrite(deliveredSink.contextView())
+ .subscribe();
disposableRef.set(connection);
transportSession.addConnection(connection);
@@ -731,8 +620,32 @@ else if (statusCode >= 400 && statusCode < 500) {
}
- private static String sessionIdOrPlaceholder(McpTransportSession> transportSession) {
- return transportSession.sessionId().orElse("[missing_session_id]");
+ /**
+ * Fails the exchange over a response with an error status. A session id the server
+ * does not recognise invalidates the session; any other failure carries the response
+ * body, which is what the server said about it.
+ */
+ private Flux statusError(HttpRequest request, HttpResponse>> response) {
+ int statusCode = response.statusCode();
+ // Classify the response against the session id that this very request carried,
+ // rather than the one currently held by the session, which can be established
+ // concurrently. Some implementations return 400 rather than 404 for a session id
+ // they do not know about.
+ // https://github.com/modelcontextprotocol/typescript-sdk/issues/389
+ Optional sessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);
+ if ((statusCode == NOT_FOUND || statusCode == BAD_REQUEST) && sessionId.isPresent()) {
+ logger.debug("Session not found for session ID: {}", sessionId.get());
+ return ResponseBodyHandlers.drainThenError(response.body(), this.maxResponseSize,
+ new McpTransportSessionNotFoundException(sessionId.get()));
+ }
+ String failure = statusCode == NOT_FOUND ? "Server Not Found. Status code:" + statusCode
+ : statusCode == BAD_REQUEST ? "Bad Request. Status code:" + statusCode
+ : "Received unexpected status code: " + statusCode;
+ return ResponseBodyHandlers.readThenError(response.body(), this.maxResponseSize, failure);
+ }
+
+ private static String sessionIdOrPlaceholder(Optional sessionId) {
+ return sessionId.orElse("[missing_session_id]");
}
@Override
diff --git a/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseBodyHandlers.java b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseBodyHandlers.java
new file mode 100644
index 000000000..01e243da5
--- /dev/null
+++ b/mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseBodyHandlers.java
@@ -0,0 +1,576 @@
+/*
+ * Copyright 2024 - 2026 the original author or authors.
+ */
+
+package io.modelcontextprotocol.client.transport;
+
+import java.net.http.HttpClient;
+import java.net.http.HttpRequest;
+import java.net.http.HttpResponse;
+import java.nio.ByteBuffer;
+import java.nio.CharBuffer;
+import java.nio.charset.CharacterCodingException;
+import java.nio.charset.CharsetDecoder;
+import java.nio.charset.CoderResult;
+import java.nio.charset.CodingErrorAction;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.CancellationException;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
+import java.util.concurrent.Flow;
+import java.util.concurrent.Flow.Publisher;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.modelcontextprotocol.spec.McpTransportException;
+import io.modelcontextprotocol.util.Utils;
+import reactor.adapter.JdkFlowAdapter;
+import reactor.core.publisher.Flux;
+import reactor.core.publisher.Mono;
+
+/**
+ * Utility class providing various operations for handling different types of HTTP
+ * response bodies in the context of Model Context Protocol (MCP) clients.
+ *
+ *
+ * Defines Flux operators for processing Server-Sent Events (SSE), aggregate responses,
+ * and bodiless responses.
+ *
+ * @author Christian Tzolov
+ * @author Dariusz Jędrzejczyk
+ * @author Daniel Garnier-Moiroux
+ */
+class ResponseBodyHandlers {
+
+ /**
+ * Bytes of SSE field framing a single line may carry on top of the message payload:
+ * {@code "event: "} is the longest field prefix this parser recognises. Line
+ * terminators are not counted, as they reset the running line length. Without this
+ * allowance, an event carrying exactly the maximum message size would be rejected
+ * because of the bytes the SSE wire format adds around it.
+ */
+ private static final int SSE_FRAMING_OVERHEAD = "event: ".length();
+
+ /**
+ * The type of an SSE event that does not name one with an {@code event:} field.
+ */
+ private static final String DEFAULT_EVENT_TYPE = "message";
+
+ record SseEvent(String id, String event, String data) {
+ }
+
+ /**
+ * Adds {@link #SSE_FRAMING_OVERHEAD} to {@code maxSize}, saturating at
+ * {@link Integer#MAX_VALUE} rather than overflowing into a negative bound that would
+ * reject everything.
+ */
+ private static int plusFramingOverhead(int maxSize) {
+ return maxSize > Integer.MAX_VALUE - SSE_FRAMING_OVERHEAD ? Integer.MAX_VALUE : maxSize + SSE_FRAMING_OVERHEAD;
+ }
+
+ /**
+ * Converts a publisher of byte-buffer chunks into a flux of decoded string lines,
+ * bounding how much memory a single line may occupy.
+ *
+ *
+ * The decoder buffers characters until it encounters a line terminator, so a peer
+ * that never terminates a line (or sends an enormous one) would force the transport
+ * to buffer it in memory. Exceeding the bound fails the flux, which cancels the
+ * subscription and so closes the connection.
+ *
+ *
+ * The bound is allowed {@link #SSE_FRAMING_OVERHEAD} extra characters so that the SSE
+ * framing around a payload does not count against the payload's own budget: only SSE
+ * streams are read line by line, so every caller of this method is parsing one.
+ *
+ *
+ * This only bounds a single line. What accumulates across lines is bounded where it
+ * accumulates: see {@link #decodeSseResponse} for multi-line SSE events and
+ * {@link #decodeAggregateResponse} for whole response bodies.
+ * @param publisher the response body
+ * @param maxSize the maximum number of bytes read for a single inbound message
+ */
+ static Flux decodeLines(Publisher> publisher, int maxSize) {
+ return Flux.defer(() -> {
+ Utf8LineDecoder dec = new Utf8LineDecoder(plusFramingOverhead(maxSize));
+ return JdkFlowAdapter.flowPublisherToFlux(publisher)
+ .concatMapIterable(dec::decode)
+ .concatWith(Flux.defer(() -> Flux.fromIterable(dec.flush())));
+ });
+ }
+
+ /**
+ * Parses a flux of SSE-formatted lines into a flux of {@link SseEvent}, bounding how
+ * much memory a single event may occupy.
+ * @param lines the SSE-formatted lines to parse
+ * @param maxSize the maximum number of bytes that may accumulate for a single SSE
+ * event
+ */
+ static Flux decodeSseResponse(Flux lines, int maxSize) {
+ return Flux.defer(() -> {
+ SseEventParser parser = new SseEventParser(maxSize);
+ return lines.handle((line, sink) -> parser.feed(line).ifPresent(sink::next))
+ .concatWith(Mono.defer(() -> parser.flush().map(Mono::just).orElseGet(Mono::empty)));
+ });
+ }
+
+ /**
+ * Collects all byte-buffer chunks from the publisher into a single UTF-8 decoded
+ * string, bounding how much memory it may occupy. A peer sending a body larger than
+ * {@code maxSize} has its response aborted instead of forcing the transport to buffer
+ * it in memory.
+ * @param publisher the response body
+ * @param maxSize the maximum number of bytes read for the response body
+ */
+ static Mono decodeAggregateResponse(Publisher> publisher, int maxSize) {
+ return boundTotalBytes(publisher, maxSize).collectList().map(buffers -> {
+ int totalSize = buffers.stream().mapToInt(ByteBuffer::remaining).sum();
+ ByteBuffer combined = ByteBuffer.allocate(totalSize);
+ buffers.forEach(combined::put);
+ combined.flip();
+ return StandardCharsets.UTF_8.decode(combined).toString();
+ }).defaultIfEmpty("");
+ }
+
+ /**
+ * Subscribes to the body publisher to release the underlying connection, discarding
+ * all bytes, then propagates the given error.
+ *
+ *
+ * Nothing accumulates here, so the bound is not protecting memory: it stops a peer
+ * from making the transport read an unbounded body only to throw it away. Should the
+ * body outgrow {@code maxSize}, or fail to be read, that failure is dropped and
+ * {@code error} is propagated all the same, as it is the reason the body is being
+ * discarded in the first place.
+ * @param body the response body
+ * @param maxSize the maximum number of bytes read for the response body
+ * @param error the error to propagate once the body has been discarded
+ */
+ static Flux drainThenError(Publisher> body, int maxSize, Throwable error) {
+ return boundTotalBytes(body, maxSize).onErrorComplete().thenMany(Mono.error(error));
+ }
+
+ /**
+ * Reads the body as text, then propagates a {@link McpTransportException} carrying
+ * {@code message} followed by that text, so that what the server said about the
+ * failure reaches the caller. The body is read under the same bound as any other.
+ * @param body the response body
+ * @param maxSize the maximum number of bytes read for the response body
+ * @param message describes the failure the body explains
+ */
+ static Flux readThenError(Publisher> body, int maxSize, String message) {
+ return decodeAggregateResponse(body, maxSize).flatMapMany(text -> Flux
+ .error(new McpTransportException(Utils.hasText(text) ? message + ", response body: " + text : message)));
+ }
+
+ /**
+ * Subscribes to the body publisher to release the underlying connection, discarding
+ * all bytes, then completes empty.
+ *
+ *
+ * As in {@link #drainThenError}, the bound caps what a peer can make the transport
+ * read rather than what it can make it hold.
+ * @param body the response body
+ * @param maxSize the maximum number of bytes read for the response body
+ */
+ static Flux drain(Publisher> body, int maxSize) {
+ return boundTotalBytes(body, maxSize).thenMany(Flux.empty());
+ }
+
+ /**
+ * Subscribes to the body publisher only to cancel it, which releases the underlying
+ * connection without reading the body, then completes empty.
+ *
+ *
+ * Unlike {@link #drain}, this suits a body that may never end, such as an SSE stream
+ * whose content is of no further interest.
+ * @param body the response body
+ */
+ static Flux cancel(Publisher> body) {
+ return Flux.defer(() -> {
+ cancelBody(body);
+ return Flux.empty();
+ });
+ }
+
+ static Mono>>> sendAsync(HttpClient httpClient, HttpRequest request) {
+ // Not Mono.fromFuture: cancelling aborts the exchange, and the HttpClient then
+ // fails the future with a CompletionException wrapping a CancellationException,
+ // which fromFuture reports as a dropped error. Only this method cna cancel the
+ // future, so that failure is ignored here. Replace with a plain fromFuture,
+ // keeping
+ // the doOnDiscard, once https://github.com/reactor/reactor-core/issues/4415 is
+ // resolved.
+ return Mono.>>>create(sink -> {
+ CompletableFuture>>> exchange = httpClient.sendAsync(request,
+ HttpResponse.BodyHandlers.ofPublisher());
+ sink.onCancel(() -> exchange.cancel(true));
+ exchange.whenComplete((response, error) -> {
+ if (error == null) {
+ // Emit the response so the body can be consumed.
+ // If the surrounding Mono was cancelled though and due to a race
+ // the headers were already parsed, the below call will simply
+ // discard the response.
+ sink.success(response);
+ return;
+ }
+ Throwable cause = error instanceof CompletionException && error.getCause() != null ? error.getCause()
+ : error;
+ if (cause instanceof CancellationException) {
+ sink.success();
+ }
+ else {
+ sink.error(cause);
+ }
+ });
+ })
+ // A body that is never subscribed to never releases its connection.
+ .doOnDiscard(HttpResponse.class, response -> {
+ if (response.body() instanceof Publisher> body) {
+ cancelBody(body);
+ }
+ });
+ }
+
+ private static void cancelBody(Publisher> body) {
+ body.subscribe(CancellingSubscriber.INSTANCE);
+ }
+
+ /**
+ * Flattens the body into its individual byte buffers, failing once more than
+ * {@code maxSize} bytes have passed through. Failing cancels the subscription, which
+ * closes the connection and so stops the peer from streaming any more.
+ */
+ private static Flux boundTotalBytes(Publisher> body, int maxSize) {
+ return Flux.defer(() -> {
+ // Held in an array because the handle callback below cannot mutate a
+ // captured local. The enclosing defer gives each subscriber its own.
+ long[] totalBytes = new long[1];
+ return JdkFlowAdapter.flowPublisherToFlux(body)
+ .flatMapIterable(list -> list)
+ .handle((buffer, sink) -> {
+ totalBytes[0] += buffer.remaining();
+ if (totalBytes[0] > maxSize) {
+ sink.error(new McpTransportException(
+ "Inbound response body exceeds the maximum allowed size of " + maxSize + " bytes"));
+ return;
+ }
+ sink.next(buffer);
+ });
+ });
+ }
+
+ /**
+ * Stateful UTF-8 decoder that splits a stream of byte-buffer chunks into complete
+ * lines. Handles multi-byte characters split across chunk boundaries, and terminates
+ * a line on {@code "\r\n"}, {@code "\r"} or {@code "\n"} alike, as the SSE wire
+ * format does. Bytes that do not decode are replaced rather than reported, so a peer
+ * sending one does not cost the stream.
+ */
+ static final class Utf8LineDecoder {
+
+ /**
+ * Undecodable input costs one replacement character rather than the stream: a
+ * decoder left on the default {@link CodingErrorAction#REPORT} fails the whole
+ * response over a single byte a peer mangled, and takes with it the lines already
+ * decoded from the same chunk, because {@link #decode(List)} throws instead of
+ * returning them. A body cut short mid-character is enough to hit it. This
+ * matches {@link java.net.http.HttpResponse.BodySubscribers#fromLineSubscriber},
+ * the path this decoder replaces, which configured the same two actions.
+ */
+ private final CharsetDecoder decoder = StandardCharsets.UTF_8.newDecoder()
+ .onMalformedInput(CodingErrorAction.REPLACE)
+ .onUnmappableCharacter(CodingErrorAction.REPLACE);
+
+ private final CharBuffer charBuffer = CharBuffer.allocate(4096);
+
+ private final StringBuilder leftover = new StringBuilder();
+
+ /**
+ * The maximum number of bytes a single line may occupy. Measured against
+ * {@link #leftover}'s length in characters, which for UTF-8 is never more than
+ * the number of bytes those characters were decoded from, so a line is only ever
+ * rejected once it has genuinely exceeded the bound in bytes.
+ */
+ private final int maxSize;
+
+ /**
+ * How many leading characters of {@link #leftover} are already known to hold no
+ * line terminator, so that the search for one resumes where the previous search
+ * ended instead of restarting at the beginning of the buffer. Without it, a long
+ * line is searched again in full for every chunk that arrives, which makes
+ * reading an event cost time proportional to the square of its length.
+ * @see #1042
+ */
+ private int scannedForLineTerminator = 0;
+
+ /**
+ * Whether the line just emitted was terminated by a CR, so that a LF opening what
+ * follows completes that terminator instead of ending a line of its own. A CR is
+ * emitted on as soon as it arrives, before it is known whether a LF follows it,
+ * and the two may be split across chunks.
+ */
+ private boolean crTerminatedPreviousLine = false;
+
+ // Holds partial UTF-8 sequences left over from a previous chunk (max 3 bytes
+ // for a BMP code point; 4 bytes for a supplementary one).
+ private ByteBuffer pendingBytes = ByteBuffer.allocate(0);
+
+ Utf8LineDecoder(int maxSize) {
+ this.maxSize = maxSize;
+ }
+
+ List decode(List chunk) {
+ List lines = new ArrayList<>();
+ for (ByteBuffer bb : chunk) {
+ ByteBuffer input = bb;
+ if (pendingBytes.hasRemaining()) {
+ ByteBuffer merged = ByteBuffer.allocate(pendingBytes.remaining() + bb.remaining());
+ merged.put(pendingBytes).put(bb);
+ merged.flip();
+ pendingBytes = ByteBuffer.allocate(0);
+ input = merged;
+ }
+ while (true) {
+ CoderResult result = decoder.decode(input, charBuffer, false);
+ drainCharBuffer();
+ extractCompletedLines(lines);
+ // Unreachable while the decoder replaces undecodable input, but kept
+ // so that an error result cannot spin this loop: it is neither an
+ // underflow nor an overflow.
+ if (result.isError()) {
+ try {
+ result.throwException();
+ }
+ catch (CharacterCodingException e) {
+ throw new RuntimeException(e);
+ }
+ }
+ if (result.isUnderflow()) {
+ if (input.hasRemaining()) {
+ pendingBytes = ByteBuffer.allocate(input.remaining());
+ pendingBytes.put(input).flip();
+ }
+ break;
+ }
+ }
+ }
+ return lines;
+ }
+
+ List flush() {
+ ByteBuffer tail = pendingBytes.hasRemaining() ? pendingBytes : ByteBuffer.allocate(0);
+ CoderResult result = decoder.decode(tail, charBuffer, true);
+ while (result.isOverflow()) {
+ drainCharBuffer();
+ result = decoder.decode(tail, charBuffer, true);
+ }
+ drainCharBuffer();
+ if (result.isError()) {
+ try {
+ result.throwException();
+ }
+ catch (CharacterCodingException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ result = decoder.flush(charBuffer);
+ while (result.isOverflow()) {
+ drainCharBuffer();
+ result = decoder.flush(charBuffer);
+ }
+ drainCharBuffer();
+ pendingBytes = ByteBuffer.allocate(0);
+
+ List lines = new ArrayList<>();
+ extractCompletedLines(lines);
+ if (leftover.length() > 0) {
+ String last = leftover.toString();
+ leftover.setLength(0);
+ this.scannedForLineTerminator = 0;
+ lines.add(last);
+ }
+ this.crTerminatedPreviousLine = false;
+ return lines;
+ }
+
+ private void drainCharBuffer() {
+ charBuffer.flip();
+ leftover.append(charBuffer);
+ charBuffer.clear();
+ }
+
+ private void extractCompletedLines(List out) {
+ while (true) {
+ if (this.crTerminatedPreviousLine) {
+ if (leftover.length() == 0) {
+ // The LF, if there is one, is in a chunk that has not arrived.
+ return;
+ }
+ if (leftover.charAt(0) == '\n') {
+ leftover.delete(0, 1);
+ }
+ this.crTerminatedPreviousLine = false;
+ }
+ int terminatorIdx = indexOfLineTerminator(this.scannedForLineTerminator);
+ if (terminatorIdx == -1) {
+ this.scannedForLineTerminator = leftover.length();
+ if (leftover.length() > this.maxSize) {
+ throw new McpTransportException(
+ "Inbound line exceeds the maximum allowed size of " + this.maxSize + " bytes");
+ }
+ return;
+ }
+ out.add(leftover.substring(0, terminatorIdx));
+ this.crTerminatedPreviousLine = leftover.charAt(terminatorIdx) == '\r';
+ leftover.delete(0, terminatorIdx + 1);
+ // What is left starts after the terminator, so none of it has been
+ // searched yet.
+ this.scannedForLineTerminator = 0;
+ }
+ }
+
+ /**
+ * Index of the first CR or LF in {@link #leftover} at or after {@code from}, or
+ * {@code -1} when there is none.
+ */
+ private int indexOfLineTerminator(int from) {
+ for (int i = from; i < leftover.length(); i++) {
+ char c = leftover.charAt(i);
+ if (c == '\n' || c == '\r') {
+ return i;
+ }
+ }
+ return -1;
+ }
+
+ }
+
+ /**
+ * Stateful SSE line parser. Accumulates {@code data:}, {@code id:} and {@code event:}
+ * fields until a blank line dispatches the event. Per the SSE spec, {@code id} is the
+ * last event ID and persists across events until re-set, with an empty value clearing
+ * it; {@code event} and {@code data} are reset by every blank line, so an event that
+ * does not name its type is a {@code message} event whatever preceded it. A blank
+ * line dispatches only when a {@code data:} field was seen, whether or not it carried
+ * a value. Comments and fields the parser does not handle, such as {@code retry:},
+ * are ignored as the spec requires.
+ *
+ * @see Interpreting
+ * an event stream
+ */
+ static final class SseEventParser {
+
+ private static final Logger logger = LoggerFactory.getLogger(SseEventParser.class);
+
+ private final StringBuilder data = new StringBuilder();
+
+ /**
+ * The maximum number of bytes that may accumulate for a single SSE event. A peer
+ * that never terminates an event (e.g. an endless stream of {@code data:} lines)
+ * has its stream aborted instead of exhausting memory. The accumulated data is
+ * measured in characters, which for UTF-8 is never more than the number of bytes
+ * it was decoded from.
+ */
+ private final int maxSize;
+
+ private String id;
+
+ private String event;
+
+ SseEventParser(int maxSize) {
+ this.maxSize = maxSize;
+ }
+
+ Optional feed(String line) {
+ if (line.isEmpty()) {
+ return flush();
+ }
+ if (line.startsWith("data:")) {
+ // Every data field appends its value followed by a separator, so a
+ // valueless `data:` line still marks the event as carrying data and gets
+ // dispatched with empty data. Servers send such an event to prime a
+ // stream, and dropping it leaves the request it answers hanging.
+ String value = line.substring(5).trim();
+ // Measured before appending, so that an event carrying exactly
+ // maxSize of data is accepted: the trailing separator below is
+ // stripped again before the event is emitted.
+ if (data.length() + value.length() > this.maxSize) {
+ throw new McpTransportException(
+ "Inbound SSE event exceeds the maximum allowed size of " + this.maxSize + " bytes");
+ }
+ data.append(value).append('\n');
+ }
+ else if (line.startsWith("id:")) {
+ String value = line.substring(3).trim();
+ // The spec ignores an id carrying a NULL, and an empty id resets the last
+ // event ID, which leaves nothing to resume from.
+ if (value.indexOf('\0') == -1) {
+ id = value.isEmpty() ? null : value;
+ }
+ }
+ else if (line.startsWith("event:")) {
+ String value = line.substring(6).trim();
+ event = value.isEmpty() ? null : value;
+ }
+ else if (line.startsWith(":")) {
+ logger.debug("Ignoring comment line: {}", line);
+ }
+ else {
+ // The SSE spec mandates that fields the client does not know about, such
+ // as `retry:`, are ignored rather than treated as a protocol error.
+ logger.debug("Ignoring unknown SSE field line: {}", line);
+ }
+ return Optional.empty();
+ }
+
+ /**
+ * Emits the pending event, if a {@code data:} field was seen, and resets the
+ * per-event state. The event type is reset even when nothing is dispatched, as
+ * the spec requires, while the id is the last event ID and so survives. An event
+ * that did not name its type is emitted as a {@code message} event.
+ */
+ Optional flush() {
+ String type = this.event;
+ this.event = null;
+ if (data.isEmpty()) {
+ return Optional.empty();
+ }
+ SseEvent result = new SseEvent(id, type != null ? type : DEFAULT_EVENT_TYPE, data.toString().trim());
+ data.setLength(0);
+ return Optional.of(result);
+ }
+
+ }
+
+ private static class CancellingSubscriber implements Flow.Subscriber