diff --git a/client/src/main/java/org/asynchttpclient/AsyncHttpClientConfig.java b/client/src/main/java/org/asynchttpclient/AsyncHttpClientConfig.java index cae3900ee..94eebf180 100644 --- a/client/src/main/java/org/asynchttpclient/AsyncHttpClientConfig.java +++ b/client/src/main/java/org/asynchttpclient/AsyncHttpClientConfig.java @@ -165,6 +165,39 @@ public interface AsyncHttpClientConfig { @Nullable ThreadFactory getThreadFactory(); + /** + * Returns whether fallback resolution with the default blocking name resolver should be offloaded from + * Netty event-loop threads. Disabled by default. + * + * @return {@code true} if fallback name resolution should be offloaded + * @since 3.0.12 + */ + default boolean isFallbackNameResolverOffloadEnabled() { + return false; + } + + /** + * Returns the number of threads used to offload fallback name resolution. + * + * @return the configured fallback resolver thread count. Values less than or equal to {@code 0} use + * {@link #getIoThreadsCount()}. + * @since 3.0.12 + */ + default int getFallbackNameResolverOffloadThreadsCount() { + return -1; + } + + /** + * Returns the fallback name resolver worker queue capacity. + * + * @return the configured fallback resolver queue size. Values less than or equal to {@code 0} use the + * default queue size. + * @since 3.0.12 + */ + default int getFallbackNameResolverOffloadQueueSize() { + return 0; + } + /** * An instance of {@link ProxyServer} used by an {@link AsyncHttpClient} * diff --git a/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java b/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java index 467e3d9ad..47467865f 100644 --- a/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java +++ b/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java @@ -64,6 +64,9 @@ import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultExpiredCookieEvictionDelay; import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultFailedIpCooldownEnabled; import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultFailedIpCooldownPeriod; +import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultFallbackNameResolverOffloadEnabled; +import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultFallbackNameResolverOffloadQueueSize; +import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultFallbackNameResolverOffloadThreadsCount; import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultFilterInsecureCipherSuites; import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultFollowRedirect; import static org.asynchttpclient.config.AsyncHttpClientConfigDefaults.defaultHandshakeTimeout; @@ -226,6 +229,9 @@ public class DefaultAsyncHttpClientConfig implements AsyncHttpClientConfig { private final @Nullable Consumer wsAdditionalChannelInitializer; private final ResponseBodyPartFactory responseBodyPartFactory; private final int ioThreadsCount; + private final boolean fallbackNameResolverOffloadEnabled; + private final int fallbackNameResolverOffloadThreadsCount; + private final int fallbackNameResolverOffloadQueueSize; private final long hashedWheelTimerTickDuration; private final int hashedWheelTimerSize; @@ -330,6 +336,9 @@ private DefaultAsyncHttpClientConfig(// http @Nullable Consumer wsAdditionalChannelInitializer, ResponseBodyPartFactory responseBodyPartFactory, int ioThreadsCount, + boolean fallbackNameResolverOffloadEnabled, + int fallbackNameResolverOffloadThreadsCount, + int fallbackNameResolverOffloadQueueSize, long hashedWheelTimerTickDuration, int hashedWheelTimerSize) { @@ -449,6 +458,9 @@ private DefaultAsyncHttpClientConfig(// http this.wsAdditionalChannelInitializer = wsAdditionalChannelInitializer; this.responseBodyPartFactory = responseBodyPartFactory; this.ioThreadsCount = ioThreadsCount; + this.fallbackNameResolverOffloadEnabled = fallbackNameResolverOffloadEnabled; + this.fallbackNameResolverOffloadThreadsCount = fallbackNameResolverOffloadThreadsCount; + this.fallbackNameResolverOffloadQueueSize = fallbackNameResolverOffloadQueueSize; this.hashedWheelTimerTickDuration = hashedWheelTimerTickDuration; this.hashedWheelTimerSize = hashedWheelTimerSize; } @@ -907,6 +919,21 @@ public int getIoThreadsCount() { return ioThreadsCount; } + @Override + public boolean isFallbackNameResolverOffloadEnabled() { + return fallbackNameResolverOffloadEnabled; + } + + @Override + public int getFallbackNameResolverOffloadThreadsCount() { + return fallbackNameResolverOffloadThreadsCount; + } + + @Override + public int getFallbackNameResolverOffloadQueueSize() { + return fallbackNameResolverOffloadQueueSize; + } + /** * Builder for an {@link AsyncHttpClient} */ @@ -1016,6 +1043,9 @@ public static class Builder { private @Nullable Consumer wsAdditionalChannelInitializer; private ResponseBodyPartFactory responseBodyPartFactory = ResponseBodyPartFactory.EAGER; private int ioThreadsCount = defaultIoThreadsCount(); + private boolean fallbackNameResolverOffloadEnabled = defaultFallbackNameResolverOffloadEnabled(); + private int fallbackNameResolverOffloadThreadsCount = defaultFallbackNameResolverOffloadThreadsCount(); + private int fallbackNameResolverOffloadQueueSize = defaultFallbackNameResolverOffloadQueueSize(); private long hashedWheelTickDuration = defaultHashedWheelTimerTickDuration(); private int hashedWheelSize = defaultHashedWheelTimerSize(); @@ -1127,6 +1157,9 @@ public Builder(AsyncHttpClientConfig config) { wsAdditionalChannelInitializer = config.getWsAdditionalChannelInitializer(); responseBodyPartFactory = config.getResponseBodyPartFactory(); ioThreadsCount = config.getIoThreadsCount(); + fallbackNameResolverOffloadEnabled = config.isFallbackNameResolverOffloadEnabled(); + fallbackNameResolverOffloadThreadsCount = config.getFallbackNameResolverOffloadThreadsCount(); + fallbackNameResolverOffloadQueueSize = config.getFallbackNameResolverOffloadQueueSize(); hashedWheelTickDuration = config.getHashedWheelTimerTickDuration(); hashedWheelSize = config.getHashedWheelTimerSize(); } @@ -1688,6 +1721,21 @@ public Builder setIoThreadsCount(int ioThreadsCount) { return this; } + public Builder setFallbackNameResolverOffloadEnabled(boolean fallbackNameResolverOffloadEnabled) { + this.fallbackNameResolverOffloadEnabled = fallbackNameResolverOffloadEnabled; + return this; + } + + public Builder setFallbackNameResolverOffloadThreadsCount(int fallbackNameResolverOffloadThreadsCount) { + this.fallbackNameResolverOffloadThreadsCount = fallbackNameResolverOffloadThreadsCount; + return this; + } + + public Builder setFallbackNameResolverOffloadQueueSize(int fallbackNameResolverOffloadQueueSize) { + this.fallbackNameResolverOffloadQueueSize = fallbackNameResolverOffloadQueueSize; + return this; + } + private ProxyServerSelector resolveProxyServerSelector() { if (proxyServerSelector != null) { return proxyServerSelector; @@ -1793,6 +1841,9 @@ public DefaultAsyncHttpClientConfig build() { wsAdditionalChannelInitializer, responseBodyPartFactory, ioThreadsCount, + fallbackNameResolverOffloadEnabled, + fallbackNameResolverOffloadThreadsCount, + fallbackNameResolverOffloadQueueSize, hashedWheelTickDuration, hashedWheelSize); } diff --git a/client/src/main/java/org/asynchttpclient/config/AsyncHttpClientConfigDefaults.java b/client/src/main/java/org/asynchttpclient/config/AsyncHttpClientConfigDefaults.java index 550381fbd..bf07e101a 100644 --- a/client/src/main/java/org/asynchttpclient/config/AsyncHttpClientConfigDefaults.java +++ b/client/src/main/java/org/asynchttpclient/config/AsyncHttpClientConfigDefaults.java @@ -90,6 +90,9 @@ public final class AsyncHttpClientConfigDefaults { public static final String USE_NATIVE_TRANSPORT_CONFIG = "useNativeTransport"; public static final String USE_ONLY_EPOLL_NATIVE_TRANSPORT = "useOnlyEpollNativeTransport"; public static final String IO_THREADS_COUNT_CONFIG = "ioThreadsCount"; + public static final String FALLBACK_NAME_RESOLVER_OFFLOAD_ENABLED_CONFIG = "fallbackNameResolverOffloadEnabled"; + public static final String FALLBACK_NAME_RESOLVER_OFFLOAD_THREADS_COUNT_CONFIG = "fallbackNameResolverOffloadThreadsCount"; + public static final String FALLBACK_NAME_RESOLVER_OFFLOAD_QUEUE_SIZE_CONFIG = "fallbackNameResolverOffloadQueueSize"; public static final String HASHED_WHEEL_TIMER_TICK_DURATION = "hashedWheelTimerTickDuration"; public static final String HASHED_WHEEL_TIMER_SIZE = "hashedWheelTimerSize"; public static final String EXPIRED_COOKIE_EVICTION_DELAY = "expiredCookieEvictionDelay"; @@ -347,6 +350,21 @@ public static int defaultIoThreadsCount() { return threads; } + public static boolean defaultFallbackNameResolverOffloadEnabled() { + return AsyncHttpClientConfigHelper.getAsyncHttpClientConfig().getBoolean( + ASYNC_CLIENT_CONFIG_ROOT + FALLBACK_NAME_RESOLVER_OFFLOAD_ENABLED_CONFIG); + } + + public static int defaultFallbackNameResolverOffloadThreadsCount() { + return AsyncHttpClientConfigHelper.getAsyncHttpClientConfig().getInt( + ASYNC_CLIENT_CONFIG_ROOT + FALLBACK_NAME_RESOLVER_OFFLOAD_THREADS_COUNT_CONFIG); + } + + public static int defaultFallbackNameResolverOffloadQueueSize() { + return AsyncHttpClientConfigHelper.getAsyncHttpClientConfig().getInt( + ASYNC_CLIENT_CONFIG_ROOT + FALLBACK_NAME_RESOLVER_OFFLOAD_QUEUE_SIZE_CONFIG); + } + public static int defaultHashedWheelTimerTickDuration() { return AsyncHttpClientConfigHelper.getAsyncHttpClientConfig().getInt(ASYNC_CLIENT_CONFIG_ROOT + HASHED_WHEEL_TIMER_TICK_DURATION); } diff --git a/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java b/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java index d5db73d54..821a9b844 100755 --- a/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java +++ b/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java @@ -59,6 +59,7 @@ import io.netty.resolver.NameResolver; import io.netty.util.Timer; import io.netty.util.concurrent.DefaultThreadFactory; +import io.netty.util.concurrent.EventExecutor; import io.netty.util.concurrent.Future; import io.netty.util.concurrent.GlobalEventExecutor; import io.netty.util.concurrent.ImmediateEventExecutor; @@ -81,6 +82,7 @@ import org.asynchttpclient.netty.handler.HttpHandler; import org.asynchttpclient.netty.handler.WebSocketHandler; import org.asynchttpclient.netty.request.NettyRequestSender; +import org.asynchttpclient.netty.resolver.NameResolverOffload; import org.asynchttpclient.netty.ssl.DefaultSslEngineFactory; import org.asynchttpclient.proxy.ProxyServer; import org.asynchttpclient.proxy.ProxyType; @@ -138,6 +140,7 @@ public class ChannelManager { private final Map.Entry, Object>[] channelOptions; private final long handshakeTimeout; private final @Nullable AddressResolverGroup addressResolverGroup; + private final NameResolverOffload nameResolverOffload; private final ChannelPool channelPool; private final ChannelGroup openChannels; @@ -176,6 +179,7 @@ private boolean isInstanceof(Object object, String name) { public ChannelManager(final AsyncHttpClientConfig config, Timer nettyTimer) { this.config = config; + nameResolverOffload = new NameResolverOffload(config); sslEngineFactory = config.getSslEngineFactory() != null ? config.getSslEngineFactory() : new DefaultSslEngineFactory(); try { @@ -695,6 +699,7 @@ private void doClose() { } public void close() { + nameResolverOffload.close(); // Fail any requests parked waiting for a sibling HTTP/2 connection to register (see the // http2ConnectionWaiters field): the client is closing, so no connection will arrive and their // request-timeout backstop is not scheduled yet. Do this synchronously up front — doClose() only @@ -886,7 +891,7 @@ public Future getBootstrap(Uri uri, NameResolver nameRes } }); } else { - nameResolver.resolve(proxy.getHost()).addListener((Future whenProxyAddress) -> { + resolveProxyHost(nameResolver, proxy.getHost()).addListener((Future whenProxyAddress) -> { if (whenProxyAddress.isSuccess()) { InetSocketAddress proxyAddress = new InetSocketAddress(whenProxyAddress.get(), proxy.getPort()); configureSocksBootstrap(socksBootstrap, httpBootstrapHandler, proxyAddress, proxy, promise); @@ -907,6 +912,33 @@ public Future getBootstrap(Uri uri, NameResolver nameRes return promise; } + private Future resolveProxyHost(NameResolver nameResolver, String host) { + EventExecutor eventLoop = currentEventLoop(); + if (eventLoop == null || !nameResolverOffload.shouldOffload(nameResolver)) { + return nameResolver.resolve(host); + } + + Promise promise = ImmediateEventExecutor.INSTANCE.newPromise(); + nameResolverOffload.execute(eventLoop, promise, () -> + nameResolver.resolve(host).addListener((Future whenResolved) -> { + if (whenResolved.isSuccess()) { + nameResolverOffload.completeSuccess(eventLoop, promise, whenResolved.getNow()); + } else { + nameResolverOffload.completeFailure(eventLoop, promise, whenResolved.cause()); + } + })); + return promise; + } + + private @Nullable EventExecutor currentEventLoop() { + for (EventExecutor executor : eventLoopGroup) { + if (executor.inEventLoop()) { + return executor; + } + } + return null; + } + private void configureSocksBootstrap(Bootstrap socksBootstrap, ChannelHandler httpBootstrapHandler, InetSocketAddress proxyAddress, ProxyServer proxy, Promise promise) { socksBootstrap.handler(new ChannelInitializer() { @@ -1121,6 +1153,18 @@ public EventLoopGroup getEventLoopGroup() { return eventLoopGroup; } + /** + * Returns the client-owned fallback resolver offload support. + * This method is public only for cross-package request-sender wiring; applications should configure + * resolver offloading through {@link AsyncHttpClientConfig}. + * + * @return fallback resolver offload support + * @since 3.0.12 + */ + public NameResolverOffload getNameResolverOffload() { + return nameResolverOffload; + } + /** * Return the {@link AddressResolverGroup} used for async DNS resolution, or {@code null} * if per-request name resolvers should be used (legacy behavior). diff --git a/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java b/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java index 36af9019b..b1fb3d97b 100755 --- a/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java +++ b/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java @@ -37,6 +37,7 @@ import io.netty.handler.codec.http2.Http2StreamChannelBootstrap; import io.netty.resolver.AddressResolver; import io.netty.resolver.AddressResolverGroup; +import io.netty.resolver.NameResolver; import io.netty.util.AsciiString; import io.netty.util.Timeout; import io.netty.util.Timer; @@ -75,6 +76,7 @@ import org.asynchttpclient.netty.handler.Http2ContentDecompressor; import org.asynchttpclient.netty.request.body.NettyBody; import org.asynchttpclient.netty.request.body.NettyDirectBody; +import org.asynchttpclient.netty.resolver.NameResolverOffload; import org.asynchttpclient.netty.timeout.TimeoutsHolder; import org.asynchttpclient.proxy.ProxyServer; import org.asynchttpclient.proxy.ProxyType; @@ -82,6 +84,7 @@ import org.asynchttpclient.uri.Uri; import org.asynchttpclient.ws.WebSocketUpgradeHandler; +import org.jetbrains.annotations.Nullable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -91,10 +94,12 @@ import java.net.SocketAddress; import java.net.UnknownHostException; import java.time.Duration; +import java.util.ArrayList; import java.util.Iterator; import java.util.List; import java.util.Locale; import java.util.Map; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Consumer; @@ -121,11 +126,13 @@ public final class NettyRequestSender { private final AsyncHttpClientState clientState; private final NettyRequestFactory requestFactory; private final RoundRobinAddressSelector rrSelector = new RoundRobinAddressSelector(); + private final NameResolverOffload nameResolverOffload; // Deprioritizes a recently-failed IP when ordering a direct connection's resolved addresses, in any // LoadBalance mode. Null when the failed-IP cooldown is disabled; call sites gate on ipCooldown != null. private final FailedIpCooldownHolder ipCooldown; - public NettyRequestSender(AsyncHttpClientConfig config, ChannelManager channelManager, Timer nettyTimer, AsyncHttpClientState clientState) { + public NettyRequestSender(AsyncHttpClientConfig config, ChannelManager channelManager, Timer nettyTimer, + AsyncHttpClientState clientState) { this.config = config; this.channelManager = channelManager; connectionSemaphore = config.getConnectionSemaphoreFactory() == null @@ -133,6 +140,7 @@ public NettyRequestSender(AsyncHttpClientConfig config, ChannelManager channelMa : config.getConnectionSemaphoreFactory().newConnectionSemaphore(config); this.nettyTimer = nettyTimer; this.clientState = clientState; + nameResolverOffload = channelManager.getNameResolverOffload(); requestFactory = new NettyRequestFactory(config); // Guard the period against a custom AsyncHttpClientConfig that enables the cooldown but returns a // null period: leave the cooldown off rather than NPE while constructing the client. @@ -596,7 +604,86 @@ private Future> resolveHostname(Request request, InetSoc AddressResolver resolver = group.getResolver(channelManager.getEventLoopGroup().next()); return RequestHostnameResolver.INSTANCE.resolve(resolver, unresolvedRemoteAddress, asyncHandler); } - return RequestHostnameResolver.INSTANCE.resolve(request.getNameResolver(), unresolvedRemoteAddress, asyncHandler); + return resolveWithRequestNameResolver(request.getNameResolver(), unresolvedRemoteAddress, asyncHandler); + } + + private Future> resolveWithRequestNameResolver(NameResolver nameResolver, + InetSocketAddress unresolvedRemoteAddress, + AsyncHandler asyncHandler) { + EventExecutor eventLoop = currentEventLoop(); + if (eventLoop == null || !nameResolverOffload.shouldOffload(nameResolver)) { + return RequestHostnameResolver.INSTANCE.resolve(nameResolver, unresolvedRemoteAddress, asyncHandler); + } + + String hostname = unresolvedRemoteAddress.getHostString(); + int port = unresolvedRemoteAddress.getPort(); + Promise> promise = ImmediateEventExecutor.INSTANCE.newPromise(); + try { + asyncHandler.onHostnameResolutionAttempt(hostname); + } catch (Exception e) { + LOGGER.error("onHostnameResolutionAttempt crashed", e); + promise.tryFailure(e); + return promise; + } + + nameResolverOffload.execute(eventLoop, promise, () -> + nameResolver.resolveAll(hostname).addListener((Future> whenResolved) -> { + if (whenResolved.isSuccess()) { + completeResolutionSuccessOnEventLoop(eventLoop, promise, hostname, port, asyncHandler, + whenResolved.getNow()); + } else { + completeResolutionFailureOnEventLoop(eventLoop, promise, hostname, asyncHandler, + whenResolved.cause()); + } + })); + return promise; + } + + private static void completeResolutionSuccessOnEventLoop(EventExecutor eventLoop, + Promise> promise, + String hostname, + int port, + AsyncHandler asyncHandler, + List addresses) { + try { + eventLoop.execute(() -> { + ArrayList socketAddresses = new ArrayList<>(addresses.size()); + for (InetAddress address : addresses) { + socketAddresses.add(new InetSocketAddress(address, port)); + } + try { + asyncHandler.onHostnameResolutionSuccess(hostname, socketAddresses); + } catch (Exception e) { + LOGGER.error("onHostnameResolutionSuccess crashed", e); + promise.tryFailure(e); + return; + } + promise.trySuccess(socketAddresses); + }); + } catch (RejectedExecutionException e) { + promise.tryFailure(e); + } + } + + private static void completeResolutionFailureOnEventLoop(EventExecutor eventLoop, + Promise promise, + String hostname, + AsyncHandler asyncHandler, + Throwable cause) { + try { + eventLoop.execute(() -> { + try { + asyncHandler.onHostnameResolutionFailure(hostname, cause); + } catch (Exception e) { + LOGGER.error("onHostnameResolutionFailure crashed", e); + promise.tryFailure(e); + return; + } + promise.tryFailure(cause); + }); + } catch (RejectedExecutionException e) { + promise.tryFailure(e); + } } private NettyResponseFuture newNettyResponseFuture(Request request, AsyncHandler asyncHandler, NettyRequest nettyRequest, ProxyServer proxyServer) { @@ -1342,12 +1429,16 @@ private Channel pollHttp2(Object h2Key) { } private boolean isOnEventLoop() { + return currentEventLoop() != null; + } + + private @Nullable EventExecutor currentEventLoop() { for (EventExecutor executor : channelManager.getEventLoopGroup()) { if (executor.inEventLoop()) { - return true; + return executor; } } - return false; + return null; } private Channel pollPooledChannel(NettyResponseFuture future, Request request, ProxyServer proxy, AsyncHandler asyncHandler) { diff --git a/client/src/main/java/org/asynchttpclient/netty/resolver/NameResolverOffload.java b/client/src/main/java/org/asynchttpclient/netty/resolver/NameResolverOffload.java new file mode 100644 index 000000000..adb1e4cfe --- /dev/null +++ b/client/src/main/java/org/asynchttpclient/netty/resolver/NameResolverOffload.java @@ -0,0 +1,196 @@ +/* + * 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.resolver; + +import io.netty.resolver.NameResolver; +import io.netty.util.concurrent.DefaultThreadFactory; +import io.netty.util.concurrent.EventExecutor; +import io.netty.util.concurrent.Promise; +import org.asynchttpclient.AsyncHttpClientConfig; +import org.jetbrains.annotations.Nullable; + +import java.net.InetAddress; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import static java.util.Objects.requireNonNull; +import static org.asynchttpclient.RequestBuilderBase.DEFAULT_NAME_RESOLVER; + +/** + * Offloads the default blocking fallback name resolver and tracks its pending promises for client shutdown. + * This type is public only for cross-package client wiring; applications should configure resolver offloading + * through {@link AsyncHttpClientConfig}. + * + * @since 3.0.12 + */ +public final class NameResolverOffload implements AutoCloseable { + + private final @Nullable ExecutorService executor; + private final Map, EventExecutor> pending = new ConcurrentHashMap<>(); + private final AtomicBoolean closed = new AtomicBoolean(); + + /** + * Creates the client-owned fallback resolver offload support. + * + * @param config client configuration + */ + public NameResolverOffload(AsyncHttpClientConfig config) { + requireNonNull(config, "config"); + if (!config.isFallbackNameResolverOffloadEnabled()) { + executor = null; + return; + } + + int threads = configuredThreads(config); + int queueSize = configuredQueueSize(config, threads); + ThreadFactory threadFactory = config.getThreadFactory() != null + ? config.getThreadFactory() + : new DefaultThreadFactory(config.getThreadPoolName() + "-resolver"); + executor = new ThreadPoolExecutor(threads, threads, 0L, TimeUnit.MILLISECONDS, + new LinkedBlockingQueue<>(queueSize), threadFactory); + } + + /** + * Returns whether the resolver is the default blocking resolver and offloading is enabled. + * + * @param nameResolver resolver selected for the request + * @return {@code true} when the resolution should be offloaded + */ + public boolean shouldOffload(NameResolver nameResolver) { + return executor != null && nameResolver == DEFAULT_NAME_RESOLVER; + } + + /** + * Submits fallback resolution work and associates it with its request promise. + * + * @param eventLoop event loop that owns the request flow + * @param promise request promise to fail on rejection or shutdown + * @param task fallback resolution work + */ + public void execute(EventExecutor eventLoop, Promise promise, Runnable task) { + requireNonNull(eventLoop, "eventLoop"); + requireNonNull(promise, "promise"); + requireNonNull(task, "task"); + + ExecutorService resolverExecutor = executor; + if (resolverExecutor == null) { + completeFailure(eventLoop, promise, + new RejectedExecutionException("Fallback name resolver offload is disabled")); + return; + } + + pending.put(promise, eventLoop); + promise.addListener(ignored -> pending.remove(promise)); + + if (closed.get()) { + completeFailure(eventLoop, promise, closedFailure()); + return; + } + + try { + resolverExecutor.execute(() -> { + if (promise.isDone()) { + return; + } + try { + task.run(); + } catch (RuntimeException e) { + completeFailure(eventLoop, promise, e); + } + }); + } catch (RejectedExecutionException e) { + completeFailure(eventLoop, promise, e); + } + } + + /** + * Completes a resolution promise on its request event loop. + * + * @param eventLoop event loop that owns the request flow + * @param promise request promise + * @param result resolution result + * @param result type + */ + public void completeSuccess(EventExecutor eventLoop, Promise promise, T result) { + if (eventLoop.inEventLoop()) { + promise.trySuccess(result); + return; + } + try { + eventLoop.execute(() -> promise.trySuccess(result)); + } catch (RejectedExecutionException e) { + promise.tryFailure(e); + } + } + + /** + * Fails a resolution promise on its request event loop. + * + * @param eventLoop event loop that owns the request flow + * @param promise request promise + * @param cause resolution failure + */ + public void completeFailure(EventExecutor eventLoop, Promise promise, Throwable cause) { + if (eventLoop.inEventLoop()) { + promise.tryFailure(cause); + return; + } + try { + eventLoop.execute(() -> promise.tryFailure(cause)); + } catch (RejectedExecutionException e) { + promise.tryFailure(cause); + } + } + + @Override + public void close() { + ExecutorService resolverExecutor = executor; + if (resolverExecutor == null || !closed.compareAndSet(false, true)) { + return; + } + + resolverExecutor.shutdownNow(); + RejectedExecutionException failure = closedFailure(); + pending.forEach((promise, eventLoop) -> completeFailure(eventLoop, promise, failure)); + } + + private static int configuredThreads(AsyncHttpClientConfig config) { + int threads = config.getFallbackNameResolverOffloadThreadsCount(); + if (threads <= 0) { + threads = config.getIoThreadsCount(); + } + return Math.max(1, threads); + } + + private static int configuredQueueSize(AsyncHttpClientConfig config, int threads) { + int queueSize = config.getFallbackNameResolverOffloadQueueSize(); + if (queueSize <= 0) { + queueSize = threads > Integer.MAX_VALUE / 16 ? Integer.MAX_VALUE : Math.max(1024, threads * 16); + } + return Math.max(1, queueSize); + } + + private static RejectedExecutionException closedFailure() { + return new RejectedExecutionException("Fallback name resolver offload is closed"); + } +} diff --git a/client/src/main/resources/org/asynchttpclient/config/ahc-default.properties b/client/src/main/resources/org/asynchttpclient/config/ahc-default.properties index 5df97add5..0e6b5d05f 100644 --- a/client/src/main/resources/org/asynchttpclient/config/ahc-default.properties +++ b/client/src/main/resources/org/asynchttpclient/config/ahc-default.properties @@ -55,6 +55,9 @@ org.asynchttpclient.shutdownTimeout=PT15S org.asynchttpclient.useNativeTransport=false org.asynchttpclient.useOnlyEpollNativeTransport=false org.asynchttpclient.ioThreadsCount=-1 +org.asynchttpclient.fallbackNameResolverOffloadEnabled=false +org.asynchttpclient.fallbackNameResolverOffloadThreadsCount=-1 +org.asynchttpclient.fallbackNameResolverOffloadQueueSize=0 org.asynchttpclient.hashedWheelTimerTickDuration=100 org.asynchttpclient.hashedWheelTimerSize=512 org.asynchttpclient.expiredCookieEvictionDelay=30000 diff --git a/client/src/test/java/org/asynchttpclient/AddressResolverGroupTest.java b/client/src/test/java/org/asynchttpclient/AddressResolverGroupTest.java index 41de0c7cf..0f8815a9c 100644 --- a/client/src/test/java/org/asynchttpclient/AddressResolverGroupTest.java +++ b/client/src/test/java/org/asynchttpclient/AddressResolverGroupTest.java @@ -16,27 +16,43 @@ package org.asynchttpclient; import io.github.artsok.RepeatedIfExceptionsTest; +import io.netty.channel.EventLoopGroup; import io.netty.channel.socket.nio.NioDatagramChannel; +import io.netty.resolver.InetNameResolver; import io.netty.resolver.dns.DnsAddressResolverGroup; import io.netty.resolver.dns.DnsServerAddressStreamProviders; +import io.netty.util.concurrent.EventExecutor; +import io.netty.util.concurrent.ImmediateEventExecutor; +import io.netty.util.concurrent.Promise; +import org.asynchttpclient.proxy.ProxyServer; +import org.asynchttpclient.proxy.ProxyType; import org.asynchttpclient.test.EventCollectingHandler; import org.asynchttpclient.testserver.HttpServer; import org.asynchttpclient.testserver.HttpTest; +import org.asynchttpclient.uri.Uri; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Tag; +import java.net.InetAddress; import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import static java.util.concurrent.TimeUnit.SECONDS; import static org.asynchttpclient.Dsl.asyncHttpClient; import static org.asynchttpclient.Dsl.config; import static org.asynchttpclient.Dsl.get; +import static org.asynchttpclient.RequestBuilderBase.DEFAULT_NAME_RESOLVER; import static org.asynchttpclient.test.TestUtils.isExternalNetworkAvailable; import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -47,6 +63,9 @@ public class AddressResolverGroupTest extends HttpTest { private static final String GOOGLE_URL = "https://www.google.com/"; private static final String EXAMPLE_URL = "https://www.example.com/"; + private static final String INITIAL_HOST = "initial.test"; + private static final String REDIRECT_HOST = "redirect.test"; + private static final String SOCKS_HOST = "socks.test"; private HttpServer server; @@ -118,6 +137,126 @@ public void defaultConfigDoesNotSetAddressResolverGroup() { "Default config should not have an AddressResolverGroup"); } + @RepeatedIfExceptionsTest(repeats = 5) + public void customNameResolverOnRedirectKeepsEventLoopContext() throws Throwable { + server.enqueueRedirect(302, "http://" + REDIRECT_HOST + ':' + server.getHttpPort() + "/target"); + server.enqueueOk(); + + EventLoopProbeNameResolver resolver = new EventLoopProbeNameResolver(); + try (DefaultAsyncHttpClient client = new DefaultAsyncHttpClient(config().setFollowRedirect(true).build())) { + resolver.setEventLoopGroup(client.channelManager().getEventLoopGroup()); + + Response response = client.executeRequest(get("http://" + INITIAL_HOST + ':' + server.getHttpPort() + "/start") + .setNameResolver(resolver) + .build()) + .get(5, SECONDS); + + assertEquals(200, response.getStatusCode()); + } + + assertTrue(resolver.redirectResolutionAttempted.get(), "redirect host should have been resolved"); + assertTrue(resolver.redirectResolutionOnEventLoop.get(), "custom resolver should retain its calling context"); + } + + @RepeatedIfExceptionsTest(repeats = 5) + public void customSocksProxyResolverKeepsEventLoopContext() throws Throwable { + EventLoopProbeNameResolver resolver = new EventLoopProbeNameResolver(); + ProxyServer proxy = new ProxyServer.Builder(SOCKS_HOST, 1080) + .setProxyType(ProxyType.SOCKS_V5) + .build(); + + try (DefaultAsyncHttpClient client = new DefaultAsyncHttpClient(config().build())) { + EventLoopGroup eventLoopGroup = client.channelManager().getEventLoopGroup(); + resolver.setEventLoopGroup(eventLoopGroup); + + CountDownLatch complete = new CountDownLatch(1); + AtomicReference failure = new AtomicReference<>(); + eventLoopGroup.next().execute(() -> client.channelManager() + .getBootstrap(Uri.create("http://target.test/"), resolver, proxy) + .addListener(bootstrap -> { + if (!bootstrap.isSuccess()) { + failure.set(bootstrap.cause()); + } + complete.countDown(); + })); + + assertTrue(complete.await(5, SECONDS), "SOCKS bootstrap resolution should complete"); + assertNull(failure.get(), "SOCKS bootstrap should resolve the proxy host"); + } + + assertTrue(resolver.socksResolutionAttempted.get(), "SOCKS proxy host should have been resolved"); + assertTrue(resolver.socksResolutionOnEventLoop.get(), "custom resolver should retain its calling context"); + } + + @RepeatedIfExceptionsTest(repeats = 5) + public void defaultSocksProxyResolverUsesOffloadPool() throws Throwable { + String poolName = "ahc-socks-resolver-enabled"; + ProxyServer proxy = new ProxyServer.Builder("localhost", 1080) + .setProxyType(ProxyType.SOCKS_V5) + .build(); + + try (DefaultAsyncHttpClient client = new DefaultAsyncHttpClient(config() + .setThreadPoolName(poolName) + .setFallbackNameResolverOffloadEnabled(true) + .build())) { + EventLoopGroup eventLoopGroup = client.channelManager().getEventLoopGroup(); + CountDownLatch complete = new CountDownLatch(1); + AtomicReference failure = new AtomicReference<>(); + eventLoopGroup.next().execute(() -> client.channelManager() + .getBootstrap(Uri.create("http://target.test/"), DEFAULT_NAME_RESOLVER, proxy) + .addListener(bootstrap -> { + if (!bootstrap.isSuccess()) { + failure.set(bootstrap.cause()); + } + complete.countDown(); + })); + + assertTrue(complete.await(5, SECONDS), "SOCKS bootstrap resolution should complete"); + assertNull(failure.get(), "SOCKS bootstrap should resolve the proxy host"); + assertTrue(hasLiveThreadNamed(poolName + "-resolver"), + "SOCKS fallback DNS should use the resolver offload pool"); + } + } + + @RepeatedIfExceptionsTest(repeats = 5) + public void defaultNameResolverOnRedirectUsesOffloadPool() throws Throwable { + server.enqueueRedirect(302, "http://localhost:" + server.getHttpPort() + "/target"); + server.enqueueOk(); + + String poolName = "ahc-fallback-resolver-enabled"; + try (AsyncHttpClient client = asyncHttpClient(config() + .setFollowRedirect(true) + .setThreadPoolName(poolName) + .setFallbackNameResolverOffloadEnabled(true))) { + Response response = client.prepareGet("http://127.0.0.1:" + server.getHttpPort() + "/start") + .execute() + .get(5, SECONDS); + + assertEquals(200, response.getStatusCode()); + assertTrue(hasLiveThreadNamed(poolName + "-resolver"), + "redirect fallback DNS should use the resolver offload pool"); + } + } + + @RepeatedIfExceptionsTest(repeats = 5) + public void defaultNameResolverOnRedirectKeepsInlineBehavior() throws Throwable { + server.enqueueRedirect(302, "http://localhost:" + server.getHttpPort() + "/target"); + server.enqueueOk(); + + String poolName = "ahc-fallback-resolver-disabled"; + try (AsyncHttpClient client = asyncHttpClient(config() + .setFollowRedirect(true) + .setThreadPoolName(poolName))) { + Response response = client.prepareGet("http://127.0.0.1:" + server.getHttpPort() + "/start") + .execute() + .get(5, SECONDS); + + assertEquals(200, response.getStatusCode()); + assertFalse(hasLiveThreadNamed(poolName + "-resolver"), + "disabled offload must not create a resolver worker"); + } + } + @RepeatedIfExceptionsTest(repeats = 5) public void unknownHostWithDnsResolverGroupFails() throws Throwable { DnsAddressResolverGroup resolverGroup = new DnsAddressResolverGroup( @@ -172,4 +311,66 @@ public void resolveMultipleRealDomainsWithDnsResolverGroup() throws Throwable { "Expected successful HTTP status for example.com but got " + response2.getStatusCode()); } } + + private static final class EventLoopProbeNameResolver extends InetNameResolver { + private final AtomicReference eventLoopGroup = new AtomicReference<>(); + private final AtomicBoolean redirectResolutionAttempted = new AtomicBoolean(); + private final AtomicBoolean redirectResolutionOnEventLoop = new AtomicBoolean(); + private final AtomicBoolean socksResolutionAttempted = new AtomicBoolean(); + private final AtomicBoolean socksResolutionOnEventLoop = new AtomicBoolean(); + + EventLoopProbeNameResolver() { + super(ImmediateEventExecutor.INSTANCE); + } + + void setEventLoopGroup(EventLoopGroup eventLoopGroup) { + this.eventLoopGroup.set(eventLoopGroup); + } + + @Override + protected void doResolve(String inetHost, Promise promise) { + recordResolution(inetHost); + promise.setSuccess(InetAddress.getLoopbackAddress()); + } + + @Override + protected void doResolveAll(String inetHost, Promise> promise) { + recordResolution(inetHost); + promise.setSuccess(Collections.singletonList(InetAddress.getLoopbackAddress())); + } + + private void recordResolution(String inetHost) { + boolean onEventLoop = isOnEventLoop(); + if (REDIRECT_HOST.equals(inetHost)) { + redirectResolutionAttempted.set(true); + redirectResolutionOnEventLoop.set(onEventLoop); + } else if (SOCKS_HOST.equals(inetHost)) { + socksResolutionAttempted.set(true); + socksResolutionOnEventLoop.set(onEventLoop); + } + } + + private boolean isOnEventLoop() { + EventLoopGroup group = eventLoopGroup.get(); + if (group == null) { + return false; + } + Thread currentThread = Thread.currentThread(); + for (EventExecutor executor : group) { + if (executor.inEventLoop(currentThread)) { + return true; + } + } + return false; + } + } + + private static boolean hasLiveThreadNamed(String namePart) { + for (Thread thread : Thread.getAllStackTraces().keySet()) { + if (thread.isAlive() && thread.getName().contains(namePart)) { + return true; + } + } + return false; + } } diff --git a/client/src/test/java/org/asynchttpclient/AsyncHttpClientDefaultsTest.java b/client/src/test/java/org/asynchttpclient/AsyncHttpClientDefaultsTest.java index d125a9fa4..d4892c2db 100644 --- a/client/src/test/java/org/asynchttpclient/AsyncHttpClientDefaultsTest.java +++ b/client/src/test/java/org/asynchttpclient/AsyncHttpClientDefaultsTest.java @@ -144,6 +144,28 @@ public void testDefaultUseInsecureTrustManager() { testBooleanSystemProperty("useInsecureTrustManager", "defaultUseInsecureTrustManager", "false"); } + @RepeatedIfExceptionsTest(repeats = 5) + public void testDefaultFallbackNameResolverOffloadEnabled() { + assertFalse(AsyncHttpClientConfigDefaults.defaultFallbackNameResolverOffloadEnabled()); + testBooleanSystemProperty("fallbackNameResolverOffloadEnabled", + "defaultFallbackNameResolverOffloadEnabled", "true"); + AsyncHttpClientConfigHelper.reloadProperties(); + } + + @RepeatedIfExceptionsTest(repeats = 5) + public void testDefaultFallbackNameResolverOffloadThreadsCount() { + assertEquals(-1, AsyncHttpClientConfigDefaults.defaultFallbackNameResolverOffloadThreadsCount()); + testIntegerSystemProperty("fallbackNameResolverOffloadThreadsCount", + "defaultFallbackNameResolverOffloadThreadsCount", "2"); + } + + @RepeatedIfExceptionsTest(repeats = 5) + public void testDefaultFallbackNameResolverOffloadQueueSize() { + assertEquals(0, AsyncHttpClientConfigDefaults.defaultFallbackNameResolverOffloadQueueSize()); + testIntegerSystemProperty("fallbackNameResolverOffloadQueueSize", + "defaultFallbackNameResolverOffloadQueueSize", "17"); + } + @RepeatedIfExceptionsTest(repeats = 5) public void testDefaultHashedWheelTimerTickDuration() { assertEquals(AsyncHttpClientConfigDefaults.defaultHashedWheelTimerTickDuration(), 100); diff --git a/client/src/test/java/org/asynchttpclient/netty/resolver/NameResolverOffloadTest.java b/client/src/test/java/org/asynchttpclient/netty/resolver/NameResolverOffloadTest.java new file mode 100644 index 000000000..92596ad39 --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/netty/resolver/NameResolverOffloadTest.java @@ -0,0 +1,125 @@ +/* + * 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.resolver; + +import io.netty.channel.DefaultEventLoop; +import io.netty.resolver.DefaultNameResolver; +import io.netty.util.concurrent.ImmediateEventExecutor; +import io.netty.util.concurrent.Promise; +import org.asynchttpclient.DefaultAsyncHttpClientConfig; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.atomic.AtomicReference; + +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.asynchttpclient.RequestBuilderBase.DEFAULT_NAME_RESOLVER; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class NameResolverOffloadTest { + + private final DefaultEventLoop eventLoop = new DefaultEventLoop(); + + @AfterEach + void closeEventLoop() { + eventLoop.shutdownGracefully(0, 5, SECONDS).syncUninterruptibly(); + } + + @Test + void offloadsDefaultResolverToConfiguredWorker() throws Exception { + DefaultAsyncHttpClientConfig config = new DefaultAsyncHttpClientConfig.Builder() + .setThreadPoolName("ahc-resolver-unit") + .setFallbackNameResolverOffloadEnabled(true) + .setFallbackNameResolverOffloadThreadsCount(1) + .setFallbackNameResolverOffloadQueueSize(1) + .build(); + + try (NameResolverOffload offload = new NameResolverOffload(config)) { + assertTrue(offload.shouldOffload(DEFAULT_NAME_RESOLVER)); + assertFalse(offload.shouldOffload(new DefaultNameResolver(ImmediateEventExecutor.INSTANCE))); + + Promise promise = ImmediateEventExecutor.INSTANCE.newPromise(); + AtomicReference threadName = new AtomicReference<>(); + offload.execute(eventLoop, promise, () -> { + threadName.set(Thread.currentThread().getName()); + offload.completeSuccess(eventLoop, promise, null); + }); + + promise.get(5, SECONDS); + assertTrue(threadName.get().contains("ahc-resolver-unit-resolver")); + } + } + + @Test + void defaultOffloadDoesNotSelectDefaultResolver() { + DefaultAsyncHttpClientConfig config = new DefaultAsyncHttpClientConfig.Builder().build(); + + try (NameResolverOffload offload = new NameResolverOffload(config)) { + assertFalse(offload.shouldOffload(DEFAULT_NAME_RESOLVER)); + } + } + + @Test + void boundedQueueRejectsOverflowAndCloseFailsPendingWork() throws Exception { + DefaultAsyncHttpClientConfig config = new DefaultAsyncHttpClientConfig.Builder() + .setFallbackNameResolverOffloadEnabled(true) + .setFallbackNameResolverOffloadThreadsCount(1) + .setFallbackNameResolverOffloadQueueSize(1) + .build(); + CountDownLatch workerStarted = new CountDownLatch(1); + CountDownLatch releaseWorker = new CountDownLatch(1); + + try (NameResolverOffload offload = new NameResolverOffload(config)) { + Promise running = ImmediateEventExecutor.INSTANCE.newPromise(); + Promise queued = ImmediateEventExecutor.INSTANCE.newPromise(); + Promise rejected = ImmediateEventExecutor.INSTANCE.newPromise(); + + offload.execute(eventLoop, running, () -> { + workerStarted.countDown(); + try { + releaseWorker.await(); + offload.completeSuccess(eventLoop, running, null); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); + assertTrue(workerStarted.await(5, SECONDS)); + + offload.execute(eventLoop, queued, () -> offload.completeSuccess(eventLoop, queued, null)); + offload.execute(eventLoop, rejected, () -> offload.completeSuccess(eventLoop, rejected, null)); + + ExecutionException overflow = assertThrows(ExecutionException.class, () -> rejected.get(5, SECONDS)); + assertInstanceOf(RejectedExecutionException.class, overflow.getCause()); + + offload.close(); + assertRejected(running); + assertRejected(queued); + } + } + + private static void assertRejected(Promise promise) { + ExecutionException failure = assertThrows(ExecutionException.class, () -> promise.get(5, SECONDS)); + assertInstanceOf(RejectedExecutionException.class, failure.getCause()); + assertEquals("Fallback name resolver offload is closed", failure.getCause().getMessage()); + } +}