Skip to content

Commit 46511cd

Browse files
committed
Surface response body in case of errors
Signed-off-by: Dariusz Jędrzejczyk <dariusz.jedrzejczyk@broadcom.com>
1 parent 9ac6286 commit 46511cd

6 files changed

Lines changed: 127 additions & 129 deletions

File tree

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -413,8 +413,8 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
413413
return ResponseBodyHandlers.decodeSseResponse(lines, this.maxResponseSize);
414414
}
415415
else {
416-
return ResponseBodyHandlers.drainThenError(response.body(), this.maxResponseSize,
417-
new McpTransportException("Failed to connect to SSE stream: " + statusCode));
416+
return ResponseBodyHandlers.readThenError(response.body(), this.maxResponseSize,
417+
"Failed to connect to SSE stream: " + statusCode);
418418
}
419419
})
420420
.flatMap(sseEvent -> {

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

Lines changed: 79 additions & 126 deletions
Original file line numberDiff line numberDiff line change
@@ -314,7 +314,7 @@ private Flux<McpSchema.JSONRPCMessage> consumeSseStream(Flow.Publisher<List<Byte
314314
}
315315
catch (IOException e) {
316316
return Flux.<McpSchema.JSONRPCMessage>error(
317-
new McpTransportException("Error parsing JSON-RPC message", e));
317+
new McpTransportException("Error parsing JSON-RPC message: " + data, e));
318318
}
319319
});
320320
}
@@ -368,73 +368,34 @@ private Mono<Disposable> reconnect(McpTransportStream<Disposable> stream) {
368368
return Mono.from(this.httpRequestCustomizer.customize(builder, "GET", uri, null, transportContext));
369369
}).flatMapMany(requestBuilder -> {
370370
var request = requestBuilder.build();
371-
// Classify the response against the session id that this very request
372-
// carried, rather than the one currently held by the session, which
373-
// can be established concurrently.
374-
Optional<String> maybeSessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);
375-
376371
return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMapMany(httpResponse -> {
377372
int statusCode = httpResponse.statusCode();
378-
Exception exception = null;
379-
boolean proceed = false;
380373
if (statusCode == 401 || statusCode == 403) {
381374
logger.debug("Authorization error in reconnect with code {}", statusCode);
382375
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
383376
request.headers());
384-
exception = new McpHttpClientTransportAuthorizationException(
385-
"Authorization error connecting to SSE stream", requestSnapshot,
386-
toResponseInfo(httpResponse));
377+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
378+
new McpHttpClientTransportAuthorizationException(
379+
"Authorization error connecting to SSE stream", requestSnapshot,
380+
toResponseInfo(httpResponse)));
387381
}
388-
else if (statusCode == METHOD_NOT_ALLOWED) {
382+
if (statusCode == METHOD_NOT_ALLOWED) {
389383
logger.debug("The server does not support SSE streams, using request-response mode.");
384+
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
390385
}
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);
396-
}
397-
else {
398-
exception = new McpTransportException("Server Not Found. Status code:" + statusCode);
399-
}
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);
409-
}
410-
else {
411-
exception = new McpTransportException("Bad Request. Status code:" + statusCode);
412-
}
386+
if (statusCode < 200 || statusCode >= 300) {
387+
return statusError(request, httpResponse);
413388
}
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;
422-
}
423-
else {
424-
exception = new McpTransportException(
425-
"Unrecognized server error when connecting to SSE stream, status code: "
426-
+ statusCode);
427-
}
389+
String contentType = httpResponse.headers()
390+
.firstValue(HttpHeaders.CONTENT_TYPE)
391+
.orElse("")
392+
.toLowerCase();
393+
if (!contentType.contains(TEXT_EVENT_STREAM)) {
394+
return ResponseBodyHandlers.readThenError(httpResponse.body(), this.maxResponseSize,
395+
"Unrecognized server error when connecting to SSE stream, status code: " + statusCode);
428396
}
429-
else {
430-
exception = new McpTransportException("Received unrecognized status code: " + statusCode);
431-
}
432-
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);
397+
logger.debug("SSE connection established successfully");
398+
return consumeSseStream(httpResponse.body(), stream, null);
438399
});
439400
})
440401
.retryWhen(authorizationErrorRetrySpec())
@@ -551,11 +512,6 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
551512
.from(this.httpRequestCustomizer.customize(builder, "POST", uri, jsonBody, transportContext));
552513
}).flatMapMany(requestBuilder -> {
553514
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);
558-
559515
return ResponseBodyHandlers.sendAsync(this.httpClient, request).flatMapMany(httpResponse -> {
560516
int statusCode = httpResponse.statusCode();
561517
if (statusCode == 401 || statusCode == 403) {
@@ -573,76 +529,49 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
573529
reconnect(null).contextWrite(deliveredSink.contextView()).subscribe();
574530
}
575531

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);
591-
}
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-
});
613-
}
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));
532+
if (statusCode < 200 || statusCode >= 300) {
533+
return statusError(request, httpResponse);
619534
}
620-
else if (statusCode == NOT_FOUND) {
621-
if (maybeSessionId.isPresent()) {
622-
logger.debug("Session not found for session ID: {}", sessionRepresentation);
623-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
624-
new McpTransportSessionNotFoundException(
625-
"Session not found for session ID: " + sessionRepresentation));
626-
}
627-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
628-
new McpTransportException("Server Not Found. Status code:" + statusCode));
535+
536+
String sessionRepresentation = sessionIdOrPlaceholder(
537+
request.headers().firstValue(HttpHeaders.MCP_SESSION_ID));
538+
String contentType = httpResponse.headers()
539+
.firstValue(HttpHeaders.CONTENT_TYPE)
540+
.orElse("")
541+
.toLowerCase();
542+
String contentLength = httpResponse.headers().firstValue(HttpHeaders.CONTENT_LENGTH).orElse(null);
543+
544+
if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
545+
logger.debug("No body returned for POST in session {}", sessionRepresentation);
546+
markDelivered.run();
547+
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
629548
}
630-
else if (statusCode == BAD_REQUEST) {
631-
if (maybeSessionId.isPresent()) {
632-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
633-
new McpTransportSessionNotFoundException(
634-
"Session not found for session ID: " + sessionRepresentation));
635-
}
636-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
637-
new McpTransportException("Bad Request. Status code:" + statusCode));
549+
else if (contentType.contains(TEXT_EVENT_STREAM)) {
550+
return consumeSseStream(httpResponse.body(), null, markDelivered);
638551
}
639-
else if (statusCode >= 400 && statusCode < 500) {
640-
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
641-
new McpTransportException("Invalid request. Status code: " + statusCode));
552+
else if (contentType.contains(APPLICATION_JSON)) {
553+
return ResponseBodyHandlers.decodeAggregateResponse(httpResponse.body(), this.maxResponseSize)
554+
.flatMapMany(data -> {
555+
markDelivered.run();
556+
if (sentMessage instanceof McpSchema.JSONRPCNotification) {
557+
logger.warn("Notification: {} received non-compliant response: {}", sentMessage,
558+
Utils.hasText(data) ? data : "[empty]");
559+
return Flux.empty();
560+
}
561+
try {
562+
return Flux.just(McpSchema.deserializeJsonRpcMessage(jsonMapper, data));
563+
}
564+
catch (IOException e) {
565+
return Flux.<McpSchema.JSONRPCMessage>error(new McpTransportException(
566+
"Error deserializing JSON-RPC message: " + data, e));
567+
}
568+
});
642569
}
643570

