Skip to content
Merged
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
Prev Previous commit
Next Next commit
Release unsubscribed HttpClient bodies and classify POST errors by se…
…nt session id

Signed-off-by: Dariusz Jędrzejczyk <dariusz.jedrzejczyk@broadcom.com>
  • Loading branch information
chemicL committed Sep 29, 2026
commit 9ac6286eaca7e700831b538cbc32d063cce06502
Original file line number Diff line number Diff line change
Expand Up @@ -397,15 +397,13 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
sink.success();
}
};
Disposable connection = Mono
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
HttpResponse.BodyHandlers.ofPublisher()))
Disposable connection = ResponseBodyHandlers.sendAsync(this.httpClient, requestBuilder.build())
.flatMapMany(response -> {
if (isClosing) {
// The body is handed over as a publisher and nothing is read off
// the wire until it is subscribed, so it has to be drained even
// when its content is of no further interest.
return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize);
// 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 = response.statusCode();
Expand Down Expand Up @@ -542,16 +540,15 @@ private Mono<Void> sendHttpPost(final String endpoint, final String body) {
return Mono.from(this.httpRequestCustomizer.customize(builder, "POST", requestUri, body, transportContext));
}).flatMap(customizedBuilder -> {
var request = customizedBuilder.build();
return Mono.fromFuture(this.httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher()))
.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)));
});
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)));
});
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -241,8 +241,7 @@ private Publisher<Void> createDelete(String sessionId) {
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
return Mono.from(this.httpRequestCustomizer.customize(builder, "DELETE", uri, null, transportContext));
})
.flatMap(requestBuilder -> Mono.fromFuture(
() -> this.httpClient.sendAsync(requestBuilder.build(), HttpResponse.BodyHandlers.ofPublisher()))
.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))
Expand Down Expand Up @@ -283,8 +282,7 @@ public Mono<Void> closeGracefully() {
});
}

