Skip to content

Commit f8271c2

Browse files
committed
Bound HttpClient inbound reads in the reactive chain
- This discards all BodySubscribers to align with a6c88fccf50277ddf705765154baae4caa906e15, which introduced a publisher-based architecture instead of using subscribers. - Bounding the reads now happens within the reactive chains, e.g. in Utf8LineDecoder, decodeAggregateResponse and drain* methods. Signed-off-by: Daniel Garnier-Moiroux <git@garnier.wf>
1 parent 578b6bb commit f8271c2

9 files changed

Lines changed: 295 additions & 658 deletions

File tree

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

Lines changed: 25 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -390,23 +390,23 @@ public Mono<Void> connect(Function<Mono<JSONRPCMessage>, Mono<JSONRPCMessage>> h
390390
}).flatMap(requestBuilder -> Mono.create(sink -> {
391391
Disposable connection = Mono
392392
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
393-
ResponseSubscribers.boundedPublisherBodyHandler(this.maxResponseSize)))
393+
HttpResponse.BodyHandlers.ofPublisher()))
394394
.flatMapMany(response -> {
395395
if (isClosing) {
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());
399+
return ResponseSubscribers.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());
405+
Flux<String> lines = ResponseSubscribers.decodeLines(response.body(), this.maxResponseSize);
406406
return ResponseSubscribers.decodeSseResponse(lines, this.maxResponseSize);
407407
}
408408
else {
409-
return ResponseSubscribers.drainThenError(response.body(),
409+
return ResponseSubscribers.drainThenError(response.body(), this.maxResponseSize,
410410
new RuntimeException("Failed to connect to SSE stream: " + statusCode));
411411
}
412412
})
@@ -491,17 +491,7 @@ public Mono<Void> sendMessage(JSONRPCMessage message) {
491491
}
492492

493493
return this.serializeMessage(message)
494-
.flatMap(body -> sendHttpPost(messageEndpointUri, body).handle((response, sink) -> {
495-
if (response.statusCode() != 200 && response.statusCode() != 201 && response.statusCode() != 202
496-
&& response.statusCode() != 206) {
497-
sink.error(new RuntimeException("Sending message failed with a non-OK HTTP code: "
498-
+ response.statusCode() + " - " + response.body()));
499-
}
500-
else {
501-
sink.next(response);
502-
sink.complete();
503-
}
504-
}))
494+
.flatMap(body -> sendHttpPost(messageEndpointUri, body))
505495
.doOnError(error -> {
506496
if (!isClosing) {
507497
logger.error("Error sending message: {}", error.getMessage());
@@ -522,7 +512,16 @@ private Mono<String> serializeMessage(final JSONRPCMessage message) {
522512
});
523513
}
524514

525-
private Mono<HttpResponse<String>> sendHttpPost(final String endpoint, final String body) {
515+
/**
516+
* POSTs {@code body} to {@code endpoint} and consumes the response, failing if the
517+
* server did not accept the message.
518+
*
519+
* <p>
520+
* The response body is streamed rather than aggregated: it is only read as text when
521+
* a non-OK status makes it part of the failure message, and discarded otherwise.
522+
* Either way it has to be consumed, or the connection is never released.
523+
*/
524+
private Mono<Void> sendHttpPost(final String endpoint, final String body) {
526525
final URI requestUri = Utils.resolveUri(baseUri, endpoint);
527526
return Mono.deferContextual(ctx -> {
528527
var builder = this.requestBuilder.copy()
@@ -534,8 +533,16 @@ private Mono<HttpResponse<String>> sendHttpPost(final String endpoint, final Str
534533
return Mono.from(this.httpRequestCustomizer.customize(builder, "POST", requestUri, body, transportContext));
535534
}).flatMap(customizedBuilder -> {
536535
var request = customizedBuilder.build();
537-
return Mono.fromFuture(this.httpClient.sendAsync(request,
538-
ResponseSubscribers.boundedStringBodyHandler(this.maxResponseSize)));
536+
return Mono.fromFuture(this.httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher()))
537+
.flatMap(response -> {
538+
int statusCode = response.statusCode();
539+
if (statusCode == 200 || statusCode == 201 || statusCode == 202 || statusCode == 206) {
540+
return ResponseSubscribers.drain(response.body(), this.maxResponseSize).then();
541+
}
542+
return ResponseSubscribers.decodeAggregateResponse(response.body(), this.maxResponseSize)
543+
.flatMap(text -> Mono.error(new RuntimeException(
544+
"Sending message failed with a non-OK HTTP code: " + statusCode + " - " + text)));
545+
});
539546
});
540547
}
541548

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

Lines changed: 22 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -241,8 +241,12 @@ 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(() -> this.httpClient.sendAsync(requestBuilder.build(),
245-
ResponseSubscribers.boundedStringBodyHandler(this.maxResponseSize))))
244+
.flatMap(requestBuilder -> Mono.fromFuture(
245+
() -> this.httpClient.sendAsync(requestBuilder.build(), HttpResponse.BodyHandlers.ofPublisher()))
246+
// The response is not inspected, but the body still has to be consumed
247+
// to release the connection.
248+
.flatMapMany(response -> ResponseSubscribers.drain(response.body(), this.maxResponseSize))
249+
.then())
246250
.then();
247251
}
248252

