Skip to content

Commit 9ac6286

Browse files
committed
Release unsubscribed HttpClient bodies and classify POST errors by sent session id
Signed-off-by: Dariusz Jędrzejczyk <dariusz.jedrzejczyk@broadcom.com>
1 parent 45f93da commit 9ac6286

3 files changed

Lines changed: 226 additions & 160 deletions

File tree

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

Lines changed: 14 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -397,15 +397,13 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
397397
sink.success();
398398
}
399399
};
400-
Disposable connection = Mono
401-
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
402-
HttpResponse.BodyHandlers.ofPublisher()))
400+
Disposable connection = ResponseBodyHandlers.sendAsync(this.httpClient, requestBuilder.build())
403401
.flatMapMany(response -> {
404402
if (isClosing) {
405-
// The body is handed over as a publisher and nothing is read off
406-
// the wire until it is subscribed, so it has to be drained even
407-
// when its content is of no further interest.
408-
return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize);
403+
// The body is handed over as a publisher and the connection is
404+
// only released once it is subscribed to. It is an SSE stream
405+
// that may never end, so it is cancelled rather than drained.
406+
return ResponseBodyHandlers.cancel(response.body());
409407
}
410408

411409
int statusCode = response.statusCode();
@@ -542,16 +540,15 @@ private Mono<Void> sendHttpPost(final String endpoint, final String body) {
542540
return Mono.from(this.httpRequestCustomizer.customize(builder, "POST", requestUri, body, transportContext));
543541
}).flatMap(customizedBuilder -> {
544542
var request = customizedBuilder.build();
545-
return Mono.fromFuture(this.httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher()))
546-
.flatMap(response -> {
547-
int statusCode = response.statusCode();
548-
if (statusCode == 200 || statusCode == 201 || statusCode == 202 || statusCode == 206) {
549-
return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize).then();
550-
}
551-
return ResponseBodyHandlers.decodeAggregateResponse(response.body(), this.maxResponseSize)
552-
.flatMap(text -> Mono.error(new McpTransportException(
553-
"Sending message failed with a non-OK HTTP code: " + statusCode + " - " + text)));
554-
});
543+
return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMap(response -> {
544+
int statusCode = response.statusCode();
545+
if (statusCode == 200 || statusCode == 201 || statusCode == 202 || statusCode == 206) {
546+
return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize).then();
547+
}
548+
return ResponseBodyHandlers.decodeAggregateResponse(response.body(), this.maxResponseSize)
549+
.flatMap(text -> Mono.error(new McpTransportException(
550+
"Sending message failed with a non-OK HTTP code: " + statusCode + " - " + text)));
551+
});
555552
});
556553
}
557554

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

Lines changed: 143 additions & 141 deletions
Original file line numberDiff line numberDiff line change
@@ -241,8 +241,7 @@ private Publisher<Void> createDelete(String sessionId) {
241241
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
242242
return Mono.from(this.httpRequestCustomizer.customize(builder, "DELETE", uri, null, transportContext));
243243
})
244-
.flatMap(requestBuilder -> Mono.fromFuture(
245-
() -> this.httpClient.sendAsync(requestBuilder.build(), HttpResponse.BodyHandlers.ofPublisher()))
244+
.flatMap(requestBuilder -> ResponseBodyHandlers.sendAsync(this.httpClient, requestBuilder.build())
246245
// The response is not inspected, but the body still has to be consumed
247246
// to release the connection.
248247
.flatMapMany(response -> ResponseBodyHandlers.drain(response.body(), this.maxResponseSize))
@@ -283,8 +282,7 @@ public Mono<Void> closeGracefully() {
283282
});
284283
}
285284

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

