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
34 changes: 34 additions & 0 deletions client/src/main/java/org/asynchttpclient/ResponseBodyControl.java
Original file line number Diff line number Diff line change
Expand Up @@ -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.
* <p>
* Unlike calls to the other methods, a task passed to {@link #execute(Runnable)} still runs after completion.
*
* @since 3.0.14
*/
Expand Down Expand Up @@ -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)}.
* <p>
* 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.
* <p>
* 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.
* <p>
* 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.
* <p>
* 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.
* <p>
* 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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -292,6 +293,29 @@ public void suspensionCannotDeferBodylessResponse() throws Exception {
}
}

@Test
public void executeFromAnotherThreadCanResumeAnHttp2Stream() throws Exception {
try (AsyncHttpClient client = http2Client()) {
RecordingHandler handler = new RecordingHandler();
ListenableFuture<RecordingHandler> 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<Thread> 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)
Expand Down
151 changes: 151 additions & 0 deletions client/src/test/java/org/asynchttpclient/ResponseBodyControlTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -557,6 +557,157 @@ protected void initChannel(Channel channel) {
tlsServerPort = ((InetSocketAddress) tlsServerChannel.localAddress()).getPort();
}

@Test
public void executeRunsInlineOnTheEventLoopAndQueuedFromAnotherThread() throws Exception {
AtomicReference<Channel> clientChannel = new AtomicReference<>();
CompletableFuture<Boolean> 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<RecordingHandler> 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<Thread> taskThread = new CompletableFuture<>();
CompletableFuture<Boolean> 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<Channel> 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<String> 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<RecordingHandler> 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<Channel> clientChannel = new AtomicReference<>();
try (AsyncHttpClient client = asyncHttpClient(config()
.setRequestTimeout(Duration.ofSeconds(10))
.setHttpAdditionalChannelInitializer(clientChannel::set))) {
RecordingHandler handler = new RecordingHandler(false);
ListenableFuture<RecordingHandler> 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<Thread> 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<Thread> 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();
Expand Down
Original file line number Diff line number Diff line change
@@ -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());
}
}
6 changes: 6 additions & 0 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
Loading