From b2f6e15a4cce7faa19bf189ca6d6a9a5545717be Mon Sep 17 00:00:00 2001 From: Matthias Kurz Date: Tue, 6 Oct 2026 16:19:43 +0200 Subject: [PATCH] Never let a queued suspend override a later resume A ResponseBodyControl call made off the event loop is queued and takes effect later on the loop, in the order the calls were queued. A suspend() queued from another thread could therefore take effect after a resume() that the event loop made in the meantime, for example from a body callback, and leave reads suspended with nothing left to resume them. Record a suspend request when suspend() is called and clear it when resume() is called. A queued suspend() is skipped if, when it would take effect, the most recent call was a resume(). A resume() is always applied, so it is never lost to a suspend() that a callback made after it was queued. Calls made on the event loop take effect before they return, as before. Calls made concurrently on different threads are not ordered beyond that. Document when control calls take effect. Claude Code on behalf of Matthias Kurz Co-Authored-By: Claude Opus 5.5 --- .../asynchttpclient/ResponseBodyControl.java | 8 + .../netty/NettyResponseBodyControl.java | 16 +- .../Http2ResponseBodyControlTest.java | 66 ++++++- .../ResponseBodyControlTest.java | 66 +++++++ .../netty/NettyResponseBodyControlTest.java | 186 ++++++++++++++++++ 5 files changed, 340 insertions(+), 2 deletions(-) create mode 100644 client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlTest.java 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); + } +}