Skip to content

Commit eb4009f

Browse files
committed
Handle request cancellation before headers arrive and avoid onErrorDropped logs
Signed-off-by: Dariusz Jędrzejczyk <dariusz.jedrzejczyk@broadcom.com>
1 parent 5eef75a commit eb4009f

3 files changed

Lines changed: 115 additions & 5 deletions

File tree

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

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -521,7 +521,10 @@ public Mono<Void> sendMessage(McpSchema.JSONRPCMessage sentMessage) {
521521

522522
if (transportSession
523523
.markInitialized(httpResponse.headers().firstValue("mcp-session-id").orElse(null))) {
524-
reconnect(null).contextWrite(deliveredSink.contextView()).subscribe();
524+
// Fails only when the transport has been closed in the meantime,
525+
// in which case there is no stream left to open.
526+
reconnect(null).contextWrite(deliveredSink.contextView()).subscribe(ignored -> {
527+
}, t -> logger.debug("Not opening the SSE stream: {}", t.getMessage()));
525528
}
526529

527530
if (statusCode < 200 || statusCode >= 300) {

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

Lines changed: 33 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,9 @@
1717
import java.util.ArrayList;
1818
import java.util.List;
1919
import java.util.Optional;
20+
import java.util.concurrent.CancellationException;
21+
import java.util.concurrent.CompletableFuture;
22+
import java.util.concurrent.CompletionException;
2023
import java.util.concurrent.Flow;
2124
import java.util.concurrent.Flow.Publisher;
2225

@@ -201,16 +204,42 @@ static <T> Flux<T> cancel(Publisher<List<ByteBuffer>> body) {
201204
* Such a body must be subscribed to, or the connection it is read from is never
202205
* released. Should the exchange be cancelled once the response has arrived but before
203206
* its body could be subscribed to, the response is discarded, and its body cancelled.
207+
*
208+
* <p>
209+
* Cancelling the exchange before the response has arrived aborts the request. The
210+
* {@link HttpClient} then fails its future with a {@link CompletionException}
211+
* wrapping a {@link CancellationException}, which {@link Mono#fromFuture} does not
212+
* recognise as the outcome of its own cancellation and reports as a dropped error.
213+
* Only this method can cancel the future, so such a failure is always the expected
214+
* outcome of cancelling, and is ignored.
204215
* @param httpClient the client to send the request with
205216
* @param request the request to send
206217
*/
207218
static Mono<HttpResponse<Publisher<List<ByteBuffer>>>> sendAsync(HttpClient httpClient, HttpRequest request) {
208-
return Mono.fromFuture(() -> httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofPublisher()))
209-
.doOnDiscard(HttpResponse.class, response -> {
210-
if (response.body() instanceof Publisher<?> body) {
211-
cancelBody(body);
219+
return Mono.<HttpResponse<Publisher<List<ByteBuffer>>>>create(sink -> {
220+
CompletableFuture<HttpResponse<Publisher<List<ByteBuffer>>>> exchange = httpClient.sendAsync(request,
221+
HttpResponse.BodyHandlers.ofPublisher());
222+
sink.onCancel(() -> exchange.cancel(true));
223+
exchange.whenComplete((response, error) -> {
224+
if (error == null) {
225+
// Once cancelled, the response is discarded rather than emitted.
226+
sink.success(response);
227+
return;
228+
}
229+
Throwable cause = error instanceof CompletionException && error.getCause() != null ? error.getCause()
230+
: error;
231+
if (cause instanceof CancellationException) {
232+
sink.success();
233+
}
234+
else {
235+
sink.error(cause);
212236
}
213237
});
238+
}).doOnDiscard(HttpResponse.class, response -> {
239+
if (response.body() instanceof Publisher<?> body) {
240+
cancelBody(body);
241+
}
242+
});
214243
}
215244

216245
private static void cancelBody(Publisher<?> body) {
Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
/*
2+
* Copyright 2026-2026 the original author or authors.
3+
*/
4+
5+
package io.modelcontextprotocol.client.transport;
6+
7+
import java.io.IOException;
8+
import java.net.InetSocketAddress;
9+
import java.net.URI;
10+
import java.net.http.HttpClient;
11+
import java.net.http.HttpRequest;
12+
import java.util.List;
13+
import java.util.concurrent.CopyOnWriteArrayList;
14+
import java.util.concurrent.CountDownLatch;
15+
import java.util.concurrent.ExecutorService;
16+
import java.util.concurrent.Executors;
17+
import java.util.concurrent.TimeUnit;
18+
19+
import com.sun.net.httpserver.HttpServer;
20+
import org.junit.jupiter.api.AfterEach;
21+
import org.junit.jupiter.api.Test;
22+
import reactor.core.Disposable;
23+
import reactor.core.publisher.Hooks;
24+
25+
import static org.assertj.core.api.Assertions.assertThat;
26+
27+
class ResponseBodyHandlersSendAsyncTests {
28+
29+
private final ExecutorService executor = Executors.newCachedThreadPool();
30+
31+
private final CountDownLatch releaseResponse = new CountDownLatch(1);
32+
33+
private final List<Throwable> dropped = new CopyOnWriteArrayList<>();
34+
35+
private HttpServer server;
36+
37+
@AfterEach
38+
void tearDown() {
39+
Hooks.resetOnErrorDropped();
40+
this.releaseResponse.countDown();
41+
if (this.server != null) {
42+
this.server.stop(0);
43+
}
44+
this.executor.shutdownNow();
45+
}
46+
47+
@Test
48+
void cancellingBeforeTheResponseArrivesDropsNoError() throws Exception {
49+
CountDownLatch requestReceived = new CountDownLatch(1);
50+
this.server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
51+
this.server.setExecutor(this.executor);
52+
this.server.createContext("/", exchange -> {
53+
// Never responds, so that the exchange is cancelled while awaiting headers.
54+
requestReceived.countDown();
55+
try {
56+
this.releaseResponse.await();
57+
}
58+
catch (InterruptedException e) {
59+
Thread.currentThread().interrupt();
60+
}
61+
exchange.close();
62+
});
63+
this.server.start();
64+
Hooks.onErrorDropped(this.dropped::add);
65+
66+
HttpRequest request = HttpRequest
67+
.newBuilder(URI.create("http://127.0.0.1:" + this.server.getAddress().getPort() + "/"))
68+
.build();
69+
Disposable exchange = ResponseBodyHandlers.sendAsync(HttpClient.newHttpClient(), request).subscribe();
70+
assertThat(requestReceived.await(5, TimeUnit.SECONDS)).isTrue();
71+
72+
// The HttpClient fails the aborted exchange within cancel() itself.
73+
exchange.dispose();
74+
75+
assertThat(this.dropped).isEmpty();
76+
}
77+
78+
}

0 commit comments

Comments
 (0)