378-
return Mono
379-
.fromFuture(() -> this.httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher()))
380-
.flatMapMany(httpResponse -> {
381-
int statusCode = httpResponse.statusCode();
382-
Exception exception = null;
383-
boolean proceed = false;
384-
if (statusCode == 401 || statusCode == 403) {
385-
logger.debug("Authorization error in reconnect with code {}", statusCode);
386-
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
387-
request.headers());
388-
exception = new McpHttpClientTransportAuthorizationException(
389-
"Authorization error connecting to SSE stream", requestSnapshot,
390-
toResponseInfo(httpResponse));
376+
return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMapMany(httpResponse -> {
377+
int statusCode = httpResponse.statusCode();
378+
Exception exception = null;
379+
boolean proceed = false;
380+
if (statusCode == 401 || statusCode == 403) {
381+
logger.debug("Authorization error in reconnect with code {}", statusCode);
382+
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
383+
request.headers());
384+
exception = new McpHttpClientTransportAuthorizationException(
385+
"Authorization error connecting to SSE stream", requestSnapshot,
386+
toResponseInfo(httpResponse));
387+
}
388+
else if (statusCode == METHOD_NOT_ALLOWED) {
389+
logger.debug("The server does not support SSE streams, using request-response mode.");
390+
}
391+
else if (statusCode == NOT_FOUND) {
392+
if (maybeSessionId.isPresent()) {
393+
logger.debug("Session not found for session ID: {}", maybeSessionId.get());
394+
String sessionIdRepresentation = sessionIdOrPlaceholder(maybeSessionId);
395+
exception = new McpTransportSessionNotFoundException(sessionIdRepresentation);
391396
}
392-
else if (statusCode == METHOD_NOT_ALLOWED) {
393-
logger.debug("The server does not support SSE streams, using request-response mode.");
397+
else {
398+
exception = new McpTransportException("Server Not Found. Status code:" + statusCode);
394399
}
395-
else if (statusCode == NOT_FOUND) {
396-
if (maybeSessionId.isPresent()) {
397-
logger.debug("Session not found for session ID: {}", maybeSessionId.get());
398-
String sessionIdRepresentation = sessionIdOrPlaceholder(maybeSessionId);
399-
exception = new McpTransportSessionNotFoundException(sessionIdRepresentation);
400-
}
401-
else {
402-
exception = new McpTransportException("Server Not Found. Status code:" + statusCode);
403-
}
400+
}
401+
else if (statusCode == BAD_REQUEST) {
402+
// Some implementations return 400 when presented with a session
403+
// id they do not know about, so the session is invalidated.
404+
// https://github.com/modelcontextprotocol/typescript-sdk/issues/389
405+
if (maybeSessionId.isPresent()) {
406+
String sessionIdRepresentation = sessionIdOrPlaceholder(maybeSessionId);
407+
exception = new McpTransportSessionNotFoundException(
408+
"Session not found for session ID: " + sessionIdRepresentation);
404409
}
405-
else if (statusCode == BAD_REQUEST) {
406-
// Unlike a POST, a GET is not treated as a session-not-found
407-
// signal on 400: servers also reject the listening stream
408-
// itself with 400, which must not cost the client its
409-
// session.
410+
else {
410411
exception = new McpTransportException("Bad Request. Status code:" + statusCode);
411412
}
412-
else if (statusCode >= 200 && statusCode < 300) {
413-
String contentType = httpResponse.headers()
414-
.firstValue(HttpHeaders.CONTENT_TYPE)
415-
.orElse("")
416-
.toLowerCase();
417-
if (contentType.contains(TEXT_EVENT_STREAM)) {
418-
logger.debug("SSE connection established successfully");
419-
proceed = true;
420-
}
421-
else {
422-
exception = new McpTransportException(
423-
"Unrecognized server error when connecting to SSE stream, status code: "
424-
+ statusCode);
425-
}
413+
}
414+
else if (statusCode >= 200 && statusCode < 300) {
415+
String contentType = httpResponse.headers()
416+
.firstValue(HttpHeaders.CONTENT_TYPE)
417+
.orElse("")
418+
.toLowerCase();
419+
if (contentType.contains(TEXT_EVENT_STREAM)) {
420+
logger.debug("SSE connection established successfully");
421+
proceed = true;
426422
}
427423
else {
428-
exception = new McpTransportException("Received unrecognized status code: " + statusCode);
424+
exception = new McpTransportException(
425+
"Unrecognized server error when connecting to SSE stream, status code: "
426+
+ statusCode);
429427
}
428+
}
429+
else {
430+
exception = new McpTransportException("Received unrecognized status code: " + statusCode);
431+
}
430432

431-
return proceed ? consumeSseStream(httpResponse.body(), stream, null)
432-
: exception != null
433-
? ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
434-
exception)
435-
: ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
436-
});
433+
return proceed ? consumeSseStream(httpResponse.body(), stream, null)
434+
: exception != null
435+
? ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
436+
exception)
437+
: ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
438+
});
437439
})
438440
.retryWhen(authorizationErrorRetrySpec())
439441
.flatMap(jsonrpcMessage -> requestHandler.apply(Mono.just(jsonrpcMessage)))
@@ -547,102 +549,102 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
547549
var transportContext = ctx.getOrDefault(McpTransportContext.KEY, McpTransportContext.EMPTY);
548550
return Mono
549551
.from(this.httpRequestCustomizer.customize(builder, "POST", uri, jsonBody, transportContext));
550-
})
551-
.flatMapMany(requestBuilder -> Mono
552-
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
553-
HttpResponse.BodyHandlers.ofPublisher()))
554-
.flatMapMany(httpResponse -> {
555-
int statusCode = httpResponse.statusCode();
556-
Optional<String> maybeSessionId = transportSession == null ? Optional.empty()
557-
: transportSession.sessionId();
558-
if (statusCode == 401 || statusCode == 403) {
559-
logger.debug("Authorization error in sendMessage with code {}", statusCode);
560-
var request = requestBuilder.build();
561-
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
562-
request.headers());
563-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
564-
new McpHttpClientTransportAuthorizationException(
565-
"Authorization error when sending message", requestSnapshot,
566-
toResponseInfo(httpResponse)));
567-
}
552+
}).flatMapMany(requestBuilder -> {
553+
var request = requestBuilder.build();
554+
// Classify the response against the session id that this very request
555+
// carried, rather than the one currently held by the session, which
556+
// can be established concurrently.
557+
Optional<String> maybeSessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);
568558

569-
if (transportSession
570-
.markInitialized(httpResponse.headers().firstValue("mcp-session-id").orElse(null))) {
571-
reconnect(null).contextWrite(deliveredSink.contextView()).subscribe();
572-
}
559+
return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMapMany(httpResponse -> {
560+
int statusCode = httpResponse.statusCode();
561+
if (statusCode == 401 || statusCode == 403) {
562+
logger.debug("Authorization error in sendMessage with code {}", statusCode);
563+
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
564+
request.headers());
565+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
566+
new McpHttpClientTransportAuthorizationException(
567+
"Authorization error when sending message", requestSnapshot,
568+
toResponseInfo(httpResponse)));
569+
}
573570

574-
String sessionRepresentation = sessionIdOrPlaceholder(maybeSessionId);
575-
576-
if (statusCode >= 200 && statusCode < 300) {
577-
String contentType = httpResponse.headers()
578-
.firstValue(HttpHeaders.CONTENT_TYPE)
579-
.orElse("")
580-
.toLowerCase();
581-
String contentLength = httpResponse.headers()
582-
.firstValue(HttpHeaders.CONTENT_LENGTH)
583-
.orElse(null);
584-
585-
if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
586-
logger.debug("No body returned for POST in session {}", sessionRepresentation);
587-
markDelivered.run();
588-
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
589-
}
590-
else if (contentType.contains(TEXT_EVENT_STREAM)) {
591-
return consumeSseStream(httpResponse.body(), null, markDelivered);
592-
}
593-
else if (contentType.contains(APPLICATION_JSON)) {
594-
return ResponseBodyHandlers
595-
.decodeAggregateResponse(httpResponse.body(), this.maxResponseSize)
596-
.flatMapMany(data -> {
597-
markDelivered.run();
598-
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
599-
logger.warn("Notification: {} received non-compliant response: {}",
600-
sentMessage, Utils.hasText(data) ? data : "[empty]");
601-
return Flux.empty();
602-
}
603-
try {
604-
return Flux.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
605-
}
606-
catch (IOException e) {
607-
return Flux.<McpSchema.JSONRPCMessage>error(new McpTransportException(
608-
"Error deserializing JSON-RPC message", e));
609-
}
610-
});
611-
}
612-
613-
logger.warn("Unknown media type {} returned for POST in session {}", contentType,
614-
sessionRepresentation);
615-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
616-
new McpTransportException("Unknown media type returned: " + contentType));
571+
if (transportSession
572+
.markInitialized(httpResponse.headers().firstValue("mcp-session-id").orElse(null))) {
573+
reconnect(null).contextWrite(deliveredSink.contextView()).subscribe();
574+
}
575+
576+
String sessionRepresentation = sessionIdOrPlaceholder(maybeSessionId);
577+
578+
if (statusCode >= 200 && statusCode < 300) {
579+
String contentType = httpResponse.headers()
580+
.firstValue(HttpHeaders.CONTENT_TYPE)
581+
.orElse("")
582+
.toLowerCase();
583+
String contentLength = httpResponse.headers()
584+
.firstValue(HttpHeaders.CONTENT_LENGTH)
585+
.orElse(null);
586+
587+
if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
588+
logger.debug("No body returned for POST in session {}", sessionRepresentation);
589+
markDelivered.run();
590+
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
617591
}
618-
else if (statusCode == NOT_FOUND) {
619-
if (maybeSessionId.isPresent()) {
620-
logger.debug("Session not found for session ID: {}", sessionRepresentation);
621-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
622-
new McpTransportSessionNotFoundException(
623-
"Session not found for session ID: " + sessionRepresentation));
624-
}
625-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
626-
new McpTransportException("Server Not Found. Status code:" + statusCode));
592+
else if (contentType.contains(TEXT_EVENT_STREAM)) {
593+
return consumeSseStream(httpResponse.body(), null, markDelivered);
594+
}
595+
else if (contentType.contains(APPLICATION_JSON)) {
596+
return ResponseBodyHandlers
597+
.decodeAggregateResponse(httpResponse.body(), this.maxResponseSize)
598+
.flatMapMany(data -> {
599+
markDelivered.run();
600+
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
601+
logger.warn("Notification: {} received non-compliant response: {}", sentMessage,
602+
Utils.hasText(data) ? data : "[empty]");
603+
return Flux.empty();
604+
}
605+
try {
606+
return Flux.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
607+
}
608+
catch (IOException e) {
609+
return Flux.<McpSchema.JSONRPCMessage>error(
610+
new McpTransportException("Error deserializing JSON-RPC message", e));
611+
}
612+
});
627613
}
628-
else if (statusCode == BAD_REQUEST) {
629-
if (maybeSessionId.isPresent()) {
630-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
631-
new McpTransportSessionNotFoundException(
632-
"Session not found for session ID: " + sessionRepresentation));
633-
}
614+
615+
logger.warn("Unknown media type {} returned for POST in session {}", contentType,
616+
sessionRepresentation);
617+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
618+
new McpTransportException("Unknown media type returned: " + contentType));
619+
}
620+
else if (statusCode == NOT_FOUND) {
621+
if (maybeSessionId.isPresent()) {
622+
logger.debug("Session not found for session ID: {}", sessionRepresentation);
634623
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
635-
new McpTransportException("Bad Request. Status code:" + statusCode));
624+
new McpTransportSessionNotFoundException(
625+
"Session not found for session ID: " + sessionRepresentation));
636626
}
637-
else if (statusCode >= 400 && statusCode < 500) {
627+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
628+
new McpTransportException("Server Not Found. Status code:" + statusCode));
629+
}
630+
else if (statusCode == BAD_REQUEST) {
631+
if (maybeSessionId.isPresent()) {
638632
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
639-
new McpTransportException("Invalid request. Status code: " + statusCode));
633+
new McpTransportSessionNotFoundException(
634+
"Session not found for session ID: " + sessionRepresentation));
640635
}
641-
642636
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
643-
new McpTransportException("Failed to send message, status code: " + statusCode));
644-
})
645-
.onErrorMap(CompletionException.class, Throwable::getCause))
637+
new McpTransportException("Bad Request. Status code:" + statusCode));
638+
}
639+
else if (statusCode >= 400 && statusCode < 500) {
640+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
641+
new McpTransportException("Invalid request. Status code: " + statusCode));
642+
}
643+
644+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
645+
new McpTransportException("Failed to send message, status code: " + statusCode));
646+
});
647+
})
646648
.retryWhen(authorizationErrorRetrySpec())
647649
.flatMap(jsonRpcMessage -> requestHandler.apply(Mono.just(jsonRpcMessage)))
648650
.onErrorMap(CompletionException.class, t -> t.getCause())

0 commit comments

Comments
 (0)