From 686cbeb69bcbf45d13a9d12dc5e588b16bad2416 Mon Sep 17 00:00:00 2001 From: Jonathan Prates Date: Thu, 27 Aug 2026 11:54:49 +0100 Subject: [PATCH] Fix WebClient response cancellation race Keep track while the response is being delivered. When cancellation happens during that call, allow WebClient to finish without subscribing again to a body that Spring is already draining. Keep the existing error when code subscribes to the body after response delivery is complete. Closes gh-37204 Signed-off-by: Jonathan Prates --- .../reactive/ReactorClientHttpConnector.java | 2 +- .../reactive/ReactorClientHttpResponse.java | 25 +++++-- .../reactive/ClientHttpConnectorTests.java | 66 +++++++++++++++++++ 3 files changed, 87 insertions(+), 6 deletions(-) diff --git a/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpConnector.java b/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpConnector.java index 3b0a14b3cb96..b15aad2f6b34 100644 --- a/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpConnector.java +++ b/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpConnector.java @@ -169,7 +169,7 @@ public Mono connect(HttpMethod method, URI uri, ReactorClientHttpResponse clientResponse = new ReactorClientHttpResponse(response, connection); responseRef.set(clientResponse); registerAttributeCallback(connection); - return Mono.just((ClientHttpResponse) clientResponse); + return Mono.create(clientResponse::deliver); }) .next() .doOnCancel(() -> { diff --git a/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpResponse.java b/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpResponse.java index 0fb14b640d2b..d003c48aaf64 100644 --- a/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpResponse.java +++ b/spring-web/src/main/java/org/springframework/http/client/reactive/ReactorClientHttpResponse.java @@ -17,6 +17,7 @@ package org.springframework.http.client.reactive; import java.util.Collection; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiFunction; @@ -26,6 +27,7 @@ import org.apache.commons.logging.LogFactory; import org.jspecify.annotations.Nullable; import reactor.core.publisher.Flux; +import reactor.core.publisher.MonoSink; import reactor.netty.ChannelOperationsId; import reactor.netty.Connection; import reactor.netty.NettyInbound; @@ -65,6 +67,7 @@ class ReactorClientHttpResponse implements ClientHttpResponse { // 0 - not subscribed, 1 - subscribed, 2 - cancelled via connector (before subscribe) private final AtomicInteger state = new AtomicInteger(); + private final AtomicBoolean responseDeliveryInProgress = new AtomicBoolean(); /** @@ -95,15 +98,17 @@ public String getId() { @Override public Flux getBody() { - return this.inbound.receive() - .doOnSubscribe(s -> { + return Flux.defer(() -> { if (this.state.compareAndSet(0, 1)) { - return; + return this.inbound.receive(); } if (this.state.get() == 2) { - throw new IllegalStateException( - "The client response body has been released already due to cancellation."); + return this.responseDeliveryInProgress.get() + ? Flux.empty() + : Flux.error(new IllegalStateException("The client response body has " + + "been released already due to cancellation.")); } + return this.inbound.receive(); }) .map(byteBuf -> { byteBuf.retain(); @@ -169,6 +174,16 @@ void releaseAfterCancel(HttpMethod method) { } } + void deliver(MonoSink sink) { + this.responseDeliveryInProgress.set(true); + try { + sink.success(this); + } + finally { + this.responseDeliveryInProgress.set(false); + } + } + private boolean mayHaveBody(HttpMethod method) { int code = getStatusCode().value(); return !((code >= 100 && code < 200) || code == 204 || code == 205 || diff --git a/spring-web/src/test/java/org/springframework/http/client/reactive/ClientHttpConnectorTests.java b/spring-web/src/test/java/org/springframework/http/client/reactive/ClientHttpConnectorTests.java index 1a5d3411ada5..039eec6fac1c 100644 --- a/spring-web/src/test/java/org/springframework/http/client/reactive/ClientHttpConnectorTests.java +++ b/spring-web/src/test/java/org/springframework/http/client/reactive/ClientHttpConnectorTests.java @@ -32,7 +32,9 @@ import java.util.List; import java.util.Random; import java.util.Set; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import java.util.function.Function; @@ -48,8 +50,11 @@ import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; +import reactor.core.Disposable; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import reactor.netty.http.client.HttpClient; +import reactor.netty.resources.ConnectionProvider; import reactor.test.StepVerifier; import org.springframework.core.io.buffer.DataBuffer; @@ -79,6 +84,8 @@ class ClientHttpConnectorTests { private final MockWebServer server = new MockWebServer(); + private final List connectionProviders = new ArrayList<>(); + @BeforeEach void startServer() throws IOException { server.start(); @@ -86,6 +93,7 @@ void startServer() throws IOException { @AfterEach void stopServer() { + this.connectionProviders.forEach(provider -> provider.disposeLater().block()); server.close(); } @@ -182,6 +190,64 @@ void cancelResponseBody(ClientHttpConnector connector) { .verify(); } + @Test + void cancelWhileResponseDeliveryIsInProgress() throws Exception { + prepareResponse(builder -> builder.body("body")); + ConnectionProvider connectionProvider = ConnectionProvider.create("cancelWhileResponseDelivery", 1); + this.connectionProviders.add(connectionProvider); + ReactorClientHttpConnector connector = new ReactorClientHttpConnector(HttpClient.create(connectionProvider)); + + CountDownLatch responseReceived = new CountDownLatch(1); + CountDownLatch continueResponse = new CountDownLatch(1); + CompletableFuture bodyCompleted = new CompletableFuture<>(); + + Mono responseMono = connector + .connect(HttpMethod.GET, this.server.url("/").uri(), ReactiveHttpOutputMessage::setComplete) + .doOnNext(response -> { + responseReceived.countDown(); + try { + continueResponse.await(5, TimeUnit.SECONDS); + } + catch (InterruptedException ex) { + Thread.currentThread().interrupt(); + } + }) + .doOnNext(response -> + response.getBody().subscribe(DataBufferUtils::release, + bodyCompleted::completeExceptionally, () -> bodyCompleted.complete(true))); + + Disposable subscription = responseMono.subscribe(); + + assertThat(responseReceived.await(5, TimeUnit.SECONDS)).isTrue(); + subscription.dispose(); + continueResponse.countDown(); + assertThat(bodyCompleted.get(5, TimeUnit.SECONDS)).isTrue(); + + prepareResponse(builder -> builder.body("body")); + StepVerifier.create(connector + .connect(HttpMethod.GET, this.server.url("/").uri(), ReactiveHttpOutputMessage::setComplete) + .flatMap(response -> response.getBody().doOnNext(DataBufferUtils::release).then())) + .verifyComplete(); + + assertThat(this.server.takeRequest().getConnectionIndex()) + .isEqualTo(this.server.takeRequest().getConnectionIndex()); + } + + @Test + void bodySubscribedAfterCancellation() { + prepareResponse(builder -> builder.body("body")); + + ClientHttpResponse response = new ReactorClientHttpConnector() + .connect(HttpMethod.GET, this.server.url("/").uri(), ReactiveHttpOutputMessage::setComplete) + .block(); + assertThat(response).isInstanceOf(ReactorClientHttpResponse.class); + ((ReactorClientHttpResponse) response).releaseAfterCancel(HttpMethod.GET); + + StepVerifier.create(response.getBody()) + .expectErrorMessage("The client response body has been released already due to cancellation.") + .verify(); + } + @ParameterizedConnectorTest void cookieExpireValueSetAsMaxAge(ClientHttpConnector connector) { ZonedDateTime tomorrow = ZonedDateTime.now(ZoneOffset.UTC).plusDays(1);