@@ -282,7 +286,7 @@ public Mono<Void> closeGracefully() {
282286
private Flux<McpSchema.JSONRPCMessage> consumeSseStream(
283287
java.util.concurrent.Flow.Publisher<List<java.nio.ByteBuffer>> body,
284288
McpTransportStream<Disposable> existingStream, Runnable onFirstMessage) {
285-
Flux<String> lines = ResponseSubscribers.decodeLines(body);
289+
Flux<String> lines = ResponseSubscribers.decodeLines(body, this.maxResponseSize);
286290
return ResponseSubscribers.decodeSseResponse(lines, this.maxResponseSize).flatMap(sseEvent -> {
287291
if (!isMessageEvent(sseEvent.event())) {
288292
logger.debug("Received SSE event with type: {}", sseEvent);
@@ -372,8 +376,7 @@ private Mono<Disposable> reconnect(McpTransportStream<Disposable> stream) {
372376
Optional<String> maybeSessionId = request.headers().firstValue(HttpHeaders.MCP_SESSION_ID);
373377

374378
return Mono
375-
.fromFuture(() -> this.httpClient.sendAsync(request,
376-
ResponseSubscribers.boundedPublisherBodyHandler(this.maxResponseSize)))
379+
.fromFuture(() -> this.httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher()))
377380
.flatMapMany(httpResponse -> {
378381
int statusCode = httpResponse.statusCode();
379382
Exception exception = null;
@@ -430,8 +433,10 @@ else if (statusCode >= 200 && statusCode < 300) {
430433
}
431434

432435
return proceed ? consumeSseStream(httpResponse.body(), stream, null)
433-
: exception != null ? ResponseSubscribers.drainThenError(httpResponse.body(), exception)
434-
: ResponseSubscribers.drain(httpResponse.body());
436+
: exception != null
437+
? ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
438+
exception)
439+
: ResponseSubscribers.drain(httpResponse.body(), this.maxResponseSize);
435440
});
436441
})
437442
.retryWhen(authorizationErrorRetrySpec())
@@ -540,7 +545,7 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
540545
})
541546
.flatMapMany(requestBuilder -> Mono
542547
.fromFuture(() -> this.httpClient.sendAsync(requestBuilder.build(),
543-
ResponseSubscribers.boundedPublisherBodyHandler(this.maxResponseSize)))
548+
HttpResponse.BodyHandlers.ofPublisher()))
544549
.flatMapMany(httpResponse -> {
545550
int statusCode = httpResponse.statusCode();
546551
Optional<String> maybeSessionId = transportSession == null ? Optional.empty()
@@ -550,7 +555,7 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
550555
var request = requestBuilder.build();
551556
var requestSnapshot = new HttpRequestSnapshot(request.uri(), request.method(),
552557
request.headers());
553-
return ResponseSubscribers.drainThenError(httpResponse.body(),
558+
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
554559
new McpHttpClientTransportAuthorizationException(
555560
"Authorization error when sending message", requestSnapshot,
556561
toResponseInfo(httpResponse)));
@@ -575,7 +580,7 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
575580
if (contentType.isBlank() || "0".equals(contentLength) || statusCode == 202) {
576581
logger.debug("No body returned for POST in session {}", sessionRepresentation);
577582
deliveredSink.success();
578-
return ResponseSubscribers.drain(httpResponse.body());
583+
return ResponseSubscribers.drain(httpResponse.body(), this.maxResponseSize);
579584
}
580585
else if (contentType.contains(TEXT_EVENT_STREAM)) {
581586
AtomicBoolean delivered = new AtomicBoolean();
@@ -607,34 +612,34 @@ else if (contentType.contains(APPLICATION_JSON)) {
607612

608613
logger.warn("Unknown media type {} returned for POST in session {}", contentType,
609614
sessionRepresentation);
610-
return ResponseSubscribers.drainThenError(httpResponse.body(),
615+
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
611616
new RuntimeException("Unknown media type returned: " + contentType));
612617
}
613618
else if (statusCode == NOT_FOUND) {
614619
if (maybeSessionId.isPresent()) {
615620
logger.debug("Session not found for session ID: {}", sessionRepresentation);
616-
return ResponseSubscribers.drainThenError(httpResponse.body(),
621+
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
617622
new McpTransportSessionNotFoundException(
618623
"Session not found for session ID: " + sessionRepresentation));
619624
}
620-
return ResponseSubscribers.drainThenError(httpResponse.body(),
625+
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
621626
new McpTransportException("Server Not Found. Status code:" + statusCode));
622627
}
623628
else if (statusCode == BAD_REQUEST) {
624629
if (maybeSessionId.isPresent()) {
625-
return ResponseSubscribers.drainThenError(httpResponse.body(),
630+
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
626631
new McpTransportSessionNotFoundException(
627632
"Session not found for session ID: " + sessionRepresentation));
628633
}
629-
return ResponseSubscribers.drainThenError(httpResponse.body(),
634+
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
630635
new McpTransportException("Bad Request. Status code:" + statusCode));
631636
}
632637
else if (statusCode >= 400 && statusCode < 500) {
633-
return ResponseSubscribers.drainThenError(httpResponse.body(),
638+
return ResponseSubscribers.drainThenError(httpResponse.body(), this.maxResponseSize,
634639
new McpTransportException("Invalid request. Status code: " + statusCode));
635640
}
636641

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

0 commit comments

Comments
 (0)