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"