Skip to content

Commit 81a3dc9

Browse files
committed
Ensure errors are priorities over body limits, rename ResponseSubscribers, refactor test code
Signed-off-by: Dariusz Jędrzejczyk <dariusz.jedrzejczyk@broadcom.com>
1 parent f8271c2 commit 81a3dc9

7 files changed

Lines changed: 32 additions & 32 deletions

File tree

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

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -396,17 +396,17 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
396396
// The body is handed over as a publisher and nothing is read off
397397
// the wire until it is subscribed, so it has to be drained even
398398
// when its content is of no further interest.
399-
return ResponseSubscribers.drain(response.body(), this.maxResponseSize);
399+
return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize);
400400
}
401401

402402
int statusCode = response.statusCode();
403403

404404
if (statusCode >= 200 && statusCode < 300) {
405-
Flux<String> lines = ResponseSubscribers.decodeLines(response.body(), this.maxResponseSize);
406-
return ResponseSubscribers.decodeSseResponse(lines, this.maxResponseSize);
405+
Flux<String> lines = ResponseBodyHandlers.decodeLines(response.body(), this.maxResponseSize);
406+
return ResponseBodyHandlers.decodeSseResponse(lines, this.maxResponseSize);
407407
}
408408
else {
409-
return ResponseSubscribers.drainThenError(response.body(), this.maxResponseSize,
409+
return ResponseBodyHandlers.drainThenError(response.body(), this.maxResponseSize,
410410
new RuntimeException("Failed to connect to SSE stream: " + statusCode));
411411
}
412412
})
@@ -537,10 +537,10 @@ private Mono<Void> sendHttpPost(final String endpoint, final String body) {
537537
.flatMap(response -> {
538538
int statusCode = response.statusCode();
539539
if (statusCode == 200 || statusCode == 201 || statusCode == 202 || statusCode == 206) {
540-
return ResponseSubscribers.drain(response.body(), this.maxResponseSize).then();
540+
return ResponseBodyHandlers.drain(response.body(), this.maxResponseSize).then();
541541
}
542-
return ResponseSubscribers.decodeAggregateResponse(response.body(), this.maxResponseSize)
543-
.flatMap(text -> Mono.error(new RuntimeException(
542+
return ResponseBodyHandlers.decodeAggregateResponse(response.body(), this.maxResponseSize)
543+
.flatMap(text -> Mono.error(new RuntimeException(
544544
"Sending message failed with a non-OK HTTP code: " + statusCode + " - " + text)));
545545
});
546546
});

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

Lines changed: 15 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -245,7 +245,7 @@ private Publisher<Void> createDelete(String sessionId) {
245245
() -> this.httpClient.sendAsync(requestBuilder.build(), HttpResponse.BodyHandlers.ofPublisher()))
246246
// The response is not inspected, but the body still has to be consumed
247247
// to release the connection.
248-
.flatMapMany(response -> ResponseSubscribers.drain(response.body(), this.maxResponseSize))
248+
.flatMapMany(response -> ResponseBodyHandlers.drain(response.body(), this.maxResponseSize))
249249
.then())
250250
.then();
251251
}
@@ -286,8 +286,8 @@ public Mono<Void> closeGracefully() {
286286
private Flux<McpSchema.JSONRPCMessage> consumeSseStream(
287287
java.util.concurrent.Flow.Publisher<List<java.nio.ByteBuffer>> body,
288288
McpTransportStream<Disposable> existingStream, Runnable onFirstMessage) {
289-
Flux<String> lines = ResponseSubscribers.decodeLines(body, this.maxResponseSize);
290-
return ResponseSubscribers.decodeSseResponse(lines, this.maxResponseSize).flatMap(sseEvent -> {
289+
Flux<String> lines = ResponseBodyHandlers.decodeLines(body, this.maxResponseSize);
290+
return ResponseBodyHandlers.decodeSseResponse(lines, this.maxResponseSize).flatMap(sseEvent -> {
291291
if (!isMessageEvent(sseEvent.event())) {
292292
logger.debug("Received SSE event with type: {}", sseEvent);
293293
if (onFirstMessage != null) {
@@ -434,9 +434,9 @@ else if (statusCode >= 200 && statusCode < 300) {
434434

435435
return proceed ? consumeSseStream(httpResponse.body(), stream, null)
436436
: exception != null
437-
? ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
437+
? ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
438438
exception)
439-
: ResponseSubscribers.drain(httpResponse.body(), this.maxResponseSize);
439+
: ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
440440
});
441441
})
442442
.retryWhen(authorizationErrorRetrySpec())
@@ -555,7 +555,7 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
555555
var request = requestBuilder.build();
556556
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
557557
request.headers());
558-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
558+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
559559
new McpHttpClientTransportAuthorizationException(
560560
"Authorization error when sending message", requestSnapshot,
561561
toResponseInfo(httpResponse)));
@@ -580,7 +580,7 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
580580
if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
581581
logger.debug("No body returned for POST in session {}", sessionRepresentation);
582582
deliveredSink.success();
583-
return ResponseSubscribers.drain(httpResponse.body(), this.maxResponseSize);
583+
return ResponseBodyHandlers.drain(httpResponse.body(), this.maxResponseSize);
584584
}
585585
else if (contentType.contains(TEXT_EVENT_STREAM)) {
586586
AtomicBoolean delivered = new AtomicBoolean();
@@ -591,7 +591,7 @@ else if (contentType.contains(TEXT_EVENT_STREAM)) {
591591
});
592592
}
593593
else if (contentType.contains(APPLICATION_JSON)) {
594-
return ResponseSubscribers
594+
return ResponseBodyHandlers
595595
.decodeAggregateResponse(httpResponse.body(), this.maxResponseSize)
596596
.flatMapMany(data -> {
597597
deliveredSink.success();
@@ -612,34 +612,34 @@ else if (contentType.contains(APPLICATION_JSON)) {
612612

613613
logger.warn("Unknown media type {} returned for POST in session {}", contentType,
614614
sessionRepresentation);
615-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
615+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
616616
new RuntimeException("Unknown media type returned: " + contentType));
617617
}
618618
else if (statusCode == NOT_FOUND) {
619619
if (maybeSessionId.isPresent()) {
620620
logger.debug("Session not found for session ID: {}", sessionRepresentation);
621-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
621+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
622622
new McpTransportSessionNotFoundException(
623623
"Session not found for session ID: " + sessionRepresentation));
624624
}
625-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
625+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
626626
new McpTransportException("Server Not Found. Status code:" + statusCode));
627627
}
628628
else if (statusCode == BAD_REQUEST) {
629629
if (maybeSessionId.isPresent()) {
630-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
630+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
631631
new McpTransportSessionNotFoundException(
632632
"Session not found for session ID: " + sessionRepresentation));
633633
}
634-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
634+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
635635
new McpTransportException("Bad Request. Status code:" + statusCode));
636636
}
637637
else if (statusCode >= 400 && statusCode < 500) {
638-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
638+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
639639
new McpTransportException("Invalid request. Status code: " + statusCode));
640640
}
641641