private Flux<McpSchema.JSONRPCMessage> consumeSseStream(
java.util.concurrent.Flow.Publisher<List<java.nio.ByteBuffer>> body,
private Flux<McpSchema.JSONRPCMessage> consumeSseStream(Flow.Publisher<List<ByteBuffer>> body,
McpTransportStream<Disposable> existingStream, Runnable onFirstMessage) {
Flux<String> lines = ResponseBodyHandlers.decodeLines(body, this.maxResponseSize);
return ResponseBodyHandlers.decodeSseResponse(lines, this.maxResponseSize).flatMap(sseEvent -> {
Expand Down Expand Up @@ -375,65 +373,69 @@ private Mono<Disposable> reconnect(McpTransportStream<Disposable> stream) {
// can be established concurrently.
Optional<String> maybeSessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);

return Mono
.fromFuture(() -> this.httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher()))
.flatMapMany(httpResponse -> {
int statusCode = httpResponse.statusCode();
Exception exception = null;
boolean proceed = false;
if (statusCode == 401 || statusCode == 403) {
logger.debug("Authorization error in reconnect with code {}", statusCode);
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
request.headers());
exception = new McpHttpClientTransportAuthorizationException(
"Authorization error connecting to SSE stream", requestSnapshot,
toResponseInfo(httpResponse));
return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMapMany(httpResponse -> {
int statusCode = httpResponse.statusCode();
Exception exception = null;
boolean proceed = false;
if (statusCode == 401 || statusCode == 403) {
logger.debug("Authorization error in reconnect with code {}", statusCode);
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
request.headers());
exception = new McpHttpClientTransportAuthorizationException(
"Authorization error connecting to SSE stream", requestSnapshot,
toResponseInfo(httpResponse));
}
else if (statusCode == METHOD_NOT_ALLOWED) {
logger.debug("The server does not support SSE streams, using request-response mode.");
}
else if (statusCode == NOT_FOUND) {
if (maybeSessionId.isPresent()) {
logger.debug("Session not found for session ID: {}", maybeSessionId.get());
String sessionIdRepresentation = sessionIdOrPlaceholder(maybeSessionId);
exception = new McpTransportSessionNotFoundException(sessionIdRepresentation);
}
else if (statusCode == METHOD_NOT_ALLOWED) {
logger.debug("The server does not support SSE streams, using request-response mode.");
else {
exception = new McpTransportException("Server Not Found. Status code:" + statusCode);
}
else if (statusCode == NOT_FOUND) {
if (maybeSessionId.isPresent()) {
logger.debug("Session not found for session ID: {}", maybeSessionId.get());
String sessionIdRepresentation = sessionIdOrPlaceholder(maybeSessionId);
exception = new McpTransportSessionNotFoundException(sessionIdRepresentation);
}
else {
exception = new McpTransportException("Server Not Found. Status code:" + statusCode);
}
}
else if (statusCode == BAD_REQUEST) {
// Some implementations return 400 when presented with a session
// id they do not know about, so the session is invalidated.
// https://github.com/modelcontextprotocol/typescript-sdk/issues/389
if (maybeSessionId.isPresent()) {
String sessionIdRepresentation = sessionIdOrPlaceholder(maybeSessionId);
exception = new McpTransportSessionNotFoundException(
"Session not found for session ID: " + sessionIdRepresentation);
}
else if (statusCode == BAD_REQUEST) {
// Unlike a POST, a GET is not treated as a session-not-found
// signal on 400: servers also reject the listening stream
// itself with 400, which must not cost the client its
// session.
else {
exception = new McpTransportException("Bad Request. Status code:" + statusCode);
}
else if (statusCode >= 200 && statusCode < 300) {
String contentType = httpResponse.headers()
.firstValue(HttpHeaders.CONTENT_TYPE)
.orElse("")
.toLowerCase();
if (contentType.contains(TEXT_EVENT_STREAM)) {
logger.debug("SSE connection established successfully");
proceed = true;
}
else {
exception = new McpTransportException(
"Unrecognized server error when connecting to SSE stream, status code: "
+ statusCode);
}
}
else if (statusCode >= 200 && statusCode < 300) {
String contentType = httpResponse.headers()
.firstValue(HttpHeaders.CONTENT_TYPE)
.orElse("")
.toLowerCase();
if (contentType.contains(TEXT_EVENT_STREAM)) {
logger.debug("SSE connection established successfully");
proceed = true;
}
else {
exception = new McpTransportException("Received unrecognized status code: " + statusCode);
exception = new McpTransportException(
"Unrecognized server error when connecting to SSE stream, status code: "
+ statusCode);
}
}
else {
exception = new McpTransportException("Received unrecognized status code: " + statusCode);
}

return proceed ? consumeSseStream(httpResponse.body(), stream, null)
: exception != null
? ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
exception)
: ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
});
return proceed ? consumeSseStream(httpResponse.body(), stream, null)
: exception != null
? ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
exception)
: ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
});
})
.retryWhen(authorizationErrorRetrySpec())
.flatMap(jsonrpcMessage -> requestHandler.apply(Mono.just(jsonrpcMessage)))
Expand Down Expand Up @@ -547,102 +549,102 @@ public Mono<Void> 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 -> Mono
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
HttpResponse.BodyHandlers.ofPublisher()))
.flatMapMany(httpResponse -> {
int statusCode = httpResponse.statusCode();
Optional<String> maybeSessionId = transportSession == null ? Optional.empty()
: transportSession.sessionId();
if (statusCode == 401 || statusCode == 403) {
logger.debug("Authorization error in sendMessage with code {}", statusCode);
var request = requestBuilder.build();
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)));
}
}).flatMapMany(requestBuilder -> {
var request = requestBuilder.build();
// 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.
Optional<String> maybeSessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);

if (transportSession
.markInitialized(httpResponse.headers().firstValue("mcp-session-id").orElse(null))) {
reconnect(null).contextWrite(deliveredSink.contextView()).subscribe();
}
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(maybeSessionId);

if (statusCode >= 200 && statusCode < 300) {
String contentType = httpResponse.headers()
.firstValue(HttpHeaders.CONTENT_TYPE)
.orElse("")
.toLowerCase();
String contentLength = httpResponse.headers()
.firstValue(HttpHeaders.CONTENT_LENGTH)
.orElse(null);

if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
logger.debug("No body returned for POST in session {}", sessionRepresentation);
markDelivered.run();
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
}
else if (contentType.contains(TEXT_EVENT_STREAM)) {
return consumeSseStream(httpResponse.body(), null, markDelivered);
}
else if (contentType.contains(APPLICATION_JSON)) {
return ResponseBodyHandlers
.decodeAggregateResponse(httpResponse.body(), this.maxResponseSize)
.flatMapMany(data -> {
markDelivered.run();
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
logger.warn("Notification: {} received non-compliant response: {}",
sentMessage, Utils.hasText(data) ? data : "[empty]");
return Flux.empty();
}
try {
return Flux.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
}
catch (IOException e) {
return Flux.<McpSchema.JSONRPCMessage>error(new McpTransportException(
"Error deserializing JSON-RPC message", e));
}
});
}

