From 0fd14889af36f01c193a45110c4837aa3f410d14 Mon Sep 17 00:00:00 2001 From: Matthias Kurz Date: Wed, 7 Oct 2026 10:45:02 +0200 Subject: [PATCH] Add ResponseBodyControl.execute A consumer that buffers body parts on another thread decides there whether to resume reads, but resume() from that thread only takes effect later on the event loop. In between, the event loop can deliver another part that already meets the consumer's demand, and the resume then lets one more read through. With a highly compressible body, that read can inflate to many megabytes beyond demand. Only the caller knows which state its decision depends on, so AHC cannot re-check it. Add execute(Runnable), which runs a task on the event loop that delivers the response body, in sequence with the callbacks made there: before returning when called on that event loop, and queued to it otherwise. suspend(), resume() and cancel() called from the task take effect immediately, so a consumer can decide there on everything delivered so far. onThrowable can run on other threads, after a cancel or a timeout, so a task is not ordered with it. Unlike the other methods, the task also runs after completion; it does not run once the event loop rejects tasks. NettyResponseBodyControl already applied its own control calls this way; that method is now public. The default implementation runs the task on the calling thread, for implementations that apply control calls synchronously. Play WS needs this: it currently reaches the event loop through Netty's internal ThreadExecutorMap to avoid the extra read. Claude Code on behalf of Matthias Kurz Co-Authored-By: Claude Opus 5.5 --- .../asynchttpclient/ResponseBodyControl.java | 34 ++++ .../netty/NettyResponseBodyControl.java | 4 +- .../Http2ResponseBodyControlTest.java | 24 +++ .../ResponseBodyControlTest.java | 151 ++++++++++++++++++ .../NettyResponseBodyControlExecuteTest.java | 53 ++++++ pom.xml | 6 + 6 files changed, 271 insertions(+), 1 deletion(-) create mode 100644 client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlExecuteTest.java diff --git a/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java b/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java index 7ada0f7fa..731b98462 100644 --- a/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java +++ b/client/src/main/java/org/asynchttpclient/ResponseBodyControl.java @@ -24,6 +24,8 @@ * 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. + *