642-
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
642+
return ResponseBodyHandlers.drainThenError(httpResponse.body(), this.maxResponseSize,
643643
new RuntimeException("Failed to send message, status code: " + statusCode));
644644
})
645645
.onErrorMap(CompletionException.class, Throwable::getCause))

mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseSubscribers.java renamed to mcp-core/src/main/java/io/modelcontextprotocol/client/transport/ResponseBodyHandlers.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -36,7 +36,7 @@
3636
* @author Dariusz Jędrzejczyk
3737
* @author Daniel Garnier-Moiroux
3838
*/
39-
class ResponseSubscribers {
39+
class ResponseBodyHandlers {
4040

4141
/**
4242
* Bytes of SSE field framing a single line may carry on top of the message payload:
@@ -137,7 +137,7 @@ static Mono<String> decodeAggregateResponse(Publisher<List<ByteBuffer>> publishe
137137
* @param error the error to propagate once the body has been discarded
138138
*/
139139
static <T> Flux<T> drainThenError(Publisher<List<ByteBuffer>> body, int maxSize, Throwable error) {
140-
return boundTotalBytes(body, maxSize).thenMany(Flux.error(error));
140+
return boundTotalBytes(body, maxSize).onErrorComplete().thenMany(Mono.error(error));
141141
}
142142

143143
/**

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

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,7 @@
1717
import reactor.adapter.JdkFlowAdapter;
1818
import reactor.core.publisher.Flux;
1919

20-
import io.modelcontextprotocol.client.transport.ResponseSubscribers.SseEvent;
20+
import io.modelcontextprotocol.client.transport.ResponseBodyHandlers.SseEvent;
2121

2222
import static org.assertj.core.api.Assertions.assertThat;
2323

@@ -48,7 +48,7 @@
4848
*
4949
* <p>
5050
* The middle column read each chunk incrementally, which is ~25x quicker than what 2.0.0
51-
* shipped, but {@link ResponseSubscribers.Utf8LineDecoder} still searched its buffered
51+
* shipped, but {@link ResponseBodyHandlers.Utf8LineDecoder} still searched its buffered
5252
* characters for a line terminator from the start of the buffer on every chunk, so eight
5353
* times the payload cost ~45x the time. Resuming that search where the previous one ended
5454
* gives the third column, which scales with the payload rather than with its square and
@@ -173,8 +173,8 @@ private static long timeDecode(byte[] body, int expectedEvents) {
173173
private static List<SseEvent> decode(byte[] body) {
174174
Flow.Publisher<List<ByteBuffer>> publisher = JdkFlowAdapter
175175
.publisherToFlowPublisher(Flux.fromIterable(chunk(body)));
176-
Flux<String> lines = ResponseSubscribers.decodeLines(publisher, Integer.MAX_VALUE);
177-
return ResponseSubscribers.decodeSseResponse(lines, MAX_SIZE).collectList().block();
176+
Flux<String> lines = ResponseBodyHandlers.decodeLines(publisher, Integer.MAX_VALUE);
177+
return ResponseBodyHandlers.decodeSseResponse(lines, MAX_SIZE).collectList().block();
178178
}
179179

180180
private static List<List<ByteBuffer>> chunk(byte[] body) {

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,8 @@
88

99
import org.junit.jupiter.api.Test;
1010

11-
import io.modelcontextprotocol.client.transport.ResponseSubscribers.SseEvent;
12-
import io.modelcontextprotocol.client.transport.ResponseSubscribers.SseEventParser;
11+
import io.modelcontextprotocol.client.transport.ResponseBodyHandlers.SseEvent;
12+
import io.modelcontextprotocol.client.transport.ResponseBodyHandlers.SseEventParser;
1313

1414
import static org.assertj.core.api.Assertions.assertThat;
1515

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88
import java.nio.charset.StandardCharsets;
99
import java.util.List;
1010

11-
import io.modelcontextprotocol.client.transport.ResponseSubscribers.Utf8LineDecoder;
11+
import io.modelcontextprotocol.client.transport.ResponseBodyHandlers.Utf8LineDecoder;
1212
import io.modelcontextprotocol.spec.McpTransportException;
1313
import org.junit.jupiter.api.Test;
1414

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@
99
import java.nio.charset.StandardCharsets;
1010
import java.util.List;
1111

12-
import io.modelcontextprotocol.client.transport.ResponseSubscribers.Utf8LineDecoder;
12+
import io.modelcontextprotocol.client.transport.ResponseBodyHandlers.Utf8LineDecoder;
1313
import org.junit.jupiter.api.Test;
1414

1515
import static org.assertj.core.api.Assertions.assertThat;

0 commit comments

Comments
 (0)