diff --git a/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java b/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java index 7ada0f7fa..a28ee7636 100644 --- a/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java +++ b/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java @@ -21,6 +21,14 @@ * The control is thread-safe and remains valid until its response completes. Calls made after completion have no * effect. *

+ * Calls take effect on the event loop that delivers the response body callbacks to the {@link AsyncHandler}. A call + * made on that event loop, for example from one of those callbacks, takes effect before it returns. A call made from + * any other thread takes effect asynchronously on the event loop, possibly after further callbacks. A + * {@link #resume()} always takes effect, even after a {@link #suspend()} that a callback makes before the resume is + * applied, so a resume is never lost. A {@link #suspend()} from another thread is skipped if, by the time it would take + * effect, the most recent call to {@link #suspend()} or {@link #resume()} was a resume. Calls made concurrently on + * different threads are not ordered otherwise. + *

* A control is also supplied when the final response headers end the response without a body. In that case, * {@link #suspend()} cannot defer completion: the control becomes inactive when * {@link AsyncHandler#onResponseBodyStart(ResponseBodyControl)} returns, and later calls have no effect. diff --git a/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java b/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java index 65fae4f8b..036465c84 100644 --- a/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java +++ b/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java @@ -39,6 +39,9 @@ public final class NettyResponseBodyControl implements ResponseBodyControl { private final boolean previousAutoRead; private final AtomicBoolean active = new AtomicBoolean(true); + // Set by suspend() and cleared by resume() before they are applied, so that a suspend() queued before a later + // resume() is dropped instead of undoing it. + private volatile boolean suspendRequested; private volatile boolean suspended; private volatile boolean bodyFullyRead; @@ -119,11 +122,13 @@ private NettyResponseBodyControl(NettyResponseFuture future, Channel channel, @Override public void suspend() { - execute(this::suspend0); + suspendRequested = true; + execute(this::applySuspend); } @Override public void resume() { + suspendRequested = false; execute(this::resume0); } @@ -132,6 +137,15 @@ public void cancel() { execute(this::cancel0); } + private void applySuspend() { + // Skipped if the most recent call was a resume(), so an older suspend() cannot undo it. A resume() is always + // applied, also after this, so a resume from another thread is never lost to a suspend() that a callback + // decided on an older view. + if (suspendRequested) { + suspend0(); + } + } + private void suspend0() { if (active.get() && !suspended) { suspensionStartedAction.run(); diff --git a/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java b/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java index 819b463f1..06069bad4 100644 --- a/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java +++ b/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java @@ -61,6 +61,7 @@ import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.SECONDS; @@ -177,6 +178,58 @@ public void suspensionAppliesHttp2FlowControlAndParentIsReused() throws Exceptio } } + @Test + public void earlierSuspendFromAnotherThreadDoesNotUndoLaterResume() throws Exception { + AtomicReference connection = new AtomicReference<>(); + try (AsyncHttpClient client = http2Client(connection::set)) { + RecordingHandler handler = new RecordingHandler(); + ListenableFuture request = client.prepareGet(url("/large")).execute(handler); + ResponseBodyControl control = handler.control.get(5, SECONDS); + assertTrue(largeResponseQueued.await(5, SECONDS)); + + // HTTP/2 streams share the event loop of their parent connection. + connection.get().eventLoop().submit(() -> { + callFromAnotherThread(control::suspend); + control.resume(); + return null; + }).get(5, SECONDS); + + largeResponseWritten.get(5, SECONDS); + assertSame(handler, request.get(5, SECONDS)); + assertEquals((long) FRAME_SIZE * FRAME_COUNT, handler.bodyBytes.get()); + assertNull(handler.throwable.get()); + } + } + + @Test + public void resumeFromAnotherThreadIsNotLostToALaterSuspend() throws Exception { + AtomicReference connection = new AtomicReference<>(); + try (AsyncHttpClient client = http2Client(connection::set)) { + RecordingHandler handler = new RecordingHandler(); + ListenableFuture request = client.prepareGet(url("/large")).execute(handler); + ResponseBodyControl control = handler.control.get(5, SECONDS); + assertTrue(largeResponseQueued.await(5, SECONDS)); + + connection.get().eventLoop().submit(() -> { + callFromAnotherThread(control::resume); + control.suspend(); + return null; + }).get(5, SECONDS); + + largeResponseWritten.get(5, SECONDS); + assertSame(handler, request.get(5, SECONDS)); + assertEquals((long) FRAME_SIZE * FRAME_COUNT, handler.bodyBytes.get()); + assertNull(handler.throwable.get()); + + Response sibling = client.prepareGet(url("/large-sibling")) + .setReadTimeout(Duration.ofSeconds(5)) + .execute() + .get(10, SECONDS); + assertEquals((long) FRAME_SIZE * SIBLING_FRAME_COUNT, sibling.getResponseBodyAsBytes().length); + assertEquals(1, connectionCount.get()); + } + } + @Test public void cancellationResetsOnlyTheHttp2Stream() throws Exception { try (AsyncHttpClient client = http2Client()) { @@ -293,13 +346,24 @@ public void suspensionCannotDeferBodylessResponse() throws Exception { } private AsyncHttpClient http2Client() { + return http2Client(ignored -> { + }); + } + + private AsyncHttpClient http2Client(Consumer connectionInitializer) { return asyncHttpClient(config() .setUseInsecureTrustManager(true) .setHttp2Enabled(true) .setHttp2InitialWindowSize(32 * 1024) .setMaxConnectionsPerHost(1) .setReadTimeout(Duration.ofMillis(100)) - .setRequestTimeout(Duration.ofSeconds(10))); + .setRequestTimeout(Duration.ofSeconds(10)) + .setHttpAdditionalChannelInitializer(connectionInitializer)); + } + + // The calling event loop stays busy until the call returns, so the control has to queue the call behind it. + private static void callFromAnotherThread(Runnable call) throws Exception { + CompletableFuture.runAsync(call).get(5, SECONDS); } private String url(String path) { diff --git a/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java b/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java index 846d8a567..a7baffbee 100644 --- a/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java +++ b/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java @@ -57,6 +57,7 @@ import java.net.InetSocketAddress; import java.time.Duration; import java.util.List; +import java.util.concurrent.Callable; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CountDownLatch; @@ -185,6 +186,62 @@ public void suspensionControlsReadsAndCompletedConnectionIsPooled() throws Excep } } + @Test + public void earlierSuspendFromAnotherThreadDoesNotUndoLaterResume() throws Exception { + AtomicReference clientChannel = new AtomicReference<>(); + try (AsyncHttpClient client = asyncHttpClient(config() + .setRequestTimeout(Duration.ofSeconds(10)) + .setHttpAdditionalChannelInitializer(clientChannel::set))) { + RecordingHandler handler = new RecordingHandler(false); + ListenableFuture request = client.prepareGet(url("/controlled")).execute(handler); + ResponseBodyControl control = handler.control.get(5, SECONDS); + ChannelHandlerContext server = responseContext.get(5, SECONDS); + Channel channel = clientChannel.get(); + + runOnEventLoop(channel, () -> { + callFromAnotherThread(control::suspend); + control.resume(); + return null; + }); + awaitEventLoop(channel); + + assertTrue(channel.config().isAutoRead(), "the queued suspend() must not override the later resume()"); + writeChunk(server, "one"); + assertEquals("one", handler.items.poll(5, SECONDS)); + writeLast(server); + assertSame(handler, request.get(5, SECONDS)); + assertNull(handler.throwable.get()); + } + } + + @Test + public void resumeFromAnotherThreadIsNotLostToALaterSuspend() throws Exception { + AtomicReference clientChannel = new AtomicReference<>(); + try (AsyncHttpClient client = asyncHttpClient(config() + .setRequestTimeout(Duration.ofSeconds(10)) + .setHttpAdditionalChannelInitializer(clientChannel::set))) { + RecordingHandler handler = new RecordingHandler(false); + ListenableFuture request = client.prepareGet(url("/controlled")).execute(handler); + ResponseBodyControl control = handler.control.get(5, SECONDS); + ChannelHandlerContext server = responseContext.get(5, SECONDS); + Channel channel = clientChannel.get(); + + runOnEventLoop(channel, () -> { + callFromAnotherThread(control::resume); + control.suspend(); + return null; + }); + awaitEventLoop(channel); + + assertTrue(channel.config().isAutoRead(), "the later suspend() must not drop the queued resume()"); + writeChunk(server, "one"); + assertEquals("one", handler.items.poll(5, SECONDS)); + writeLast(server); + assertSame(handler, request.get(5, SECONDS)); + assertNull(handler.throwable.get()); + } + } + @Test public void suspensionPausesReadTimeoutAndCancellationClosesConnection() throws Exception { try (AsyncHttpClient client = asyncHttpClient(config() @@ -562,6 +619,15 @@ private static void awaitEventLoop(Channel channel) throws InterruptedException }).sync(); } + private static void runOnEventLoop(Channel channel, Callable task) throws Exception { + channel.eventLoop().submit(task).get(5, SECONDS); + } + + // The calling event loop stays busy until the call returns, so the control has to queue the call behind it. + private static void callFromAnotherThread(Runnable call) throws Exception { + CompletableFuture.runAsync(call).get(5, SECONDS); + } + private static void writeChunk(ChannelHandlerContext ctx, String value) throws InterruptedException { ctx.executor().submit(() -> ctx.writeAndFlush( new DefaultHttpContent(Unpooled.copiedBuffer(value, CharsetUtil.US_ASCII)))).sync(); diff --git a/client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlTest.java b/client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlTest.java new file mode 100644 index 000000000..2d7ca592a --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlTest.java @@ -0,0 +1,186 @@ +/* + * Copyright (c) 2026 AsyncHttpClient Project. All rights reserved. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.asynchttpclient.netty; + +import io.netty.channel.Channel; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelOutboundHandlerAdapter; +import io.netty.channel.EventLoopGroup; +import io.netty.channel.MultiThreadIoEventLoopGroup; +import io.netty.channel.local.LocalChannel; +import io.netty.channel.local.LocalIoHandler; +import org.asynchttpclient.AsyncHandler; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.util.concurrent.Callable; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicInteger; + +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; + +public class NettyResponseBodyControlTest { + + private final NettyResponseFuture future = + new NettyResponseFuture<>(null, mock(AsyncHandler.class), null, 3, null, null, null); + private final AtomicInteger reads = new AtomicInteger(); + private final AtomicInteger suspensionsStarted = new AtomicInteger(); + private final AtomicInteger suspensionsEnded = new AtomicInteger(); + + private EventLoopGroup group; + private Channel channel; + + @BeforeEach + public void registerChannel() throws InterruptedException { + group = new MultiThreadIoEventLoopGroup(1, LocalIoHandler.newFactory()); + channel = new LocalChannel(); + channel.pipeline().addLast(new ChannelOutboundHandlerAdapter() { + @Override + public void read(ChannelHandlerContext ctx) { + reads.incrementAndGet(); + ctx.read(); + } + }); + group.register(channel).sync(); + } + + @AfterEach + public void closeChannel() throws InterruptedException { + channel.close().sync(); + group.shutdownGracefully(0, 100, MILLISECONDS).sync(); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void earlierSuspendFromAnotherThreadDoesNotUndoLaterResume(boolean autoRead) throws Exception { + NettyResponseBodyControl control = suspendedControl(autoRead); + + onEventLoop(() -> { + callFromAnotherThread(control::suspend); + control.resume(); + return null; + }); + awaitEventLoop(); + + assertFalse(NettyResponseBodyControl.isSuspended(future)); + assertEquals(autoRead, channel.config().isAutoRead()); + assertEquals(1, reads.get(), "resuming must request exactly one read"); + assertEquals(1, suspensionsStarted.get()); + assertEquals(1, suspensionsEnded.get()); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void resumeFromAnotherThreadIsNotLostToALaterSuspend(boolean autoRead) throws Exception { + NettyResponseBodyControl control = suspendedControl(autoRead); + control.resume(); + awaitEventLoop(); + + onEventLoop(() -> { + // A callback buffered a part and decided to suspend. Before it did, a consumer took that part and resumed. + callFromAnotherThread(control::resume); + control.suspend(); + return null; + }); + awaitEventLoop(); + + assertFalse(NettyResponseBodyControl.isSuspended(future), "nothing else would resume reads"); + assertEquals(autoRead, channel.config().isAutoRead()); + assertEquals(2, reads.get()); + assertEquals(2, suspensionsStarted.get()); + assertEquals(2, suspensionsEnded.get()); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void queuedCallsFromAnotherThreadEndInTheStateOfTheLatestCall(boolean autoRead) throws Exception { + NettyResponseBodyControl control = suspendedControl(autoRead); + control.resume(); + awaitEventLoop(); + + onEventLoop(() -> { + callFromAnotherThread(control::suspend); + callFromAnotherThread(control::resume); + callFromAnotherThread(control::suspend); + return null; + }); + awaitEventLoop(); + assertTrue(NettyResponseBodyControl.isSuspended(future)); + assertFalse(channel.config().isAutoRead()); + assertEquals(suspensionsEnded.get() + 1, suspensionsStarted.get()); + + onEventLoop(() -> { + callFromAnotherThread(control::resume); + callFromAnotherThread(control::suspend); + callFromAnotherThread(control::resume); + return null; + }); + awaitEventLoop(); + assertFalse(NettyResponseBodyControl.isSuspended(future)); + assertEquals(autoRead, channel.config().isAutoRead()); + assertEquals(suspensionsStarted.get(), suspensionsEnded.get()); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void suspendAfterASkippedOneStillTakesEffect(boolean autoRead) throws Exception { + NettyResponseBodyControl control = suspendedControl(autoRead); + onEventLoop(() -> { + callFromAnotherThread(control::suspend); + control.resume(); + return null; + }); + awaitEventLoop(); + assertFalse(NettyResponseBodyControl.isSuspended(future)); + + control.suspend(); + awaitEventLoop(); + assertTrue(NettyResponseBodyControl.isSuspended(future)); + assertFalse(channel.config().isAutoRead()); + } + + private NettyResponseBodyControl suspendedControl(boolean autoRead) throws Exception { + return onEventLoop(() -> { + channel.config().setAutoRead(autoRead); + NettyResponseBodyControl control = NettyResponseBodyControl.create(future, channel, + suspensionsStarted::incrementAndGet, suspensionsEnded::incrementAndGet, () -> { + }, ignored -> { + }); + control.suspend(); + return control; + }); + } + + private T onEventLoop(Callable task) throws Exception { + return channel.eventLoop().submit(task).get(5, SECONDS); + } + + private void awaitEventLoop() throws Exception { + onEventLoop(() -> null); + } + + // The calling event loop stays busy until the call returns, so the control has to queue the call behind it. + private static void callFromAnotherThread(Runnable call) throws Exception { + CompletableFuture.runAsync(call).get(5, SECONDS); + } +}