From 86eee99fd68c707ed9451195784a07540105c6a8 Mon Sep 17 00:00:00 2001 From: Oleg Kalnichevski Date: Wed, 30 Sep 2026 14:31:23 +0200 Subject: [PATCH] Improved handling of idle sessions by AbstractIOSessionPool / H2ConnPool --- .../hc/core5/http2/nio/pool/H2ConnPool.java | 42 ++++++++++--- .../apache/hc/core5/http/HttpConnection.java | 12 +++- .../impl/nio/HttpConnectionEventHandler.java | 7 --- .../core5/reactor/AbstractIOSessionPool.java | 50 ++++++++++++--- .../reactor/TestAbstractIOSessionPool.java | 63 ++++++++++++++++++- 5 files changed, 146 insertions(+), 28 deletions(-) diff --git a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java index fae155bd2..dab0026f7 100644 --- a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java +++ b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/nio/pool/H2ConnPool.java @@ -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; @@ -35,9 +36,9 @@ 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; @@ -45,6 +46,7 @@ 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; @@ -65,16 +67,27 @@ public final class H2ConnPool extends AbstractIOSessionPool { private volatile TimeValue validateAfterInactivity = TimeValue.NEG_ONE_MILLISECOND; + /** + * @since 5.5 + */ public H2ConnPool( + final Clock clock, final ConnectionInitiator connectionInitiator, final Resolver 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 addressResolver, + final TlsStrategy tlsStrategy) { + this(Clock.systemUTC(), connectionInitiator, addressResolver, tlsStrategy); + } + public TimeValue getValidateAfterInactivity() { return validateAfterInactivity; } @@ -99,7 +112,9 @@ protected Future connectSession( final HttpHost namedEndpoint, final Timeout connectTimeout, final FutureCallback 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, @@ -140,11 +155,11 @@ protected void validateSession( final IOSession ioSession, final Callback 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; } @@ -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 diff --git a/httpcore5/src/main/java/org/apache/hc/core5/http/HttpConnection.java b/httpcore5/src/main/java/org/apache/hc/core5/http/HttpConnection.java index 26d42fcdb..4d3a33d4a 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/http/HttpConnection.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/http/HttpConnection.java @@ -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; @@ -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; + } + } diff --git a/httpcore5/src/main/java/org/apache/hc/core5/http/impl/nio/HttpConnectionEventHandler.java b/httpcore5/src/main/java/org/apache/hc/core5/http/impl/nio/HttpConnectionEventHandler.java index 0eb7c8bd8..101005a1b 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/http/impl/nio/HttpConnectionEventHandler.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/http/impl/nio/HttpConnectionEventHandler.java @@ -38,11 +38,4 @@ @Internal public interface HttpConnectionEventHandler extends IOEventHandler, HttpConnection { - /** - * @since 5.5 - */ - default boolean isIdle() { - return false; - } - } diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractIOSessionPool.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractIOSessionPool.java index c9d2947f1..21c0d2df1 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractIOSessionPool.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/AbstractIOSessionPool.java @@ -26,6 +26,7 @@ */ package org.apache.hc.core5.reactor; +import java.time.Clock; import java.util.ArrayDeque; import java.util.HashSet; import java.util.Queue; @@ -33,6 +34,7 @@ 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; @@ -56,11 +58,20 @@ @Contract(threading = ThreadingBehavior.SAFE) public abstract class AbstractIOSessionPool implements ModalCloseable { + private final Clock clock; private final ConcurrentMap 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(); } @@ -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)) { @@ -263,18 +284,27 @@ public final void enumAvailable(final Callback 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 getRoutes() { diff --git a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestAbstractIOSessionPool.java b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestAbstractIOSessionPool.java index 53ad45c4a..dfe844659 100644 --- a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestAbstractIOSessionPool.java +++ b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestAbstractIOSessionPool.java @@ -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; @@ -63,14 +71,21 @@ class TestAbstractIOSessionPool { AutoCloseable closeable; + Clock clock; + AbstractIOSessionPool 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 @@ -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()); + } + }