Skip to content

Commit c1af3a4

Browse files
Arm request timeouts on an event loop (#2313)
## Problem Request and read timeouts are armed on the client's `HashedWheelTimer`. That has two properties that only show up on short deadlines: * **A wheel quantizes.** It fires on the first tick at or after the deadline, so a deadline near or below `hashedWheelTimerTickDuration` is rounded up to it. * **One thread carries every expiry** for the whole client, and `HashedWheelTimer`'s default `taskExecutor` is `ImmediateExecutor`, so each expiry runs inline on the wheel thread - including `future.completeExceptionally(...)` and therefore whatever the caller chained onto the response future. On a one-second budget the first costs 0.3% and nobody notices. On a budget of tens of milliseconds a tick is a large fraction of it, and a burst of expiries has no headroom to absorb before the wheel starts running late. ## Measured 2000 timeouts armed as one burst on Netty 4.2.16, JDK 17, tasks doing nothing but recording their own lag. This is the floor; real work on the firing thread only adds to it. | instrument | deadline | mean | p50 | p99 | max | |---|---:|---:|---:|---:|---:| | wheel, tick 5 ms | 20 ms | +2.7 | +2 | +5 | +5 | | wheel, tick 1 ms | 20 ms | +1.3 | +1 | +2 | +2 | | `EventLoop.schedule` | 20 ms | +0.0 | +0 | +0 | +0 | | wheel, tick 5 ms | 1000 ms | +3.0 | +3 | +3 | +3 | | wheel, tick 1 ms | 1000 ms | +1.9 | +2 | +2 | +2 | | `EventLoop.schedule` | 1000 ms | +1.6 | +2 | +2 | +2 | An event loop shows zero overshoot because it schedules by deadline and derives its own `select()` timeout from the nearest one. There is no quantum to round to. This was a throwaway probe rather than JMH - `client/src/jmh/java` is not currently wired into the build, so its benchmarks do not compile. Happy to add a proper benchmark if that is fixed first, or as part of this. ## Change `AsyncHttpClientConfig#isUseEventLoopTimeouts()`, **off by default**, arms the request and read timeouts on an event loop instead of the timer. * **Always the channel's own loop.** On the pooled path the channel is already in hand, so its loop is used and the timeout expires on the thread that would have to close it. On the connect path there is no channel yet - deliberately, so that the timeout also bounds address resolution and the connect - so it is armed on the timer and moved onto the loop once the connect succeeds. No other loop is ever used; see the review round below for why that matters. * **Arming allocates nothing extra.** The cancellation handle lives on the task rather than in a wrapper, and the existing `done` flag stands in for the scheduler's already-expired flag, which the two schedulers spell differently. * **Shutdown race closed.** `isShuttingDown()` can return false and `schedule` reject immediately after. Netty answers a rejected timeout with a logged warning rather than an exception, which would leave the exchange with nothing to end it, so a rejection falls back to the timer. * The connection-pool cleaner stays on the timer either way. Off by default because the expiry - and so whatever the caller chained onto the future - then runs on an I/O thread, and blocking one stalls every connection it serves. The javadoc says so and points callers at `handleAsync`. ## Why not a wheel per event loop That is how the Aerospike client solves the same problem: `EventLoopBase` owns a `HashedWheelTimer` that is a `Runnable` the loop ticks itself. Deliberately not copied here. A wheel arms in O(1) against O(log n) for a deadline queue, but at a few thousand timeouts per loop that is a dozen comparisons, while the quantization it reintroduces costs milliseconds on a 20 ms budget - the third row above is the whole point. A wheel also has to be ticked forever, waking every loop even with nothing armed. Aerospike wrote its own because its `EventLoop` abstracts over NIO, Netty and direct NIO and needed one timer; AHC is Netty-only and gets a per-loop deadline queue for free. ## Review round 1 Most of the substance of this PR changed in review, so the sections above describe the current shape rather than what was first pushed. Two things are worth calling out here because they were design errors, not polish: **The loop is now only ever the channel's own.** It used to come from `EventLoopGroup#next()`, which is almost never the loop the channel ends up on: `initAndRegister` draws from the same chooser, so the two agreed about one time in N. Every completion then cancelled an entry on a foreign loop, and until the original deadline that entry sat in a queue whose loop it would wake for a request that had long finished. Drawing from the chooser also shifted which loops connections land on. The pooled path has the channel in hand; the connect path arms on the timer, as before this branch, and `NettyConnectListener` moves the timeouts onto the loop once the connect succeeds, next to the `attachChannel` that publishes it on the future. That listener already runs on the channel's loop, so the move costs a same-thread schedule and no wakeup. **Arming left the `TimeoutsHolder` constructor.** The task holds the holder and can run the moment it is armed, and an event loop does not round a short deadline up to a tick, so the expiry could reach a holder whose fields were not yet frozen, a future that had not been handed the holder, and on the pooled path a future with no channel attached - which aborted with `null` and left the pooled socket open. The caller now publishes the holder, attaches the channel, and calls `start()` last. Also from the review: `arm` re-checks `cancelled` after recording its handle, so an exchange that finishes mid-arming cannot leave behind an entry nobody will cancel; `cancelArmed` catches the `RejectedExecutionException` that Netty's off-loop cancellation path can raise on a closing client, which had never escaped `ListenableFuture#cancel` before; the cancellation handle is two typed fields rather than an `Object` and `instanceof`; `requestTimeoutArmed` is gone in favour of a null test on the task; the rationale lives on `isUseEventLoopTimeouts()` alone; and the new option sits in the `// timeouts` group everywhere rather than splitting the two `failedIpCooldown` entries. ## API compatibility No `revapi` entries. `implements Runnable` and `run()` live on the two subclasses rather than on `TimeoutTimerTask`, which leaves that class's surface unchanged: both subclasses already declared `run(Timeout)` without a throws clause, so inside a subclass `run()` calls its own override and has nothing to catch. No dead handler, and nothing narrowed. The knock-on is that only the concrete classes are both a `TimerTask` and a `Runnable`, so `TimeoutsHolder#arm` takes that intersection as a type parameter. Everything else is additive: the existing `TimeoutsHolder` constructor is kept and delegates, and nothing is removed. `start()` is called from `NettyResponseFuture#setTimeoutsHolder` rather than by the sender. `TimeoutsHolder` has a public constructor in an exported package, and splitting the arming out of it would otherwise leave an outside caller free to install a holder and get an exchange with no request timeout at all. ## Tests Four cases in `EventLoopTimeoutTest`, asserting where an expiry is delivered from rather than what it does. They hand the config their own `Timer` and `EventLoopGroup` so the assertions are against those objects and not against thread names, which a pool name containing `timer` or a configured thread factory would have broken with no bug present: - the timer default, by identity against the timer's own thread; - a connecting exchange, on the loop of the channel `onTcpConnectSuccess` reported; - an exchange on a pooled channel, on the loop of the channel `onConnectionPooled` reported, which is also what says it reused the connection rather than opening one of its own; - a read timeout under the new mode, which is armed after the request is written and so exercises a different arming path. The group has eight loops. With two, a timeout armed on the wrong loop is on the right one half the time, and these assertions would have passed about half the runs against the bug they exist to catch; at eight, reverting `timeoutExecutor` to `next()` fails the pooled case. The connecting case, and the pooled case's first request, get a one second budget on purpose. A deadline reached before connecting would be delivered from the timer quite correctly, there being no channel to deliver it from, and would prove nothing either way; and the pooled case's first request is the cold one - class loading, the connect, the server's first response - and is meant to succeed. One gap, called out rather than papered over: the `RejectedExecutionException` fallback in `arm` has no test. Reaching it needs a loop that answers `isShuttingDown()` with `false` and then rejects the schedule, and the executor comes from the channel, so there is no way in through the config. A test double for `EventExecutor` would do it if that is acceptable. ## Verification `mvnw clean verify` - BUILD SUCCESS, 1468 tests, 0 failures, 0 errors, 21 skipped. Error Prone, NullAway clean; `revapi` clean with the one scoped entry above. Caveat on the testing gate: `AGENTS.md` requires the build to run on JDK 11 and no JDK 11 is installed on this machine, so it was run on **JDK 17** (also in the CI matrix). The JDK 11 leg of CI on this PR is the real gate. Claude Code on behalf of @pavel-ptashyts 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
1 parent 793aae9 commit c1af3a4

12 files changed

Lines changed: 552 additions & 29 deletions

File tree

‎client/src/main/java/org/asynchttpclient/AsyncHttpClientConfig.java‎

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -114,6 +114,34 @@ public interface AsyncHttpClientConfig {
114114
*/
115115
Duration getRequestTimeout();
116116

117+
/**
118+
* Whether request and read timeouts are armed on an event loop rather than on {@link #getNettyTimer()}.
119+
* <p>
120+
* The timer is a hashed wheel: it fires on the first tick at or after the deadline, so a deadline near or
121+
* below {@link #getHashedWheelTimerTickDuration()} is rounded up to it, and one thread carries every expiry
122+
* for the whole client. An event loop instead schedules by deadline and derives its own select timeout from
123+
* the nearest one, so nothing is rounded up, and the loops share the load rather than funnelling it through
124+
* a single thread. Both effects matter most to short deadlines, where a tick is a large fraction of the
125+
* budget and a burst of expiries has no headroom to absorb.
126+
* <p>
127+
* The cost is where the expiry runs. On the timer it runs on the timer thread; on an event loop it runs on
128+
* an I/O thread, and so does whatever the caller chained onto the response future, because that future is
129+
* completed from there. Blocking an I/O thread stalls every connection it serves, so a caller enabling this
130+
* should hand its own work off with {@code handleAsync} or an {@code AsyncHandler} that does the same.
131+
* That is why this is opt-in rather than the default.
132+
* <p>
133+
* The loop is always the one that owns the exchange's channel. Until there is a channel -- while an address
134+
* is being resolved and a connection made -- the timer carries the timeout, and the exchange moves it onto
135+
* the loop once the connection succeeds.
136+
* <p>
137+
* The connection-pool cleaner stays on the timer either way.
138+
*
139+
* @return {@code true} to arm request and read timeouts on an event loop
140+
*/
141+
default boolean isUseEventLoopTimeouts() {
142+
return false;
143+
}
144+
117145
/**
118146
* Is HTTP redirect enabled
119147
*

‎client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -97,6 +97,7 @@
9797
import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultStrict302Handling;
9898
import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultTcpNoDelay;
9999
import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultThreadPoolName;
100+
import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultUseEventLoopTimeouts;
100101
import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultUseInsecureTrustManager;
101102
import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultUseLaxCookieEncoder;
102103
import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultUseNativeTransport;
@@ -157,6 +158,7 @@ public class DefaultAsyncHttpClientConfig implements AsyncHttpClientConfig {
157158
private final Duration connectTimeout;
158159
private final Duration requestTimeout;
159160
private final Duration readTimeout;
161+
private final boolean useEventLoopTimeouts;
160162
private final Duration shutdownQuietPeriod;
161163
private final Duration shutdownTimeout;
162164

@@ -258,6 +260,7 @@ private DefaultAsyncHttpClientConfig(// http
258260
Duration connectTimeout,
259261
Duration requestTimeout,
260262
Duration readTimeout,
263+
boolean useEventLoopTimeouts,
261264
Duration shutdownQuietPeriod,
262265
Duration shutdownTimeout,
263266

@@ -367,6 +370,7 @@ private DefaultAsyncHttpClientConfig(// http
367370
this.connectTimeout = connectTimeout;
368371
this.requestTimeout = requestTimeout;
369372
this.readTimeout = readTimeout;
373+
this.useEventLoopTimeouts = useEventLoopTimeouts;
370374
this.shutdownQuietPeriod = shutdownQuietPeriod;
371375
this.shutdownTimeout = shutdownTimeout;
372376

@@ -585,6 +589,11 @@ public Duration getReadTimeout() {
585589
return readTimeout;
586590
}
587591

592+
@Override
593+
public boolean isUseEventLoopTimeouts() {
594+
return useEventLoopTimeouts;
595+
}
596+
588597
@Override
589598
public Duration getShutdownQuietPeriod() {
590599
return shutdownQuietPeriod;
@@ -958,6 +967,7 @@ public static class Builder {
958967
private Duration connectTimeout = defaultConnectTimeout();
959968
private Duration requestTimeout = defaultRequestTimeout();
960969
private Duration readTimeout = defaultReadTimeout();
970+
private boolean useEventLoopTimeouts = defaultUseEventLoopTimeouts();
961971
private Duration shutdownQuietPeriod = defaultShutdownQuietPeriod();
962972
private Duration shutdownTimeout = defaultShutdownTimeout();
963973

@@ -1064,6 +1074,7 @@ public Builder(AsyncHttpClientConfig config) {
10641074
connectTimeout = config.getConnectTimeout();
10651075
requestTimeout = config.getRequestTimeout();
10661076
readTimeout = config.getReadTimeout();
1077+
useEventLoopTimeouts = config.isUseEventLoopTimeouts();
10671078
shutdownQuietPeriod = config.getShutdownQuietPeriod();
10681079
shutdownTimeout = config.getShutdownTimeout();
10691080

@@ -1355,6 +1366,17 @@ public Builder setReadTimeout(Duration readTimeout) {
13551366
return this;
13561367
}
13571368

1369+
/**
1370+
* @param useEventLoopTimeouts whether to arm request and read timeouts on an event loop instead of on
1371+
* the client's timer; see {@link AsyncHttpClientConfig#isUseEventLoopTimeouts()}
1372+
* for the trade-off this makes
1373+
* @return this
1374+
*/
1375+
public Builder setUseEventLoopTimeouts(boolean useEventLoopTimeouts) {
1376+
this.useEventLoopTimeouts = useEventLoopTimeouts;
1377+
return this;
1378+
}
1379+
13581380
public Builder setShutdownQuietPeriod(Duration shutdownQuietPeriod) {
13591381
this.shutdownQuietPeriod = shutdownQuietPeriod;
13601382
return this;
@@ -1764,6 +1786,7 @@ public DefaultAsyncHttpClientConfig build() {
17641786
connectTimeout,
17651787
requestTimeout,
17661788
readTimeout,
1789+
useEventLoopTimeouts,
17671790
shutdownQuietPeriod,
17681791
shutdownTimeout,
17691792
keepAlive,

‎client/src/main/java/org/asynchttpclient/config/AsyncHttpClientConfigDefaults.java‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@ public final class AsyncHttpClientConfigDefaults {
4141
public static final String CONNECTION_POOL_CLEANER_PERIOD_CONFIG = "connectionPoolCleanerPeriod";
4242
public static final String READ_TIMEOUT_CONFIG = "readTimeout";
4343
public static final String REQUEST_TIMEOUT_CONFIG = "requestTimeout";
44+
public static final String USE_EVENT_LOOP_TIMEOUTS_CONFIG = "useEventLoopTimeouts";
4445
public static final String CONNECTION_TTL_CONFIG = "connectionTtl";
4546
public static final String FOLLOW_REDIRECT_CONFIG = "followRedirect";
4647
public static final String MAX_REDIRECTS_CONFIG = "maxRedirects";
@@ -154,6 +155,10 @@ public static Duration defaultRequestTimeout() {
154155
return AsyncHttpClientConfigHelper.getAsyncHttpClientConfig().getDuration(ASYNC_CLIENT_CONFIG_ROOT + REQUEST_TIMEOUT_CONFIG);
155156
}
156157

158+
public static boolean defaultUseEventLoopTimeouts() {
159+
return AsyncHttpClientConfigHelper.getAsyncHttpClientConfig().getBoolean(ASYNC_CLIENT_CONFIG_ROOT + USE_EVENT_LOOP_TIMEOUTS_CONFIG);
160+
}
161+
157162
public static Duration defaultConnectionTtl() {
158163
return AsyncHttpClientConfigHelper.getAsyncHttpClientConfig().getDuration(ASYNC_CLIENT_CONFIG_ROOT + CONNECTION_TTL_CONFIG);
159164
}

‎client/src/main/java/org/asynchttpclient/netty/NettyResponseFuture.java‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -466,6 +466,11 @@ public void setTimeoutsHolder(TimeoutsHolder timeoutsHolder) {
466466
if (ref != null) {
467467
ref.cancel();
468468
}
469+
if (timeoutsHolder != null) {
470+
// Armed here rather than by the caller: a holder can run its timeout the moment it is armed, so it
471+
// has to be reachable from this future first, and no caller can then install one that never arms.
472+
timeoutsHolder.start();
473+
}
469474
}
470475

471476
public boolean isInAuth() {

‎client/src/main/java/org/asynchttpclient/netty/channel/NettyConnectListener.java‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -124,6 +124,11 @@ public void onSuccess(Channel channel, InetSocketAddress remoteAddress) {
124124
// mid-handshake could not close the socket, stranding it until handshakeTimeout (issue #2189).
125125
future.attachChannel(channel, false);
126126

127+
// The timeouts were armed before there was a channel to arm them on; hand them the one the exchange
128+
// ended up with. This listener runs on that channel's own loop, so the move needs no wakeup, and from
129+
// here on an expiry runs on the thread that would have to close the socket.
130+
timeoutsHolder.rehomeOn(channel.eventLoop());
131+
127132
Request request = future.getTargetRequest();
128133
Uri uri = request.getUri();
129134
// don't set a null resolved address - if the remoteAddress is null we keep

‎client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java‎

Lines changed: 33 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,7 @@
8282
import org.asynchttpclient.resolver.RequestHostnameResolver;
8383
import org.asynchttpclient.uri.Uri;
8484
import org.asynchttpclient.ws.WebSocketUpgradeHandler;
85+
import org.jetbrains.annotations.Nullable;
8586

8687
import org.slf4j.Logger;
8788
import org.slf4j.LoggerFactory;
@@ -401,15 +402,17 @@ private <T> ListenableFuture<T> sendRequestWithOpenChannel(NettyResponseFuture<T
401402
return future;
402403
}
403404

405+
future.setChannelState(ChannelState.POOLED);
406+
// Before the timeout is armed, not after: an expiry reaches the channel only through the future, and on
407+
// an event loop a short enough deadline can be delivered before the next statement would have run.
408+
future.attachChannel(channel, false);
409+
404410
SocketAddress channelRemoteAddress = channel.remoteAddress();
405411
if (channelRemoteAddress != null) {
406412
// otherwise, bad luck, the channel was closed, see bellow
407-
scheduleRequestTimeout(future, (InetSocketAddress) channelRemoteAddress);
413+
scheduleRequestTimeout(future, (InetSocketAddress) channelRemoteAddress, channel);
408414
}
409415

410-
future.setChannelState(ChannelState.POOLED);
411-
future.attachChannel(channel, false);
412-
413416
if (LOGGER.isDebugEnabled()) {
414417
HttpRequest httpRequest = future.getNettyRequest().getHttpRequest();
415418
LOGGER.debug("Using open Channel {} for {} '{}'", channel, httpRequest.method(), httpRequest.uri());
@@ -1080,12 +1083,36 @@ private static void configureTransferAdapter(AsyncHandler<?> handler, HttpReques
10801083

10811084
private void scheduleRequestTimeout(NettyResponseFuture<?> nettyResponseFuture,
10821085
InetSocketAddress originalRemoteAddress) {
1086+
scheduleRequestTimeout(nettyResponseFuture, originalRemoteAddress, null);
1087+
}
1088+
1089+
/**
1090+
* @param channel the channel the exchange will run on when it is already known, so the timeout can be armed
1091+
* on the loop that owns it. Null on the connect path: the timeout is armed before the channel
1092+
* exists, deliberately, so that it also bounds address resolution and the connect itself, and
1093+
* {@code TimeoutsHolder#rehomeOn} moves it onto the loop once there is one.
1094+
*/
1095+
private void scheduleRequestTimeout(NettyResponseFuture<?> nettyResponseFuture,
1096+
InetSocketAddress originalRemoteAddress,
1097+
@Nullable Channel channel) {
10831098
nettyResponseFuture.touch();
1084-
TimeoutsHolder timeoutsHolder = new TimeoutsHolder(nettyTimer, nettyResponseFuture, this, config,
1085-
originalRemoteAddress);
1099+
TimeoutsHolder timeoutsHolder = new TimeoutsHolder(nettyTimer, timeoutExecutor(channel), nettyResponseFuture,
1100+
this, config, originalRemoteAddress);
1101+
// Arms the timeout as a part of installing the holder, which is why the pooled path attaches the
1102+
// channel first: an expiry that lands immediately reaches the channel only through the future.
10861103
nettyResponseFuture.setTimeoutsHolder(timeoutsHolder);
10871104
}
10881105

1106+
/**
1107+
* The loop to arm an exchange's timeouts on, or null to leave them on the client's timer. Only ever the
1108+
* exchange's own channel's loop: any other loop would be woken by an entry it has no interest in, and the
1109+
* group's chooser hands out channels from the same counter, so drawing from it here would shift which loops
1110+
* connections land on.
1111+
*/
1112+
private @Nullable EventExecutor timeoutExecutor(@Nullable Channel channel) {
1113+
return config.isUseEventLoopTimeouts() && channel != null ? channel.eventLoop() : null;
1114+
}
1115+
10891116
private static void scheduleReadTimeout(NettyResponseFuture<?> nettyResponseFuture) {
10901117
TimeoutsHolder timeoutsHolder = nettyResponseFuture.getTimeoutsHolder();
10911118
if (timeoutsHolder != null) {

‎client/src/main/java/org/asynchttpclient/netty/timeout/ReadTimeoutTimerTask.java‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222

2323
import static org.asynchttpclient.util.DateUtils.unpreciseMillisTime;
2424

25-
public class ReadTimeoutTimerTask extends TimeoutTimerTask {
25+
public class ReadTimeoutTimerTask extends TimeoutTimerTask implements Runnable {
2626

2727
private final long readTimeout;
2828

@@ -31,6 +31,15 @@ public class ReadTimeoutTimerTask extends TimeoutTimerTask {
3131
this.readTimeout = readTimeout;
3232
}
3333

34+
/**
35+
* The event-loop entry point. Nothing below reads the {@link Timeout}, which is the timer's own handle on
36+
* this task and something an event loop has no equivalent of, so both entry points share one body.
37+
*/
38+
@Override
39+
public void run() {
40+
run(null);
41+
}
42+
3443
@Override
3544
public void run(Timeout timeout) {
3645
if (done.getAndSet(true) || requestSender.isClosed()) {

‎client/src/main/java/org/asynchttpclient/netty/timeout/RequestTimeoutTimerTask.java‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222

2323
import static org.asynchttpclient.util.DateUtils.unpreciseMillisTime;
2424

25-
public class RequestTimeoutTimerTask extends TimeoutTimerTask {
25+
public class RequestTimeoutTimerTask extends TimeoutTimerTask implements Runnable {
2626

2727
private final long requestTimeout;
2828

@@ -34,6 +34,15 @@ public class RequestTimeoutTimerTask extends TimeoutTimerTask {
3434
this.requestTimeout = requestTimeout;
3535
}
3636

37+
/**
38+
* The event-loop entry point. Nothing below reads the {@link Timeout}, which is the timer's own handle on
39+
* this task and something an event loop has no equivalent of, so both entry points share one body.
40+
*/
41+
@Override
42+
public void run() {
43+
run(null);
44+
}
45+
3746
@Override
3847
public void run(Timeout timeout) {
3948
if (done.getAndSet(true) || requestSender.isClosed()) {

‎client/src/main/java/org/asynchttpclient/netty/timeout/TimeoutTimerTask.java‎

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,16 +15,27 @@
1515
*/
1616
package org.asynchttpclient.netty.timeout;
1717

18+
import io.netty.util.Timeout;
1819
import io.netty.util.TimerTask;
20+
import io.netty.util.concurrent.ScheduledFuture;
1921
import org.asynchttpclient.netty.NettyResponseFuture;
2022
import org.asynchttpclient.netty.request.NettyRequestSender;
23+
import org.jetbrains.annotations.Nullable;
2124
import org.slf4j.Logger;
2225
import org.slf4j.LoggerFactory;
2326

2427
import java.net.InetSocketAddress;
28+
import java.util.concurrent.RejectedExecutionException;
2529
import java.util.concurrent.TimeoutException;
2630
import java.util.concurrent.atomic.AtomicBoolean;
2731

32+
/**
33+
* A timeout that can be armed either on a {@link io.netty.util.Timer} or on an event loop; which one an
34+
* exchange uses is {@link org.asynchttpclient.AsyncHttpClientConfig#isUseEventLoopTimeouts()}. An event loop
35+
* schedules {@link Runnable}s, so each subclass implements that as a second entry point into the same body.
36+
* This class stays a {@link TimerTask} alone: {@code Runnable#run} declares no checked exception, so a
37+
* {@code run()} here would have to catch what {@link TimerTask#run(Timeout)} declares and no subclass throws.
38+
*/
2839
public abstract class TimeoutTimerTask implements TimerTask {
2940

3041
private static final Logger LOGGER = LoggerFactory.getLogger(TimeoutTimerTask.class);
@@ -33,13 +44,67 @@ public abstract class TimeoutTimerTask implements TimerTask {
3344
protected final NettyRequestSender requestSender;
3445
final TimeoutsHolder timeoutsHolder;
3546
volatile NettyResponseFuture<?> nettyResponseFuture;
47+
// The scheduled entry this task is armed on, one field per scheduler so that a scheduler changing its
48+
// return type is a compile error rather than a cancellation that silently stops working. At most one is
49+
// ever set. Held here rather than in a wrapper so arming allocates nothing beyond what the scheduler needs.
50+
private volatile @Nullable Timeout timerHandle;
51+
private volatile @Nullable ScheduledFuture<?> loopHandle;
3652

3753
TimeoutTimerTask(NettyResponseFuture<?> nettyResponseFuture, NettyRequestSender requestSender, TimeoutsHolder timeoutsHolder) {
3854
this.nettyResponseFuture = nettyResponseFuture;
3955
this.requestSender = requestSender;
4056
this.timeoutsHolder = timeoutsHolder;
4157
}
4258

59+
void armedOn(Timeout handle) {
60+
// Each clears the other, so a handle left over from a previous arming cannot mask the live one and
61+
// leave its entry sitting in a scheduler, holding the future until a deadline nobody is waiting for.
62+
loopHandle = null;
63+
timerHandle = handle;
64+
}
65+
66+
void armedOn(ScheduledFuture<?> handle) {
67+
timerHandle = null;
68+
loopHandle = handle;
69+
}
70+
71+
/**
72+
* Cancels the scheduled entry this task was armed on, if any. Never interrupts: on the event-loop path the
73+
* task may be running on the very thread this is called from, and nothing in it answers interruption.
74+
*
75+
* @return whether an entry was taken back out of its scheduler before it could run
76+
*/
77+
boolean cancelArmed() {
78+
Timeout timer = timerHandle;
79+
if (timer != null) {
80+
timerHandle = null;
81+
return timer.cancel();
82+
}
83+
ScheduledFuture<?> scheduled = loopHandle;
84+
if (scheduled != null) {
85+
loopHandle = null;
86+
try {
87+
return scheduled.cancel(false);
88+
} catch (RejectedExecutionException e) {
89+
// Cancelling from off the loop enqueues the removal, which a loop that is already shutting down
90+
// rejects. The entry dies with the loop either way, and this runs under
91+
// ListenableFuture#cancel, which has never thrown for a client that is closing.
92+
LOGGER.debug("Event loop rejected a timeout cancellation", e);
93+
return false;
94+
}
95+
}
96+
return false;
97+
}
98+
99+
/**
100+
* Whether this task has been claimed, either by firing or by {@link #clean()}. Stands in for the
101+
* scheduler's own already-expired flag, which the two schedulers spell differently, and is if anything the
102+
* more precise of the two: it flips when {@code run} is entered rather than when the entry is marked.
103+
*/
104+
boolean isClaimed() {
105+
return done.get();
106+
}
107+
43108
void expire(String message, long time) {
44109
LOGGER.debug("{} for {} after {} ms", message, nettyResponseFuture, time);
45110
requestSender.abort(nettyResponseFuture.channel(), nettyResponseFuture, new TimeoutException(message));

0 commit comments

Comments
 (0)