logger.warn("Unknown media type {} returned for POST in session {}", contentType,
sessionRepresentation);
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Unknown media type returned: " + contentType));
if (transportSession
.markInitialized(httpResponse.headers().firstValue("mcp-session-id").orElse(null))) {
reconnect(null).contextWrite(deliveredSink.contextView()).subscribe();
}

String sessionRepresentation = sessionIdOrPlaceholder(maybeSessionId);

if (statusCode >= 200 && statusCode < 300) {
String contentType = httpResponse.headers()
.firstValue(HttpHeaders.CONTENT_TYPE)
.orElse("")
.toLowerCase();
String contentLength = httpResponse.headers()
.firstValue(HttpHeaders.CONTENT_LENGTH)
.orElse(null);

if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
logger.debug("No body returned for POST in session {}", sessionRepresentation);
markDelivered.run();
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
}
else if (statusCode == NOT_FOUND) {
if (maybeSessionId.isPresent()) {
logger.debug("Session not found for session ID: {}", sessionRepresentation);
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportSessionNotFoundException(
"Session not found for session ID: " + sessionRepresentation));
}
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Server Not Found. Status code:" + statusCode));
else if (contentType.contains(TEXT_EVENT_STREAM)) {
return consumeSseStream(httpResponse.body(), null, markDelivered);
}
else if (contentType.contains(APPLICATION_JSON)) {
return ResponseBodyHandlers
.decodeAggregateResponse(httpResponse.body(), this.maxResponseSize)
.flatMapMany(data -> {
markDelivered.run();
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
logger.warn("Notification: {} received non-compliant response: {}", sentMessage,
Utils.hasText(data) ? data : "[empty]");
return Flux.empty();
}
try {
return Flux.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
}
catch (IOException e) {
return Flux.<McpSchema.JSONRPCMessage>error(
new McpTransportException("Error deserializing JSON-RPC message", e));
}
});
}
else if (statusCode == BAD_REQUEST) {
if (maybeSessionId.isPresent()) {
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportSessionNotFoundException(
"Session not found for session ID: " + sessionRepresentation));
}

logger.warn("Unknown media type {} returned for POST in session {}", contentType,
sessionRepresentation);
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Unknown media type returned: " + contentType));
}
else if (statusCode == NOT_FOUND) {
if (maybeSessionId.isPresent()) {
logger.debug("Session not found for session ID: {}", sessionRepresentation);
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Bad Request. Status code:" + statusCode));
new McpTransportSessionNotFoundException(
"Session not found for session ID: " + sessionRepresentation));
}
else if (statusCode >= 400 && statusCode < 500) {
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Server Not Found. Status code:" + statusCode));
}
else if (statusCode == BAD_REQUEST) {
if (maybeSessionId.isPresent()) {
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Invalid request. Status code: " + statusCode));
new McpTransportSessionNotFoundException(
"Session not found for session ID: " + sessionRepresentation));
}

return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Failed to send message, status code: " + statusCode));
})
.onErrorMap(CompletionException.class, Throwable::getCause))
new McpTransportException("Bad Request. Status code:" + statusCode));
}
else if (statusCode >= 400 && statusCode < 500) {
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Invalid request. Status code: " + statusCode));
}

return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
new McpTransportException("Failed to send message, status code: " + statusCode));
});
})
.retryWhen(authorizationErrorRetrySpec())
.flatMap(jsonRpcMessage -> requestHandler.apply(Mono.just(jsonRpcMessage)))
.onErrorMap(CompletionException.class, t -> t.getCause())
Expand Down
Loading
Loading