Skip to content
Closed
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 @@ -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}
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -226,6 +229,9 @@ public class DefaultAsyncHttpClientConfig implements AsyncHttpClientConfig {
private final @Nullable Consumer<Channel> 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;

Expand Down Expand Up @@ -330,6 +336,9 @@ private DefaultAsyncHttpClientConfig(// http
@Nullable Consumer<Channel> wsAdditionalChannelInitializer,
ResponseBodyPartFactory responseBodyPartFactory,
int ioThreadsCount,
boolean fallbackNameResolverOffloadEnabled,
int fallbackNameResolverOffloadThreadsCount,
int fallbackNameResolverOffloadQueueSize,
long hashedWheelTimerTickDuration,
int hashedWheelTimerSize) {

Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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}
*/
Expand Down Expand Up @@ -1016,6 +1043,9 @@ public static class Builder {
private @Nullable Consumer<Channel> 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();

Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -1793,6 +1841,9 @@ public DefaultAsyncHttpClientConfig build() {
wsAdditionalChannelInitializer,
responseBodyPartFactory,
ioThreadsCount,
fallbackNameResolverOffloadEnabled,
fallbackNameResolverOffloadThreadsCount,
fallbackNameResolverOffloadQueueSize,
hashedWheelTickDuration,
hashedWheelSize);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -138,6 +140,7 @@ public class ChannelManager {
private final Map.Entry<ChannelOption<?>, Object>[] channelOptions;
private final long handshakeTimeout;
private final @Nullable AddressResolverGroup<InetSocketAddress> addressResolverGroup;
private final NameResolverOffload nameResolverOffload;

private final ChannelPool channelPool;
private final ChannelGroup openChannels;
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -886,7 +891,7 @@ public Future<Bootstrap> getBootstrap(Uri uri, NameResolver<InetAddress> nameRes
}
});
} else {
nameResolver.resolve(proxy.getHost()).addListener((Future<InetAddress> whenProxyAddress) -> {
resolveProxyHost(nameResolver, proxy.getHost()).addListener((Future<InetAddress> whenProxyAddress) -> {
if (whenProxyAddress.isSuccess()) {
InetSocketAddress proxyAddress = new InetSocketAddress(whenProxyAddress.get(), proxy.getPort());
configureSocksBootstrap(socksBootstrap, httpBootstrapHandler, proxyAddress, proxy, promise);
Expand All @@ -907,6 +912,33 @@ public Future<Bootstrap> getBootstrap(Uri uri, NameResolver<InetAddress> nameRes
return promise;
}

private Future<InetAddress> resolveProxyHost(NameResolver<InetAddress> nameResolver, String host) {
EventExecutor eventLoop = currentEventLoop();
if (eventLoop == null || !nameResolverOffload.shouldOffload(nameResolver)) {
return nameResolver.resolve(host);
}

Promise<InetAddress> promise = ImmediateEventExecutor.INSTANCE.newPromise();
nameResolverOffload.execute(eventLoop, promise, () ->
nameResolver.resolve(host).addListener((Future<InetAddress> 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<Bootstrap> promise) {
socksBootstrap.handler(new ChannelInitializer<Channel>() {
Expand Down Expand Up @@ -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).
Expand Down
Loading
Loading