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
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,14 @@
* The control is thread-safe and remains valid until its response completes. Calls made after completion have no
* effect.
* <p>
* 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.
* <p>
* 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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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);
}

Expand All @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -177,6 +178,58 @@ public void suspensionAppliesHttp2FlowControlAndParentIsReused() throws Exceptio
}
}

@Test
public void earlierSuspendFromAnotherThreadDoesNotUndoLaterResume() throws Exception {
AtomicReference<Channel> connection = new AtomicReference<>();
try (AsyncHttpClient client = http2Client(connection::set)) {
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));

// 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<Channel> connection = new AtomicReference<>();
try (AsyncHttpClient client = http2Client(connection::set)) {
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));

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()) {
Expand Down Expand Up @@ -293,13 +346,24 @@ public void suspensionCannotDeferBodylessResponse() throws Exception {
}

private AsyncHttpClient http2Client() {
return http2Client(ignored -> {
});
}

private AsyncHttpClient http2Client(Consumer<Channel> 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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -185,6 +186,62 @@ public void suspensionControlsReadsAndCompletedConnectionIsPooled() throws Excep
}
}

@Test
public void earlierSuspendFromAnotherThreadDoesNotUndoLaterResume() 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);
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<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);
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()
Expand Down Expand Up @@ -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();
Expand Down
Loading
Loading