+ * Unlike calls to the other methods, a task passed to {@link #execute(Runnable)} still runs after completion. * * @since 3.0.14 */ @@ -59,4 +61,36 @@ public interface ResponseBodyControl { * intended for cancellation after the callback has returned, including from another thread. */ void cancel(); + + /** + * Runs {@code task} on the event loop that delivers this response's body, in sequence with the {@link AsyncHandler} + * callbacks that AHC makes on that event loop, such as + * {@link AsyncHandler#onBodyPartReceived(HttpResponseBodyPart)}. + *

+ * Called on that event loop, for example from one of those callbacks, the task runs before this method returns. + * Called from any other thread, the task is queued to the event loop and runs there after the work already pending + * on it; the event loop can deliver further body parts before the task runs. Either way, {@link #suspend()}, + * {@link #resume()} and {@link #cancel()} called from the task take effect immediately, and the task sees every + * body part delivered before it runs. + *

+ * Not every callback runs on that event loop. {@link AsyncHandler#onThrowable(Throwable)} can run on the thread + * that cancels the request, or on a timer thread when a timeout fires. A task passed to this method from there is + * queued like one from any other thread, and can run while that callback is still running. + *

+ * This lets a consumer that buffers body parts on another thread decide whether to resume reads based on everything + * delivered so far. Calling {@link #resume()} directly from that thread instead acts on a decision that a body part + * delivered in the meantime can make stale, and can let another read through after the consumer's demand is met. + *

+ * The task can run on a transport thread and must not block. It also runs after the response has completed. If the + * event loop no longer accepts tasks, for example because it is shutting down, the task does not run. + *

+ * The default implementation runs the task on the calling thread. That suits implementations that apply control + * calls synchronously on the calling thread; it does not order the task with callbacks made on other threads. + * + * @param task the task to run + * @since 3.0.15 + */ + default void execute(Runnable task) { + task.run(); + } } diff --git a/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java b/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java index 65fae4f8b..43f2c1a4a 100644 --- a/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java +++ b/client/src/main/java/org/asynchttpclient/netty/NettyResponseBodyControl.java @@ -186,7 +186,9 @@ private void endSuspension() { } } - private void execute(Runnable task) { + @Override + public void execute(Runnable task) { + Objects.requireNonNull(task, "task"); if (channel.eventLoop().inEventLoop()) { task.run(); } else { diff --git a/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java b/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java index 819b463f1..ea8c35834 100644 --- a/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java +++ b/client/src/test/java/org/asynchttpclient/Http2ResponseBodyControlTest.java @@ -68,6 +68,7 @@ import static org.asynchttpclient.Dsl.config; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotSame; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -292,6 +293,29 @@ public void suspensionCannotDeferBodylessResponse() throws Exception { } } + @Test + public void executeFromAnotherThreadCanResumeAnHttp2Stream() throws Exception { + try (AsyncHttpClient client = http2Client()) { + RecordingHandler handler = new RecordingHandler(); + ListenableFuture request = client.prepareGet(url("/large")).execute(handler); + ResponseBodyControl control = handler.control.get(5, SECONDS); + assertTrue(largeResponseQueued.await(5, SECONDS)); + assertThrows(TimeoutException.class, () -> request.get(250, MILLISECONDS)); + + CompletableFuture taskThread = new CompletableFuture<>(); + control.execute(() -> { + taskThread.complete(Thread.currentThread()); + control.resume(); + }); + + assertSame(handler, request.get(10, SECONDS)); + assertNotSame(Thread.currentThread(), taskThread.get(5, SECONDS), + "a task from another thread must run on the stream's event loop"); + assertEquals((long) FRAME_SIZE * FRAME_COUNT, handler.bodyBytes.get()); + assertNull(handler.throwable.get()); + } + } + private AsyncHttpClient http2Client() { return asyncHttpClient(config() .setUseInsecureTrustManager(true) diff --git a/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java b/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java index 846d8a567..def93302a 100644 --- a/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java +++ b/client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java @@ -557,6 +557,157 @@ protected void initChannel(Channel channel) { tlsServerPort = ((InetSocketAddress) tlsServerChannel.localAddress()).getPort(); } + @Test + public void executeRunsInlineOnTheEventLoopAndQueuedFromAnotherThread() throws Exception { + AtomicReference clientChannel = new AtomicReference<>(); + CompletableFuture ranInlineInCallback = new CompletableFuture<>(); + try (AsyncHttpClient client = asyncHttpClient(config() + .setRequestTimeout(Duration.ofSeconds(10)) + .setHttpAdditionalChannelInitializer(clientChannel::set))) { + RecordingHandler handler = new RecordingHandler(true) { + @Override + public State onBodyPartReceived(HttpResponseBodyPart bodyPart) throws IOException { + AtomicBoolean ran = new AtomicBoolean(); + super.responseBodyControl.execute(() -> ran.set(true)); + ranInlineInCallback.complete(ran.get()); + return super.onBodyPartReceived(bodyPart); + } + }; + 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(); + + CompletableFuture taskThread = new CompletableFuture<>(); + CompletableFuture autoReadSeenByTask = new CompletableFuture<>(); + control.resume(); + control.execute(() -> { + taskThread.complete(Thread.currentThread()); + autoReadSeenByTask.complete(channel.config().isAutoRead()); + }); + assertTrue(channel.eventLoop().inEventLoop(taskThread.get(5, SECONDS)), + "a task from another thread must run on the channel event loop"); + assertTrue(autoReadSeenByTask.get(5, SECONDS), "a task must run after control calls queued before it"); + + writeChunk(server, "one"); + assertEquals("one", handler.items.poll(5, SECONDS)); + assertTrue(ranInlineInCallback.get(5, SECONDS), "a task from a callback must run before execute returns"); + + writeLast(server); + control.resume(); + assertSame(handler, request.get(5, SECONDS)); + assertNull(handler.throwable.get()); + } + } + + @Test + public void executeLetsAnotherThreadDecideOnEverythingDeliveredBeforeTheTask() throws Exception { + AtomicReference clientChannel = new AtomicReference<>(); + CountDownLatch firstPartBuffered = new CountDownLatch(1); + try (AsyncHttpClient client = asyncHttpClient(config() + .setRequestTimeout(Duration.ofSeconds(10)) + .setHttpAdditionalChannelInitializer(clientChannel::set))) { + RecordingHandler handler = new RecordingHandler(true) { + @Override + public State onBodyPartReceived(HttpResponseBodyPart bodyPart) throws IOException { + if (firstPartBuffered.getCount() == 0) { + return super.onBodyPartReceived(bodyPart); + } + // While this part is being delivered, a consumer on another thread finds its buffer empty. It + // decides in a task whether to resume, so the decision sees the part buffered below. + LinkedBlockingQueue buffer = super.items; + ResponseBodyControl control = super.responseBodyControl; + runOnAnotherThread(() -> control.execute(() -> { + if (buffer.isEmpty()) { + control.resume(); + } + })); + State state = super.onBodyPartReceived(bodyPart); + firstPartBuffered.countDown(); + return state; + } + }; + 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(); + + control.resume(); + writeChunk(server, "one"); + assertTrue(firstPartBuffered.await(5, SECONDS)); + awaitEventLoop(channel); + + assertFalse(channel.config().isAutoRead(), "the task must see the part buffered after it was queued"); + writeChunk(server, "two"); + assertEquals("one", handler.items.poll(5, SECONDS)); + assertNull(handler.items.poll(250, MILLISECONDS), "a suspended response must not read new body bytes"); + + control.execute(() -> { + if (handler.items.isEmpty()) { + control.resume(); + } + }); + assertEquals("two", handler.items.poll(5, SECONDS)); + writeLast(server); + control.resume(); + assertSame(handler, request.get(5, SECONDS)); + assertNull(handler.throwable.get()); + } + } + + @Test + public void executeStillRunsTheTaskAfterTheResponseCompleted() 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); + control.resume(); + writeLast(server); + assertSame(handler, request.get(5, SECONDS)); + + CompletableFuture taskThread = new CompletableFuture<>(); + control.execute(() -> { + // A control call has no effect once the response has completed. + control.resume(); + taskThread.complete(Thread.currentThread()); + }); + assertTrue(clientChannel.get().eventLoop().inEventLoop(taskThread.get(5, SECONDS))); + assertNull(handler.throwable.get()); + } + } + + @Test + public void defaultExecuteRunsTheTaskOnTheCallingThread() { + ResponseBodyControl control = new ResponseBodyControl() { + @Override + public void suspend() { + } + + @Override + public void resume() { + } + + @Override + public void cancel() { + } + }; + AtomicReference taskThread = new AtomicReference<>(); + + control.execute(() -> taskThread.set(Thread.currentThread())); + assertSame(Thread.currentThread(), taskThread.get()); + } + + private static void runOnAnotherThread(Runnable action) { + Thread thread = new Thread(action, "response-body-control-consumer"); + thread.start(); + assertDoesNotThrow(() -> thread.join(SECONDS.toMillis(5))); + assertFalse(thread.isAlive(), "the other thread must finish"); + } + private static void awaitEventLoop(Channel channel) throws InterruptedException { channel.eventLoop().submit(() -> { }).sync(); diff --git a/client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlExecuteTest.java b/client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlExecuteTest.java new file mode 100644 index 000000000..fb5cc7a6f --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/netty/NettyResponseBodyControlExecuteTest.java @@ -0,0 +1,53 @@ +/* + * 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.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.Test; + +import java.util.concurrent.atomic.AtomicBoolean; + +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.mockito.Mockito.mock; + +public class NettyResponseBodyControlExecuteTest { + + @Test + public void executeDoesNotRunTheTaskOnceTheEventLoopRejectsIt() throws Exception { + EventLoopGroup group = new MultiThreadIoEventLoopGroup(1, LocalIoHandler.newFactory()); + Channel channel = new LocalChannel(); + group.register(channel).sync(); + NettyResponseFuture future = + new NettyResponseFuture<>(null, mock(AsyncHandler.class), null, 3, null, null, null); + NettyResponseBodyControl control = channel.eventLoop().submit(() -> NettyResponseBodyControl.create( + future, channel, () -> { + }, ignored -> { + })).get(5, SECONDS); + channel.close().sync(); + group.shutdownGracefully(0, 0, SECONDS).sync(); + + AtomicBoolean ran = new AtomicBoolean(); + assertDoesNotThrow(() -> control.execute(() -> ran.set(true))); + assertFalse(ran.get()); + } +} diff --git a/pom.xml b/pom.xml index f34f516d9..c5be9ad02 100644 --- a/pom.xml +++ b/pom.xml @@ -506,6 +506,12 @@ "configuration": { "ignore": true, "differences": [ + { + "code": "java.method.visibilityIncreased", + "old": "method void org.asynchttpclient.netty.NettyResponseBodyControl::execute(java.lang.Runnable)", + "new": "method void org.asynchttpclient.netty.NettyResponseBodyControl::execute(java.lang.Runnable)", + "justification": "NettyResponseBodyControl implements the new ResponseBodyControl.execute default method with its existing event-loop dispatch, which was private. Making a private method public is binary and source compatible, and the class is final, so no subclass can be affected. Scoped to this exact method." + }, { "code": "java.class.externalClassExposedInAPI", "justification": "Netty types are part of the public API by design"