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 @@ -27,6 +27,7 @@
package org.apache.hc.core5.http2.nio.pool;

import java.net.InetSocketAddress;
import java.time.Clock;
import java.util.concurrent.Future;

import org.apache.hc.core5.annotation.Contract;
Expand All @@ -35,16 +36,17 @@
import org.apache.hc.core5.concurrent.FutureCallback;
import org.apache.hc.core5.function.Callback;
import org.apache.hc.core5.function.Resolver;
import org.apache.hc.core5.http.HttpConnection;
import org.apache.hc.core5.http.HttpHost;
import org.apache.hc.core5.http.URIScheme;
import org.apache.hc.core5.http.impl.DefaultAddressResolver;
import org.apache.hc.core5.http.nio.command.ShutdownCommand;
import org.apache.hc.core5.http.nio.command.StaleCheckCommand;
import org.apache.hc.core5.http.nio.ssl.TlsStrategy;
import org.apache.hc.core5.io.CloseMode;
import org.apache.hc.core5.reactor.AbstractIOSessionPool;
import org.apache.hc.core5.reactor.Command;
import org.apache.hc.core5.reactor.ConnectionInitiator;
import org.apache.hc.core5.reactor.IOEventHandler;
import org.apache.hc.core5.reactor.IOSession;
import org.apache.hc.core5.reactor.ssl.TransportSecurityLayer;
import org.apache.hc.core5.util.Args;
Expand All @@ -65,16 +67,27 @@ public final class H2ConnPool extends AbstractIOSessionPool<HttpHost> {

private volatile TimeValue validateAfterInactivity = TimeValue.NEG_ONE_MILLISECOND;

/**
* @since 5.5
*/
public H2ConnPool(
final Clock clock,
final ConnectionInitiator connectionInitiator,
final Resolver<HttpHost, InetSocketAddress> addressResolver,
final TlsStrategy tlsStrategy) {
super();
super(clock);
this.connectionInitiator = Args.notNull(connectionInitiator, "Connection initiator");
this.addressResolver = addressResolver != null ? addressResolver : DefaultAddressResolver.INSTANCE;
this.addressResolver = addressResolver;
this.tlsStrategy = tlsStrategy;
}

public H2ConnPool(
final ConnectionInitiator connectionInitiator,
final Resolver<HttpHost, InetSocketAddress> addressResolver,
final TlsStrategy tlsStrategy) {
this(Clock.systemUTC(), connectionInitiator, addressResolver, tlsStrategy);
}

public TimeValue getValidateAfterInactivity() {
return validateAfterInactivity;
}
Expand All @@ -99,7 +112,9 @@ protected Future<IOSession> connectSession(
final HttpHost namedEndpoint,
final Timeout connectTimeout,
final FutureCallback<IOSession> callback) {
final InetSocketAddress remoteAddress = addressResolver.resolve(namedEndpoint);
final InetSocketAddress remoteAddress = addressResolver != null ?
addressResolver.resolve(namedEndpoint) :
new InetSocketAddress(namedEndpoint.getHostName(), namedEndpoint.getPort());
return connectionInitiator.connect(
namedEndpoint,
remoteAddress,
Expand Down Expand Up @@ -140,11 +155,11 @@ protected void validateSession(
final IOSession ioSession,
final Callback<Boolean> callback) {
if (ioSession.isOpen()) {
final TimeValue timeValue = validateAfterInactivity;
if (TimeValue.isNonNegative(timeValue)) {
final long lastAccessTime = Math.min(ioSession.getLastReadTime(), ioSession.getLastWriteTime());
final long deadline = lastAccessTime + timeValue.toNanoseconds();
if (deadline <= System.nanoTime()) {
final TimeValue inactivityTime = validateAfterInactivity;
if (TimeValue.isNonNegative(inactivityTime)) {
// Last tolerable point of inactivity
final long deadline = inactivityDeadline(inactivityTime);
if (ioSession.getLastEventTime() <= deadline) {
ioSession.enqueue(new StaleCheckCommand(callback::execute), Command.Priority.NORMAL);
return;
}
Expand All @@ -155,6 +170,15 @@ protected void validateSession(
}
}

@Override
protected boolean isIdle(final IOSession session) {
final IOEventHandler handler = session.getHandler();
if (handler instanceof HttpConnection) {
return ((HttpConnection) handler).isIdle();
}
return false;
}

/**
* Evict expired (closed) sessions.
* @since 5.5
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,8 @@ public interface HttpConnection extends SocketModalCloseable {
/**
* Closes this connection gracefully. This method will attempt to flush the internal output
* buffer prior to closing the underlying socket. This method MUST NOT be called from a
* different thread to force shutdown of the connection. Use {@link #close(CloseMode) close(CloseMode.IMMEDIATE)} instead.
* different thread to force shutdown of the connection.
* Use {@link #close(org.apache.hc.core5.io.CloseMode) close(CloseMode.IMMEDIATE)} instead.
*/
@Override
void close() throws IOException;
Expand Down Expand Up @@ -94,4 +95,13 @@ public interface HttpConnection extends SocketModalCloseable {
*/
boolean isOpen();

/**
* Checks if this connection is idle.
*
* @since 5.5
*/
default boolean isIdle() {
return false;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,4 @@
@Internal
public interface HttpConnectionEventHandler extends IOEventHandler, HttpConnection {

/**
* @since 5.5
*/
default boolean isIdle() {
return false;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -26,13 +26,15 @@
*/
package org.apache.hc.core5.reactor;

import java.time.Clock;
import java.util.ArrayDeque;
import java.util.HashSet;
import java.util.Queue;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.locks.ReentrantLock;

Expand All @@ -56,11 +58,20 @@
@Contract(threading = ThreadingBehavior.SAFE)
public abstract class AbstractIOSessionPool<T> implements ModalCloseable {

private final Clock clock;
private final ConcurrentMap<T, PoolEntry> sessionPool;
private final AtomicBoolean closed;

public AbstractIOSessionPool() {
this(Clock.systemUTC());
}

/**
* @since 5.5
*/
public AbstractIOSessionPool(final Clock clock) {
super();
this.clock = Args.notNull(clock, "Clock");
this.sessionPool = new ConcurrentHashMap<>();
this.closed = new AtomicBoolean();
}
Expand All @@ -78,6 +89,16 @@ protected abstract void closeSession(
IOSession ioSession,
CloseMode closeMode);

/**
* Tests if the given session is idle. This method is expected to be communication protocol
* aware.
*
* @since 5.5
*/
protected boolean isIdle(final IOSession session) {
return false;
}

@Override
public final void close(final CloseMode closeMode) {
if (closed.compareAndSet(false, true)) {
Expand Down Expand Up @@ -263,18 +284,27 @@ public final void enumAvailable(final Callback<IOSession> callback) {
}
}

protected long inactivityDeadline(final TimeValue inactivityTime) {
return TimeUnit.MILLISECONDS.toNanos(TimeValue.isPositive(inactivityTime) ?
this.clock.millis() - inactivityTime.toMilliseconds() :
this.clock.millis());
}

/**
* This method has no effect. Its initial implementation has been removed, as it can
* cause premature termination of sessions with multiplexing message exchanges such
* as HTTP/2.
*
* @deprecated This method has no effect as of version 5.5 and should not be used.
* Use {@link #enumAvailable(Callback)} method and implement idle detection
* logic that correctly takes into account specifics of the underlying
* communication protocol.
* Close sessions idle loner than the given period of inactivity. This method will also
* Evict expired (closed) sessions.
*/
@Deprecated
public final void closeIdle(final TimeValue idleTime) {
public final void closeIdle(final TimeValue inactivityTime) {
// Last tolerable point of inactivity
// Millisecond precision is good enough
final long deadline = inactivityDeadline(inactivityTime);
enumAvailable(session -> {
if (isIdle(session)) {
if (session.getLastEventTime() <= deadline) {
closeSession(session, CloseMode.GRACEFUL);
}
}
});
}

public final Set<T> getRoutes() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,11 +28,19 @@


import java.net.UnknownHostException;
import java.time.Clock;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.Month;
import java.time.ZoneId;
import java.time.ZoneOffset;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;

import org.apache.hc.core5.concurrent.FutureCallback;
import org.apache.hc.core5.function.Callback;
import org.apache.hc.core5.io.CloseMode;
import org.apache.hc.core5.util.TimeValue;
import org.apache.hc.core5.util.Timeout;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Assertions;
Expand Down Expand Up @@ -63,14 +71,21 @@ class TestAbstractIOSessionPool {

AutoCloseable closeable;

Clock clock;

AbstractIOSessionPool<String> impl;

@BeforeEach
void prepareMocks() {
closeable = MockitoAnnotations.openMocks(this);

clock = Clock.fixed(
LocalDateTime.of(2026, Month.SEPTEMBER, 30, 13, 30).toInstant(ZoneOffset.UTC),
ZoneId.of("UTC"));

impl = Mockito.mock(AbstractIOSessionPool.class, Mockito.withSettings()
.defaultAnswer(Answers.CALLS_REAL_METHODS)
.useConstructor());
.useConstructor(clock));
}

@AfterEach
Expand Down Expand Up @@ -270,4 +285,50 @@ void testGetSessionConnectUnknownHost() {
Assertions.assertTrue(future2.isDone());
}

@Test
void testCloseIdleOnly() {
final AbstractIOSessionPool.PoolEntry entry1 = impl.getPoolEntry("host1");
Assertions.assertNotNull(entry1);
entry1.session = ioSession1;

final AbstractIOSessionPool.PoolEntry entry2 = impl.getPoolEntry("host2");
Assertions.assertNotNull(entry2);
entry2.session = ioSession2;

Mockito.doReturn(true).when(impl).isIdle(ioSession1);
Mockito.doReturn(TimeUnit.MILLISECONDS.toNanos(Instant.now(clock).minusSeconds(1).toEpochMilli()))
.when(ioSession1).getLastEventTime();
Mockito.doReturn(false).when(impl).isIdle(ioSession2);
Mockito.doReturn(TimeUnit.MILLISECONDS.toNanos(Instant.now(clock).minusSeconds(2).toEpochMilli()))
.when(ioSession2).getLastEventTime();

impl.closeIdle(null);

Mockito.verify(impl).closeSession(ioSession1, CloseMode.GRACEFUL);
Mockito.verify(impl, Mockito.never()).closeSession(Mockito.same(ioSession2), Mockito.any());
}

@Test
void testCloseIdlePastDeadline() {
final AbstractIOSessionPool.PoolEntry entry1 = impl.getPoolEntry("host1");
Assertions.assertNotNull(entry1);
entry1.session = ioSession1;

final AbstractIOSessionPool.PoolEntry entry2 = impl.getPoolEntry("host2");
Assertions.assertNotNull(entry2);
entry2.session = ioSession2;

Mockito.doReturn(true).when(impl).isIdle(ioSession1);
Mockito.doReturn(TimeUnit.MILLISECONDS.toNanos(Instant.now(clock).minusSeconds(1).toEpochMilli()))
.when(ioSession1).getLastEventTime();
Mockito.doReturn(true).when(impl).isIdle(ioSession2);
Mockito.doReturn(TimeUnit.MILLISECONDS.toNanos(Instant.now(clock).toEpochMilli()))
.when(ioSession2).getLastEventTime();

impl.closeIdle(TimeValue.ofSeconds(1));

Mockito.verify(impl).closeSession(ioSession1, CloseMode.GRACEFUL);
Mockito.verify(impl, Mockito.never()).closeSession(Mockito.same(ioSession2), Mockito.any());
}

}
Loading