Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,7 @@ public Mono<ClientHttpResponse> 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(() -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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;
Expand Down Expand Up @@ -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();


/**
Expand Down Expand Up @@ -95,15 +98,17 @@ public String getId() {

@Override
public Flux<DataBuffer> 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();
Expand Down Expand Up @@ -169,6 +174,16 @@ void releaseAfterCancel(HttpMethod method) {
}
}

void deliver(MonoSink<ClientHttpResponse> 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 ||
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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;
Expand Down Expand Up @@ -79,13 +84,16 @@ class ClientHttpConnectorTests {

private final MockWebServer server = new MockWebServer();

private final List<ConnectionProvider> connectionProviders = new ArrayList<>();

@BeforeEach
void startServer() throws IOException {
server.start();
}

@AfterEach
void stopServer() {
this.connectionProviders.forEach(provider -> provider.disposeLater().block());
server.close();
}

Expand Down Expand Up @@ -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<Boolean> bodyCompleted = new CompletableFuture<>();

Mono<ClientHttpResponse> 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);
Expand Down