571+
logger.warn("Unknown media type {} returned for POST in session {}", contentType,
572+
sessionRepresentation);
644573
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
645-
new McpTransportException("Failed to send message, status code: " + statusCode));
574+
new McpTransportException("Unknown media type returned: " + contentType));
646575
});
647576
})
648577
.retryWhen(authorizationErrorRetrySpec())
@@ -681,6 +610,30 @@ else if (statusCode >= 400 && statusCode < 500) {
681610

682611
}
683612

613+
/**
614+
* Fails the exchange over a response with an error status. A session id the server
615+
* does not recognise invalidates the session; any other failure carries the response
616+
* body, which is what the server said about it.
617+
*/
618+
private <T> Flux<T> statusError(HttpRequest request, HttpResponse<Flow.Publisher<List<ByteBuffer>>> response) {
619+
int statusCode = response.statusCode();
620+
// Classify the response against the session id that this very request carried,
621+
// rather than the one currently held by the session, which can be established
622+
// concurrently. Some implementations return 400 rather than 404 for a session id
623+
// they do not know about.
624+
// https://github.com/modelcontextprotocol/typescript-sdk/issues/389
625+
Optional<String> sessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);
626+
if ((statusCode == NOT_FOUND || statusCode == BAD_REQUEST) && sessionId.isPresent()) {
627+
logger.debug("Session not found for session ID: {}", sessionId.get());
628+
return ResponseBodyHandlers.drainThenError(response.body(), this.maxResponseSize,
629+
new McpTransportSessionNotFoundException(sessionId.get()));
630+
}
631+
String failure = statusCode == NOT_FOUND ? "Server Not Found. Status code:" + statusCode
632+
: statusCode == BAD_REQUEST ? "Bad Request. Status code:" + statusCode
633+
: "Received unexpected status code: " + statusCode;
634+
return ResponseBodyHandlers.readThenError(response.body(), this.maxResponseSize, failure);
635+
}
636+
684637
private static String sessionIdOrPlaceholder(Optional<String> sessionId) {
685638
return sessionId.orElse("[missing_session_id]");
686639
}

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

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
import org.slf4j.LoggerFactory;
2525

2626
import io.modelcontextprotocol.spec.McpTransportException;
27+
import io.modelcontextprotocol.util.Utils;
2728
import reactor.adapter.JdkFlowAdapter;
2829
import reactor.core.publisher.Flux;
2930
import reactor.core.publisher.Mono;
@@ -150,6 +151,19 @@ static <T> Flux<T> drainThenError(Publisher<List<ByteBuffer>> body, int maxSize,
150151
return boundTotalBytes(body, maxSize).onErrorComplete().thenMany(Mono.error(error));
151152
}
152153

154+
/**
155+
* Reads the body as text, then propagates a {@link McpTransportException} carrying
156+
* {@code message} followed by that text, so that what the server said about the
157+
* failure reaches the caller. The body is read under the same bound as any other.
158+
* @param body the response body
159+
* @param maxSize the maximum number of bytes read for the response body
160+
* @param message describes the failure the body explains
161+
*/
162+
static <T> Flux<T> readThenError(Publisher<List<ByteBuffer>> body, int maxSize, String message) {
163+
return decodeAggregateResponse(body, maxSize).flatMapMany(text -> Flux
164+
.error(new McpTransportException(Utils.hasText(text) ? message + ", response body: " + text : message)));
165+
}
166+
153167
/**
154168
* Subscribes to the body publisher to release the underlying connection, discarding
155169
* all bytes, then completes empty.

‎mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientBoundedReadTestSupport.java‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,15 @@ void stopServer() {
7777
*/
7878
protected CompletableFuture<IOException> respondWith(String method, String path, String contentType,
7979
Responder responder) {
80+
return respondWith(method, path, 200, contentType, responder);
81+
}
82+
83+
/**
84+
* Like {@link #respondWith(String, String, String, Responder)}, but answering with
85+
* {@code status} rather than 200.
86+
*/
87+
protected CompletableFuture<IOException> respondWith(String method, String path, int status, String contentType,
88+
Responder responder) {
8089
CompletableFuture<IOException> response = new CompletableFuture<>();
8190
this.server.createContext(path, exchange -> {
8291
try {
@@ -85,7 +94,7 @@ protected CompletableFuture<IOException> respondWith(String method, String path,
8594
return;
8695
}
8796
exchange.getResponseHeaders().set("Content-Type", contentType);
88-
exchange.sendResponseHeaders(200, 0);
97+
exchange.sendResponseHeaders(status, 0);
8998
try (OutputStream body = exchange.getResponseBody()) {
9099
responder.respond(body);
91100
response.complete(null);

‎mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientSseClientTransportBoundedReadTests.java‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,17 @@ void shouldRejectPostResponseExceedingMaxSize() {
8585
assertHungUp(response);
8686
}
8787

88+
@Test
89+
void shouldIncludeConnectErrorResponseBodyInError() {
90+
// What the server says about a failure is the most useful part of it to report.
91+
respondWith("GET", endpoint(), 500, "text/plain",
92+
body -> body.write("upstream unavailable".getBytes(StandardCharsets.UTF_8)));
93+
94+
StepVerifier.create(connect())
95+
.verifyErrorMatches(t -> messageContains(t,
96+
"Failed to connect to SSE stream: 500, response body: upstream unavailable"));
97+
}
98+
8899
private void awaitTeardown() {
89100
try {
90101
this.keepStreamOpen.await(10, TimeUnit.SECONDS);

‎mcp-test/src/test/java/io/modelcontextprotocol/client/transport/HttpClientStreamableHttpTransportBoundedReadTests.java‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,17 @@ void shouldRejectDiscardedResponseExceedingMaxSizeButReportProperError() {
7676
assertHungUp(response);
7777
}
7878

79+
@Test
80+
void shouldIncludeErrorResponseBodyInError() {
81+
// What the server says about a failure is the most useful part of it to report.
82+
respondWith("POST", endpoint(), 404, "text/plain",
83+
body -> body.write("no MCP server here".getBytes(StandardCharsets.UTF_8)));
84+
85+
StepVerifier.create(sendMessage())
86+
.verifyErrorMatches(
87+
t -> messageContains(t, "Server Not Found. Status code:404, response body: no MCP server here"));
88+
}
89+
7990
@Test
8091
void shouldAcceptEventOfExactlyMaxSize() {
8192
// The bound is inclusive and the SSE framing around the payload is given its own

0 commit comments

Comments
 (0)