From 19433093d12ae5a3ac9a240c083438bc709b120a Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Wed, 30 Sep 2026 11:10:17 +0300 Subject: [PATCH 1/4] [#1119] Serve each LDAP connection handler with a transport of its own, built from its configuration LDAPConnectionHandler2 ran every listener on the transport the SDK shares across the JVM, so use-tcp-keep-alive, use-tcp-no-delay, buffer-size and num-request-handlers had no effect and allow-tcp-reuse-address reached only the pre-check. - GrizzlyLDAPListener takes a transport through the new GRIZZLY_TRANSPORT option and leaves it running when it is closed. - LDAPConnectionHandler2 starts a transport per listener: num-request-handlers selector threads (else org.forgerock.opendj.transport.selectors, else chosen from the processors), SO_REUSEADDR on the listen socket from allow-tcp-reuse-address, and buffer-size as the write buffer. SO_KEEPALIVE and TCP_NODELAY reach every accepted socket through the listener options. The transports share one pooled memory manager. - Stopping the listener leaves the connections it accepted open: its transport keeps serving them and is shut down after the last one closes. When the server shuts down, they are ended as a server shutdown, with a notice of disconnection, before the transport is shut down. - use-tcp-keep-alive, use-tcp-no-delay and buffer-size of the LDAP connection handler now require a component restart. Fixes #1119 --- .../opendj/grizzly/GrizzlyLDAPListener.java | 14 +- .../grizzly/GrizzlyLDAPListenerTestCase.java | 40 ++ .../LDAPConnectionHandlerConfiguration.xml | 71 +++- opendj-server-legacy/pom.xml | 5 + .../reactive/LDAPConnectionHandler2.java | 189 ++++++++- ...APConnectionHandler2TransportTestCase.java | 391 ++++++++++++++++++ 6 files changed, 700 insertions(+), 10 deletions(-) create mode 100644 opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java diff --git a/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListener.java b/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListener.java index 19e36fea4c..127283d17b 100644 --- a/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListener.java +++ b/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListener.java @@ -13,7 +13,7 @@ * * Copyright 2010 Sun Microsystems, Inc. * Portions copyright 2011-2016 ForgeRock AS. - * Portions copyright 2025 3A Systems, LLC. + * Portions copyright 2025-2026 3A Systems, LLC. */ package org.forgerock.opendj.grizzly; @@ -37,6 +37,7 @@ import org.forgerock.opendj.ldap.spi.LDAPListenerImpl; import org.forgerock.opendj.ldap.spi.LdapMessages.LdapRequestEnvelope; import org.forgerock.util.Function; +import org.forgerock.util.Option; import org.forgerock.util.Options; import org.glassfish.grizzly.filterchain.FilterChain; import org.glassfish.grizzly.nio.transport.TCPNIOBindingHandler; @@ -52,6 +53,13 @@ * LDAP listener implementation using Grizzly for transport. */ public final class GrizzlyLDAPListener implements LDAPListenerImpl { + /** + * Grizzly TCP Transport NIO implementation to bind the listener to and to serve its connections with. If + * {@code null}, the default server transport shared by every listener in the JVM will be used. A transport + * provided here is not shut down when the listener is closed: its owner remains responsible for it. + */ + public static final Option GRIZZLY_TRANSPORT = Option.of(TCPNIOTransport.class, null); + private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass(); private final ReferenceCountedObject.Reference transport; private final Collection serverConnections; @@ -66,7 +74,7 @@ public final class GrizzlyLDAPListener implements LDAPListenerImpl { * @param addresses * The addresses to listen on. * @param options - * The LDAP listener options. + * The LDAP listener options, including the optional {@link #GRIZZLY_TRANSPORT}. * @param requestHandlerFactory * The server connection factory which will be used to create server connections. * @throws IOException @@ -76,7 +84,7 @@ public GrizzlyLDAPListener(final Set addresses, final Options final Function>, LdapException> requestHandlerFactory) throws IOException { - this(addresses, requestHandlerFactory, options, null); + this(addresses, requestHandlerFactory, options, options.get(GRIZZLY_TRANSPORT)); } /** diff --git a/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListenerTestCase.java b/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListenerTestCase.java index 89b7787dc1..fae5bbd806 100644 --- a/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListenerTestCase.java +++ b/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListenerTestCase.java @@ -29,6 +29,7 @@ import static org.mockito.Mockito.mock; import java.io.IOException; +import java.lang.reflect.Field; import java.net.InetSocketAddress; import java.net.ServerSocket; import java.util.Arrays; @@ -72,6 +73,8 @@ import org.forgerock.opendj.ldap.responses.Result; import org.forgerock.util.Options; import org.forgerock.util.promise.PromiseImpl; +import org.glassfish.grizzly.nio.transport.TCPNIOTransport; +import org.glassfish.grizzly.nio.transport.TCPNIOTransportBuilder; import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; @@ -280,6 +283,43 @@ public void testLDAPListenerBasic() throws Exception { } } + /** + * A listener given a transport with {@link GrizzlyLDAPListener#GRIZZLY_TRANSPORT} serves its connections with that + * transport, and leaves it running when it is closed. + */ + @Test(timeOut = 10000) + public void testLDAPListenerWithProvidedTransport() throws Exception { + final TCPNIOTransport transport = TCPNIOTransportBuilder.newInstance().build(); + transport.start(); + try { + final MockServerConnection serverConnection = new MockServerConnection(); + final Options options = defaultOptions().set(GrizzlyLDAPListener.GRIZZLY_TRANSPORT, transport); + final LDAPListener listener = new LDAPListener(Collections.singleton(loopbackWithDynamicPort()), + new ServerConnectionFactoryAdapter(options.get(LDAP_DECODE_OPTIONS), + new MockServerConnectionFactory(serverConnection)), + options); + try { + final InetSocketAddress addr = listener.firstSocketAddress(); + final Connection connection = + new LDAPConnectionFactory(addr.getHostName(), addr.getPort()).getConnection(); + try { + final LDAPClientContext context = serverConnection.context.get(10, TimeUnit.SECONDS); + final Field field = context.getClass().getDeclaredField("connection"); + field.setAccessible(true); + assertThat(((org.glassfish.grizzly.Connection) field.get(context)).getTransport()) + .isSameAs(transport); + } finally { + connection.close(); + } + } finally { + listener.close(); + } + assertThat(transport.isStopped()).isFalse(); + } finally { + transport.shutdownNow(); + } + } + /** * Tests LDAP listener which attempts to open a connection to a remote * offline server at the point when the listener accepts the client diff --git a/opendj-maven-plugin/src/main/resources/config/xml/org/forgerock/opendj/server/config/LDAPConnectionHandlerConfiguration.xml b/opendj-maven-plugin/src/main/resources/config/xml/org/forgerock/opendj/server/config/LDAPConnectionHandlerConfiguration.xml index a21dbbafe5..0e596cc42a 100644 --- a/opendj-maven-plugin/src/main/resources/config/xml/org/forgerock/opendj/server/config/LDAPConnectionHandlerConfiguration.xml +++ b/opendj-maven-plugin/src/main/resources/config/xml/org/forgerock/opendj/server/config/LDAPConnectionHandlerConfiguration.xml @@ -88,8 +88,72 @@ - - + + + Indicates whether the + + should use TCP keep-alive. + + + If enabled, the SO_KEEPALIVE socket option is used to indicate that TCP + keepalive messages should periodically be sent to the client to + verify that the associated connection is still valid. This may + also help prevent cases in which intermediate network hardware + could silently drop an otherwise idle client connection, provided + that the keepalive interval configured in the underlying operating + system is smaller than the timeout enforced by the network + hardware. + + + + + + + true + + + + + + + + ds-cfg-use-tcp-keep-alive + + + + + + Indicates whether the + + should use TCP no-delay. + + + If enabled, the TCP_NODELAY socket option is used to ensure + that response messages to the client are sent immediately rather + than potentially waiting to determine whether additional response + messages can be sent in the same packet. In most cases, using the + TCP_NODELAY socket option provides better performance and + lower response times, but disabling it may help for some cases in + which the server sends a large number of entries to a client + in response to a search request. + + + + + + + true + + + + + + + + ds-cfg-use-tcp-no-delay + + + @@ -338,6 +402,9 @@ each client connection and used to buffer LDAP response messages data when writing. + + + 4096 bytes diff --git a/opendj-server-legacy/pom.xml b/opendj-server-legacy/pom.xml index c883bc7a5d..81ba8ee273 100644 --- a/opendj-server-legacy/pom.xml +++ b/opendj-server-legacy/pom.xml @@ -95,6 +95,11 @@ opendj-server + + org.openidentityplatform.opendj + opendj-grizzly + + org.openidentityplatform.opendj opendj-ldap-toolkit diff --git a/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java b/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java index 5f40bd1215..dc1de10e0e 100644 --- a/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java +++ b/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java @@ -41,9 +41,11 @@ import java.util.SortedSet; import java.util.TreeSet; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import javax.net.ssl.KeyManager; import javax.net.ssl.SSLContext; @@ -55,6 +57,7 @@ import org.forgerock.opendj.config.server.ConfigChangeResult; import org.forgerock.opendj.config.server.ConfigException; import org.forgerock.opendj.config.server.ConfigurationChangeListener; +import org.forgerock.opendj.grizzly.GrizzlyLDAPListener; import org.forgerock.opendj.ldap.AddressMask; import org.forgerock.opendj.ldap.DN; import org.forgerock.opendj.ldap.LDAPClientContext; @@ -69,6 +72,12 @@ import org.forgerock.opendj.server.config.server.LDAPConnectionHandlerCfg; import org.forgerock.util.Function; import org.forgerock.util.Options; +import org.glassfish.grizzly.memory.MemoryManager; +import org.glassfish.grizzly.memory.PooledMemoryManager; +import org.glassfish.grizzly.nio.transport.TCPNIOTransport; +import org.glassfish.grizzly.nio.transport.TCPNIOTransportBuilder; +import org.glassfish.grizzly.strategies.SameThreadIOStrategy; +import org.glassfish.grizzly.threadpool.ThreadPoolConfig; import org.glassfish.grizzly.utils.ArrayUtils; import org.opends.server.api.AlertGenerator; import org.opends.server.api.ClientConnection; @@ -125,8 +134,85 @@ public void run() { } } + /** + * Holds the memory manager shared by the transports of every LDAP connection handler. It pre-allocates a share of + * the heap, so one instance per transport would multiply that share by the number of handlers. + */ + private static final class MemoryManagerHolder { + /** Pooled, as in the SDK server transport, so that Grizzly buffers can be used across threads. */ + private static final MemoryManager INSTANCE = new PooledMemoryManager(true); + } + + /** + * Keeps the transport of a stopped listener running for the connections it has accepted, so that disabling, + * deleting or restarting the handler leaves them open, as the transport the SDK shares between listeners did. The + * transport is shut down once the last of them is closed, or when the server shuts down. + */ + private final class TransportDrain implements ServerShutdownListener { + private final TCPNIOTransport drained; + private final AtomicBoolean finished = new AtomicBoolean(); + + private TransportDrain(TCPNIOTransport drained) { + this.drained = drained; + } + + private void start() { + drains.add(this); + DirectoryServer.registerShutdownListener(this); + connectionClosed(); + } + + /** Shuts the transport down if no connection of this handler is left. */ + private void connectionClosed() { + if (clientConnections.isEmpty() && finish()) { + // The last connection may be closed on a selector thread of the transport being shut down. + new DirectoryThread(this::shutdownTransport, getShutdownListenerName()).start(); + } + } + + @Override + public String getShutdownListenerName() { + return "Transport drain of " + handlerName; + } + + @Override + public void processServerShutdown(LocalizableMessage reason) { + if (finish()) { + // Shutting the transport down closes every connection it accepted: end them as a server shutdown + // first, as the legacy connection handler does, rather than let them fail as protocol errors. + for (ClientConnection clientConnection : clientConnections) { + clientConnection.disconnect(DisconnectReason.SERVER_SHUTDOWN, true, reason); + } + shutdownTransport(); + } + } + + private boolean finish() { + if (!finished.compareAndSet(false, true)) { + return false; + } + drains.remove(this); + DirectoryServer.deregisterShutdownListener(this); + return true; + } + + private void shutdownTransport() { + try { + drained.shutdownNow(); + } catch (IOException e) { + logger.traceException(e); + } + } + } + private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass(); + /** + * The system property that sized the transport shared by every listener in the JVM. It is still honoured when the + * configuration lets the server choose the number of request handlers. + */ + private static final String SELECTORS_PROPERTY = "org.forgerock.opendj.transport.selectors"; + /** Default friendly name for the LDAP connection handler. */ private static final String DEFAULT_FRIENDLY_NAME = "LDAP Connection Handler"; @@ -135,6 +221,15 @@ public void run() { private LDAPListener listener; + /** + * The transport serving the connections of this connection handler only. It is started with the listener, and + * handed over to a {@link TransportDrain} when the listener stops. + */ + private TCPNIOTransport transport; + + /** The transports of stopped listeners still serving the connections they accepted. */ + private final Collection drains = new CopyOnWriteArrayList<>(); + /** The current configuration state. */ private LDAPConnectionHandlerCfg currentConfig; @@ -152,9 +247,24 @@ public void run() { /** Indicates whether to allow the reuse address socket option. */ private boolean allowReuseAddress; + /** Indicates whether to use the SO_KEEPALIVE socket option on client connections. */ + private boolean useTCPKeepAlive; + + /** Indicates whether to use the TCP_NODELAY socket option on client connections. */ + private boolean useTCPNoDelay; + + /** The size in bytes of the write buffer of each client connection. */ + private int bufferSize; + + /** The number of threads reading requests from the client connections. */ + private int numRequestHandlers; + /** Indicates whether the Directory Server is in the process of shutting down. */ private volatile boolean shutdownRequested; + /** The reason given to the client connections still open when the server shuts down. */ + private volatile LocalizableMessage closeReason; + /* Internal LDAP connection handler state */ /** Indicates whether this connection handler is enabled. */ @@ -280,7 +390,8 @@ public ConfigChangeResult applyConfigurationChange(LDAPConnectionHandlerCfg conf // * ssl policy // * ssl cert nickname // * accept backlog - // * tcp reuse address + // * tcp reuse address, keep alive and no delay + // * buffer size // * num request handler // Clear the stat tracker if LDAPv2 is being enabled. @@ -326,6 +437,7 @@ private void configureSSL(LDAPConnectionHandlerCfg config) throws DirectoryExcep @Override public void finalizeConnectionHandler(LocalizableMessage finalizeReason) { + closeReason = finalizeReason; shutdownRequested = true; currentConfig.removeLDAPChangeListener(this); @@ -502,6 +614,10 @@ public void initializeConnectionHandler(ServerContext serverContext, LDAPConnect // Save properties that cannot be dynamically modified. allowReuseAddress = config.isAllowTCPReuseAddress(); backlog = config.getAcceptBacklog(); + useTCPKeepAlive = config.isUseTCPKeepAlive(); + useTCPNoDelay = config.isUseTCPNoDelay(); + bufferSize = (int) config.getBufferSize(); + numRequestHandlers = getNumRequestHandlers(config); listenAddresses = new HashSet<>(); for (InetAddress addr : config.getListenAddress()) { listenAddresses.add(new InetSocketAddress(addr, config.getListenPort())); @@ -561,6 +677,21 @@ public void initializeConnectionHandler(ServerContext serverContext, LDAPConnect config.addLDAPChangeListener(this); } + /** + * Returns the number of request handlers set in the configuration, or when the configuration lets the server decide, + * the value of {@link #SELECTORS_PROPERTY}, or else a number chosen from the number of processors. + */ + private int getNumRequestHandlers(LDAPConnectionHandlerCfg config) { + Integer configured = config.getNumRequestHandlers(); + if (configured == null) { + final Integer selectors = Integer.getInteger(SELECTORS_PROPERTY); + if (selectors != null && selectors > 0) { + configured = selectors; + } + } + return getNumRequestHandlers(configured, friendlyName); + } + @Override public boolean isConfigurationAcceptable(ConnectionHandlerCfg configuration, List unacceptableReasons) { @@ -640,9 +771,53 @@ void stopListener() { listener = null; logger.info(NOTE_CONNHANDLER_STOPPED_LISTENING, handlerName); } + if (transport != null) { + final TransportDrain drain = new TransportDrain(transport); + transport = null; + if (DirectoryServer.getInstance().isShuttingDown()) { + drain.processServerShutdown(closeReason); + } else { + drain.start(); + } + } + } + + private void connectionClosed(ClientConnection connection) { + clientConnections.remove(connection); + for (TransportDrain drain : drains) { + drain.connectionClosed(); + } + } + + /** + * Creates and starts the transport of this connection handler. Each handler has its own, so that its selector + * threads and its socket options follow its configuration rather than the JVM-wide settings of the transport the + * SDK shares between listeners. + */ + private TCPNIOTransport newTransport() throws IOException { + final TCPNIOTransport newTransport = TCPNIOTransportBuilder.newInstance() + .setIOStrategy(SameThreadIOStrategy.getInstance()) + .setSelectorThreadPoolConfig(ThreadPoolConfig.defaultConfig() + .setCorePoolSize(numRequestHandlers) + .setMaxPoolSize(numRequestHandlers) + .setPoolName(handlerName + " Request Handler")) + .setMemoryManager(MemoryManagerHolder.INSTANCE) + .setReuseAddress(allowReuseAddress) + .setWriteBufferSize(bufferSize) + .build(); + // As in the SDK server transport: fewer selector runners than selector threads cause deadlocks. + newTransport.setSelectorRunnersCount(numRequestHandlers); + try { + newTransport.start(); + } catch (IOException e) { + newTransport.shutdownNow(); + throw e; + } + return newTransport; } private void startListener() throws IOException { + transport = newTransport(); listener = new LDAPListener( listenAddresses, new Function> clientContext.addListener(new LDAPClientContextEventListener() { @Override public void handleConnectionError(final LDAPClientContext context, final Throwable error) { - clientConnections.remove(conn); + connectionClosed(conn); } @Override public void handleConnectionDisconnected(final LDAPClientContext context, final ResultCode resultCode, String diagnosticMessage) { - clientConnections.remove(conn); + connectionClosed(conn); } @Override public void handleConnectionClosed(final LDAPClientContext context, final UnbindRequest unbindRequest) { - clientConnections.remove(conn); + connectionClosed(conn); } }); return new ReactiveHandler>() { @@ -682,7 +857,11 @@ public Stream handle(final LDAPClientContext context, } }, Options.defaultOptions() .set(LDAPListener.CONNECT_MAX_BACKLOG, backlog) - .set(LDAPListener.REQUEST_MAX_SIZE_IN_BYTES, (int) currentConfig.getMaxRequestSize())); + .set(LDAPListener.REQUEST_MAX_SIZE_IN_BYTES, (int) currentConfig.getMaxRequestSize()) + // Set by the listener on every accepted connection, over the values of the transport. + .set(LDAPListener.SO_KEEPALIVE, useTCPKeepAlive) + .set(LDAPListener.TCP_NO_DELAY, useTCPNoDelay) + .set(GrizzlyLDAPListener.GRIZZLY_TRANSPORT, transport)); logger.info(NOTE_CONNHANDLER_STARTED_LISTENING, handlerName); } diff --git a/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java new file mode 100644 index 0000000000..3dd2b5d2ae --- /dev/null +++ b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java @@ -0,0 +1,391 @@ +/* + * The contents of this file are subject to the terms of the Common Development and + * Distribution License (the License). You may not use this file except in compliance with the + * License. + * + * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the + * specific language governing permission and limitations under the License. + * + * When distributing Covered Software, include this CDDL Header Notice in each file and include + * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL + * Header, with the fields enclosed by brackets [] replaced by your own identifying + * information: "Portions copyright [year] [name of copyright owner]". + * + * Copyright 2026 3A Systems, LLC. + */ +package org.opends.server.protocols.ldap; + +import static org.testng.Assert.*; + +import java.io.ByteArrayOutputStream; +import java.io.InputStream; +import java.lang.reflect.Field; +import java.net.Socket; +import java.net.SocketTimeoutException; +import java.nio.channels.ServerSocketChannel; +import java.nio.channels.SocketChannel; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; + +import org.forgerock.i18n.LocalizableMessage; +import org.forgerock.opendj.reactive.LDAPConnectionHandler2; +import org.forgerock.opendj.server.config.meta.LDAPConnectionHandlerCfgDefn; +import org.forgerock.opendj.server.config.server.LDAPConnectionHandlerCfg; +import org.glassfish.grizzly.nio.transport.TCPNIOConnection; +import org.glassfish.grizzly.nio.transport.TCPNIOServerConnection; +import org.glassfish.grizzly.nio.transport.TCPNIOTransport; +import org.opends.server.DirectoryServerTestCase; +import org.opends.server.TestCaseUtils; +import org.opends.server.api.ClientConnection; +import org.opends.server.api.ServerShutdownListener; +import org.opends.server.core.DirectoryServer; +import org.opends.server.extensions.InitializationUtils; +import org.opends.server.types.Entry; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.DataProvider; +import org.testng.annotations.Test; + +/** + * {@link LDAPConnectionHandler2} serves its connections with a transport of its own, built from its configuration: + * {@code use-tcp-keep-alive} and {@code use-tcp-no-delay} reach every accepted socket, {@code buffer-size} is the write + * buffer of every connection, {@code allow-tcp-reuse-address} reaches the listen socket and {@code num-request-handlers} + * is the number of selector threads. Stopping the handler leaves its connections open until they are closed, or until + * the server shuts down, which ends them with a notice of disconnection; then its threads stop. + */ +@SuppressWarnings("javadoc") +@Test(groups = { "precommit" }, sequential = true) +public class LDAPConnectionHandler2TransportTestCase extends DirectoryServerTestCase +{ + private static final LocalizableMessage STOP_REASON = LocalizableMessage.raw("Stopped by the transport test."); + private static final String NOTICE_OF_DISCONNECTION_OID = "1.3.6.1.4.1.1466.20036"; + private static final String SELECTORS_PROPERTY = "org.forgerock.opendj.transport.selectors"; + /** Unlike any socket buffer size a system would choose. */ + private static final int WRITE_BUFFER_SIZE = 12345; + private static final long TIMEOUT_MS = 10000; + /** How long the connection of a stopped handler must stay open. */ + private static final long KEEPS_OPEN_MS = 3000; + + @BeforeClass + public void setUp() throws Exception + { + TestCaseUtils.startServer(); + } + + @DataProvider + public Object[][] socketOptions() + { + return new Object[][] { { true, true }, { true, false }, { false, true }, { false, false } }; + } + + @Test(dataProvider = "socketOptions") + public void acceptedSocketUsesTheConfiguredOptions(boolean keepAlive, boolean noDelay) throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, keepAlive, noDelay, true, 2)); + try (Socket client = new Socket("127.0.0.1", port)) + { + final Socket accepted = ((SocketChannel) acceptedConnection(handler).getChannel()).socket(); + assertEquals(accepted.getKeepAlive(), keepAlive, "SO_KEEPALIVE"); + assertEquals(accepted.getTcpNoDelay(), noDelay, "TCP_NODELAY"); + } + finally + { + stop(handler); + } + } + + @Test + public void connectionIsServedByTheTransportOfItsHandlerWithTheConfiguredWriteBuffer() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, true, true, true, 2)); + try (Socket client = new Socket("127.0.0.1", port)) + { + final TCPNIOConnection accepted = acceptedConnection(handler); + assertSame(accepted.getTransport(), transport(handler), "the connection is not served by its handler's transport"); + assertEquals(accepted.getWriteBufferSize(), WRITE_BUFFER_SIZE); + } + finally + { + stop(handler); + } + } + + @DataProvider + public Object[][] reuseAddress() + { + return new Object[][] { { true }, { false } }; + } + + @Test(dataProvider = "reuseAddress") + public void listenSocketUsesTheConfiguredReuseAddress(boolean reuseAddress) throws Exception + { + final LDAPConnectionHandler2 handler = + start(configuration(TestCaseUtils.findFreePort(), true, true, reuseAddress, 2)); + try + { + final Collection serverConnections = (Collection) field(field(handler, "listener"), "impl", "serverConnections"); + assertFalse(serverConnections.isEmpty(), "the handler does not listen"); + for (Object serverConnection : serverConnections) + { + final ServerSocketChannel channel = (ServerSocketChannel) ((TCPNIOServerConnection) serverConnection).getChannel(); + assertEquals(channel.socket().getReuseAddress(), reuseAddress, "SO_REUSEADDR"); + } + } + finally + { + stop(handler); + } + } + + @Test + public void selectorThreadsFollowNumRequestHandlers() throws Exception + { + assertSelectorThreads(configuration(TestCaseUtils.findFreePort(), true, true, true, 3), 3); + } + + /** The system property that sized the transport shared by every listener still applies when nothing is set. */ + @Test + public void selectorThreadsFollowTheSelectorsPropertyWhenNumRequestHandlersIsUnset() throws Exception + { + final String saved = System.getProperty(SELECTORS_PROPERTY); + System.setProperty(SELECTORS_PROPERTY, "13"); + try + { + assertSelectorThreads(configuration(TestCaseUtils.findFreePort(), true, true, true, null), 13); + } + finally + { + restore(saved); + } + } + + @Test + public void selectorThreadsAreChosenFromTheProcessorsWhenNothingIsSet() throws Exception + { + final String saved = System.getProperty(SELECTORS_PROPERTY); + System.clearProperty(SELECTORS_PROPERTY); + try + { + assertSelectorThreads(configuration(TestCaseUtils.findFreePort(), true, true, true, null), + Math.max(2, Runtime.getRuntime().availableProcessors() / 2)); + } + finally + { + restore(saved); + } + } + + /** + * Stopping the handler, as disabling, deleting or restarting it does, leaves the connections it has accepted open: + * its transport keeps serving them, and is shut down once the last of them is closed. + */ + @Test + public void stoppedHandlerKeepsItsConnectionsUntilTheyAreClosed() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, true, true, true, 2)); + final String threadPrefix = selectorThreadPrefix(handler); + try (Socket client = new Socket("127.0.0.1", port)) + { + acceptedConnection(handler); + stop(handler); + + client.setSoTimeout((int) KEEPS_OPEN_MS); + try + { + final int read = client.getInputStream().read(); + fail("the connection of the stopped handler was " + (read == -1 ? "closed" : "written to")); + } + catch (SocketTimeoutException expected) + { + // still open + } + assertFalse(threadsNamed(threadPrefix).isEmpty(), "the transport stopped while a connection was open"); + } + finally + { + stop(handler); + } + assertThreadsStop(threadPrefix); + } + + /** The server shutting down ends the connections a stopped handler left open, and shuts its transport down. */ + @Test + public void serverShutdownEndsTheConnectionsOfAStoppedHandler() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, true, true, true, 2)); + final String threadPrefix = selectorThreadPrefix(handler); + try (Socket client = new Socket("127.0.0.1", port)) + { + acceptedConnection(handler); + stop(handler); + final Collection drains = (Collection) field(handler, "drains"); + assertEquals(drains.size(), 1, "no transport is left serving the connection of the stopped handler"); + + ((ServerShutdownListener) drains.iterator().next()).processServerShutdown(STOP_REASON); + + client.setSoTimeout((int) TIMEOUT_MS); + final ByteArrayOutputStream received = new ByteArrayOutputStream(); + final InputStream in = client.getInputStream(); + final byte[] buffer = new byte[256]; + for (int read; (read = in.read(buffer)) != -1;) + { + received.write(buffer, 0, read); + } + assertTrue(new String(received.toByteArray(), StandardCharsets.ISO_8859_1).contains(NOTICE_OF_DISCONNECTION_OID), + "the connection was closed without a notice of disconnection"); + } + finally + { + stop(handler); + } + assertThreadsStop(threadPrefix); + } + + /** Returns the prefix of the names of the selector threads of the handler, after checking that some run. */ + private static String selectorThreadPrefix(LDAPConnectionHandler2 handler) + { + final String prefix = handler.getConnectionHandlerName() + " Request Handler"; + assertFalse(threadsNamed(prefix).isEmpty(), "no thread is named after the handler: " + prefix); + return prefix; + } + + private static void assertThreadsStop(String prefix) throws InterruptedException + { + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + List left; + while (!(left = threadsNamed(prefix)).isEmpty()) + { + assertTrue(System.currentTimeMillis() < deadline, "threads of the stopped handler are still alive: " + left); + Thread.sleep(100); + } + } + + private static void assertSelectorThreads(LDAPConnectionHandlerCfg config, int expected) throws Exception + { + final LDAPConnectionHandler2 handler = start(config); + try + { + final TCPNIOTransport transport = transport(handler); + assertEquals(transport.getSelectorRunnersCount(), expected, "selector runners"); + assertEquals(transport.getKernelThreadPoolConfig().getMaxPoolSize(), expected, "selector threads"); + } + finally + { + stop(handler); + } + } + + private static void restore(String selectors) + { + if (selectors != null) + { + System.setProperty(SELECTORS_PROPERTY, selectors); + } + else + { + System.clearProperty(SELECTORS_PROPERTY); + } + } + + private static List threadsNamed(String prefix) + { + final List names = new ArrayList<>(); + for (Thread thread : Thread.getAllStackTraces().keySet()) + { + if (thread.isAlive() && thread.getName().startsWith(prefix)) + { + names.add(thread.getName()); + } + } + return names; + } + + private static TCPNIOTransport transport(LDAPConnectionHandler2 handler) throws Exception + { + final TCPNIOTransport transport = (TCPNIOTransport) field(handler, "transport"); + assertNotNull(transport, "the handler has no transport"); + return transport; + } + + /** Waits for the handler to accept a connection and returns the Grizzly connection behind it. */ + private static TCPNIOConnection acceptedConnection(LDAPConnectionHandler2 handler) throws Exception + { + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + while (handler.getClientConnections().isEmpty()) + { + assertTrue(System.currentTimeMillis() < deadline, "the handler did not accept the connection"); + Thread.sleep(50); + } + final ClientConnection connection = handler.getClientConnections().iterator().next(); + return (TCPNIOConnection) field(connection, "clientContext", "connection"); + } + + /** Follows a chain of fields declared by the class of each object met. */ + private static Object field(Object target, String... names) throws Exception + { + Object value = target; + for (String name : names) + { + final Field field = value.getClass().getDeclaredField(name); + field.setAccessible(true); + value = field.get(value); + } + return value; + } + + private static LDAPConnectionHandlerCfg configuration(int port, boolean keepAlive, boolean noDelay, + boolean reuseAddress, Integer numRequestHandlers) throws Exception + { + final List lines = new ArrayList<>(); + lines.add("dn: cn=Transport Test Handler,cn=Connection Handlers,cn=config"); + lines.add("objectClass: top"); + lines.add("objectClass: ds-cfg-connection-handler"); + lines.add("objectClass: ds-cfg-ldap-connection-handler"); + lines.add("cn: Transport Test Handler"); + lines.add("ds-cfg-java-class: " + LDAPConnectionHandler2.class.getName()); + lines.add("ds-cfg-enabled: true"); + lines.add("ds-cfg-listen-address: 127.0.0.1"); + lines.add("ds-cfg-listen-port: " + port); + lines.add("ds-cfg-accept-backlog: 128"); + lines.add("ds-cfg-keep-stats: false"); + lines.add("ds-cfg-use-tcp-keep-alive: " + keepAlive); + lines.add("ds-cfg-use-tcp-no-delay: " + noDelay); + lines.add("ds-cfg-allow-tcp-reuse-address: " + reuseAddress); + lines.add("ds-cfg-buffer-size: " + WRITE_BUFFER_SIZE + " bytes"); + lines.add("ds-cfg-use-ssl: false"); + lines.add("ds-cfg-allow-start-tls: false"); + lines.add("ds-cfg-allow-ldap-v2: false"); + lines.add("ds-cfg-send-rejection-notice: true"); + if (numRequestHandlers != null) + { + lines.add("ds-cfg-num-request-handlers: " + numRequestHandlers); + } + final Entry entry = TestCaseUtils.makeEntry(lines.toArray(new String[0])); + return InitializationUtils.getConfiguration(LDAPConnectionHandlerCfgDefn.getInstance(), entry); + } + + private static LDAPConnectionHandler2 start(LDAPConnectionHandlerCfg config) throws Exception + { + final LDAPConnectionHandler2 handler = new LDAPConnectionHandler2(); + handler.initializeConnectionHandler(DirectoryServer.getInstance().getServerContext(), config); + handler.start(); + return handler; + } + + private static void stop(LDAPConnectionHandler2 handler) throws InterruptedException + { + if (!handler.isAlive()) + { + return; + } + handler.processServerShutdown(STOP_REASON); + handler.finalizeConnectionHandler(STOP_REASON); + handler.join(TIMEOUT_MS); + assertFalse(handler.isAlive(), "the connection handler thread is still running"); + } +} From b2c3d7743e1b8b6d3c99da2372d8ab2fd81d09ee Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Wed, 30 Sep 2026 15:45:31 +0300 Subject: [PATCH 2/4] [#1119] Drain each stopped transport on its own connections, and let queued notices out before shutting it down Review round 1 of #1125: - A TransportDrain tests and disconnects only the connections its own transport accepted, not every connection of the handler, so a handler that stops and starts listening again in one instance no longer keeps the old transport alive for the new one's clients. - Whether the server is shutting down is recorded in finalizeConnectionHandler instead of being read when the handler thread stops the listener, up to a second later, when an in-core restart may already have made the next server instance current. A drain ended on that path never touches the shutdown listeners. - Before shutdownNow(), a transport waits up to 2 seconds for the Grizzly connections it accepted to close, tracked by a ConnectionProbe, so a notice of disconnection queued behind unread data is not dropped. - Tests: two transports in one handler, a notice queued behind unread data, an in-core restart of the test server, a failed listen, num-request-handlers over the system property, the idle stop in the selector cases, the drain's registration, and the decoded notice (OID, result code, message). --- .../reactive/LDAPConnectionHandler2.java | 109 ++++++- ...APConnectionHandler2TransportTestCase.java | 307 +++++++++++++++++- 2 files changed, 391 insertions(+), 25 deletions(-) diff --git a/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java b/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java index dc1de10e0e..a1989df53b 100644 --- a/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java +++ b/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java @@ -72,6 +72,8 @@ import org.forgerock.opendj.server.config.server.LDAPConnectionHandlerCfg; import org.forgerock.util.Function; import org.forgerock.util.Options; +import org.glassfish.grizzly.Connection; +import org.glassfish.grizzly.ConnectionProbe; import org.glassfish.grizzly.memory.MemoryManager; import org.glassfish.grizzly.memory.PooledMemoryManager; import org.glassfish.grizzly.nio.transport.TCPNIOTransport; @@ -143,6 +145,41 @@ private static final class MemoryManagerHolder { private static final MemoryManager INSTANCE = new PooledMemoryManager(true); } + /** + * Tracks the connections a transport has accepted and not yet closed, down to the socket. A connection leaves the + * set of client connections as soon as it is disconnected, before its notice of disconnection is written, so only + * this tracking tells when shutting the transport down can no longer drop a write. + */ + private static final class OpenConnections extends ConnectionProbe.Adapter { + private final Collection> open = ConcurrentHashMap.newKeySet(); + + @Override + public void onAcceptEvent(Connection serverConnection, Connection clientConnection) { + open.add(clientConnection); + } + + @Override + public void onCloseEvent(Connection connection) { + if (open.remove(connection)) { + synchronized (this) { + notifyAll(); + } + } + } + + /** Waits until every accepted connection is closed, for at most the given time. */ + private synchronized void awaitClosed(long timeoutMs) throws InterruptedException { + final long deadlineNanos = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeoutMs); + while (!open.isEmpty()) { + final long remainingMs = TimeUnit.NANOSECONDS.toMillis(deadlineNanos - System.nanoTime()); + if (remainingMs <= 0) { + return; + } + wait(remainingMs); + } + } + } + /** * Keeps the transport of a stopped listener running for the connections it has accepted, so that disabling, * deleting or restarting the handler leaves them open, as the transport the SDK shares between listeners did. The @@ -150,21 +187,29 @@ private static final class MemoryManagerHolder { */ private final class TransportDrain implements ServerShutdownListener { private final TCPNIOTransport drained; + /** The client connections accepted by the drained transport, and only those. */ + private final Collection accepted; + private final OpenConnections open; private final AtomicBoolean finished = new AtomicBoolean(); + /** Whether this drain is registered as a shutdown listener, which only {@link #start()} does. */ + private volatile boolean registered; - private TransportDrain(TCPNIOTransport drained) { + private TransportDrain(TCPNIOTransport drained, Collection accepted, OpenConnections open) { this.drained = drained; + this.accepted = accepted; + this.open = open; } private void start() { drains.add(this); + registered = true; DirectoryServer.registerShutdownListener(this); connectionClosed(); } - /** Shuts the transport down if no connection of this handler is left. */ + /** Shuts the transport down if no connection it has accepted is left. */ private void connectionClosed() { - if (clientConnections.isEmpty() && finish()) { + if (accepted.isEmpty() && finish()) { // The last connection may be closed on a selector thread of the transport being shut down. new DirectoryThread(this::shutdownTransport, getShutdownListenerName()).start(); } @@ -180,7 +225,7 @@ public void processServerShutdown(LocalizableMessage reason) { if (finish()) { // Shutting the transport down closes every connection it accepted: end them as a server shutdown // first, as the legacy connection handler does, rather than let them fail as protocol errors. - for (ClientConnection clientConnection : clientConnections) { + for (ClientConnection clientConnection : accepted) { clientConnection.disconnect(DisconnectReason.SERVER_SHUTDOWN, true, reason); } shutdownTransport(); @@ -192,11 +237,22 @@ private boolean finish() { return false; } drains.remove(this); - DirectoryServer.deregisterShutdownListener(this); + if (registered) { + // An unregistered drain must not touch the listeners: during an in-core restart they may already + // belong to the next server instance, which has none yet. + DirectoryServer.deregisterShutdownListener(this); + } return true; } private void shutdownTransport() { + try { + // A notice of disconnection queued behind a response the client has not read yet is written once the + // client reads, and a shutdown drops it: give the connections a bounded time to take it and close. + open.awaitClosed(NOTICE_DELIVERY_TIMEOUT_MS); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } try { drained.shutdownNow(); } catch (IOException e) { @@ -213,6 +269,9 @@ private void shutdownTransport() { */ private static final String SELECTORS_PROPERTY = "org.forgerock.opendj.transport.selectors"; + /** How long a transport being shut down waits for its connections to write what they have queued and close. */ + private static final long NOTICE_DELIVERY_TIMEOUT_MS = 2000; + /** Default friendly name for the LDAP connection handler. */ private static final String DEFAULT_FRIENDLY_NAME = "LDAP Connection Handler"; @@ -227,6 +286,12 @@ private void shutdownTransport() { */ private TCPNIOTransport transport; + /** The client connections accepted by {@link #transport}, handed over with it to a {@link TransportDrain}. */ + private Collection transportConnections; + + /** The connections {@link #transport} has accepted and not yet closed, handed over with it. */ + private OpenConnections transportOpenConnections; + /** The transports of stopped listeners still serving the connections they accepted. */ private final Collection drains = new CopyOnWriteArrayList<>(); @@ -265,6 +330,12 @@ private void shutdownTransport() { /** The reason given to the client connections still open when the server shuts down. */ private volatile LocalizableMessage closeReason; + /** + * Whether the server was shutting down when this handler was finalized. Read then rather than when the listener + * stops, up to a second later: after an in-core restart the server instance is by then the next one. + */ + private volatile boolean serverShuttingDown; + /* Internal LDAP connection handler state */ /** Indicates whether this connection handler is enabled. */ @@ -437,6 +508,8 @@ private void configureSSL(LDAPConnectionHandlerCfg config) throws DirectoryExcep @Override public void finalizeConnectionHandler(LocalizableMessage finalizeReason) { + // DirectoryServer.shutDown sets its flag before it finalizes the connection handlers, on this thread. + serverShuttingDown = DirectoryServer.getInstance().isShuttingDown(); closeReason = finalizeReason; shutdownRequested = true; currentConfig.removeLDAPChangeListener(this); @@ -772,9 +845,11 @@ void stopListener() { logger.info(NOTE_CONNHANDLER_STOPPED_LISTENING, handlerName); } if (transport != null) { - final TransportDrain drain = new TransportDrain(transport); + final TransportDrain drain = new TransportDrain(transport, transportConnections, transportOpenConnections); transport = null; - if (DirectoryServer.getInstance().isShuttingDown()) { + transportConnections = null; + transportOpenConnections = null; + if (serverShuttingDown) { drain.processServerShutdown(closeReason); } else { drain.start(); @@ -782,7 +857,8 @@ void stopListener() { } } - private void connectionClosed(ClientConnection connection) { + private void connectionClosed(ClientConnection connection, Collection accepted) { + accepted.remove(connection); clientConnections.remove(connection); for (TransportDrain drain : drains) { drain.connectionClosed(); @@ -794,7 +870,7 @@ private void connectionClosed(ClientConnection connection) { * threads and its socket options follow its configuration rather than the JVM-wide settings of the transport the * SDK shares between listeners. */ - private TCPNIOTransport newTransport() throws IOException { + private TCPNIOTransport newTransport(OpenConnections openConnections) throws IOException { final TCPNIOTransport newTransport = TCPNIOTransportBuilder.newInstance() .setIOStrategy(SameThreadIOStrategy.getInstance()) .setSelectorThreadPoolConfig(ThreadPoolConfig.defaultConfig() @@ -807,6 +883,8 @@ private TCPNIOTransport newTransport() throws IOException { .build(); // As in the SDK server transport: fewer selector runners than selector threads cause deadlocks. newTransport.setSelectorRunnersCount(numRequestHandlers); + // Before the transport binds: a connection takes the probes of its transport when it is created. + newTransport.getConnectionMonitoringConfig().addProbes(openConnections); try { newTransport.start(); } catch (IOException e) { @@ -817,7 +895,11 @@ private TCPNIOTransport newTransport() throws IOException { } private void startListener() throws IOException { - transport = newTransport(); + final OpenConnections openConnections = new OpenConnections(); + final Collection accepted = ConcurrentHashMap.newKeySet(); + transport = newTransport(openConnections); + transportConnections = accepted; + transportOpenConnections = openConnections; listener = new LDAPListener( listenAddresses, new Function> LDAPClientContext clientContext) throws LdapException { final LDAPClientConnection2 conn = canAccept(clientContext); clientConnections.add(conn); + accepted.add(conn); logConnect(conn); clientContext.addListener(new LDAPClientContextEventListener() { @Override public void handleConnectionError(final LDAPClientContext context, final Throwable error) { - connectionClosed(conn); + connectionClosed(conn, accepted); } @Override public void handleConnectionDisconnected(final LDAPClientContext context, final ResultCode resultCode, String diagnosticMessage) { - connectionClosed(conn); + connectionClosed(conn, accepted); } @Override public void handleConnectionClosed(final LDAPClientContext context, final UnbindRequest unbindRequest) { - connectionClosed(conn); + connectionClosed(conn, accepted); } }); return new ReactiveHandler>() { diff --git a/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java index 3dd2b5d2ae..848f1aecf0 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java @@ -15,16 +15,17 @@ */ package org.opends.server.protocols.ldap; +import static org.opends.messages.CoreMessages.INFO_CONNHANDLER_CLOSED_BY_SHUTDOWN; import static org.testng.Assert.*; -import java.io.ByteArrayOutputStream; import java.io.InputStream; import java.lang.reflect.Field; +import java.net.InetSocketAddress; +import java.net.ServerSocket; import java.net.Socket; import java.net.SocketTimeoutException; import java.nio.channels.ServerSocketChannel; import java.nio.channels.SocketChannel; -import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collection; import java.util.List; @@ -33,15 +34,18 @@ import org.forgerock.opendj.reactive.LDAPConnectionHandler2; import org.forgerock.opendj.server.config.meta.LDAPConnectionHandlerCfgDefn; import org.forgerock.opendj.server.config.server.LDAPConnectionHandlerCfg; +import org.glassfish.grizzly.memory.Buffers; import org.glassfish.grizzly.nio.transport.TCPNIOConnection; import org.glassfish.grizzly.nio.transport.TCPNIOServerConnection; import org.glassfish.grizzly.nio.transport.TCPNIOTransport; import org.opends.server.DirectoryServerTestCase; import org.opends.server.TestCaseUtils; import org.opends.server.api.ClientConnection; +import org.opends.server.api.ConnectionHandler; import org.opends.server.api.ServerShutdownListener; import org.opends.server.core.DirectoryServer; import org.opends.server.extensions.InitializationUtils; +import org.opends.server.tools.LDAPReader; import org.opends.server.types.Entry; import org.testng.annotations.BeforeClass; import org.testng.annotations.DataProvider; @@ -66,6 +70,10 @@ public class LDAPConnectionHandler2TransportTestCase extends DirectoryServerTest private static final long TIMEOUT_MS = 10000; /** How long the connection of a stopped handler must stay open. */ private static final long KEEPS_OPEN_MS = 3000; + /** What is written at a time to fill the socket buffers of both ends and then the write queue of the server. */ + private static final int UNREAD_CHUNK = 64 * 1024; + /** Far more than any socket buffers hold. */ + private static final int MAX_UNREAD_BYTES = 64 * 1024 * 1024; @BeforeClass public void setUp() throws Exception @@ -162,6 +170,24 @@ public void selectorThreadsFollowTheSelectorsPropertyWhenNumRequestHandlersIsUns } } + @Test + public void numRequestHandlersTakesPrecedenceOverTheSelectorsProperty() throws Exception + { + final String saved = System.getProperty(SELECTORS_PROPERTY); + System.setProperty(SELECTORS_PROPERTY, "13"); + try + { + assertSelectorThreads(configuration(TestCaseUtils.findFreePort(), true, true, true, 3), 3); + } + finally + { + restore(saved); + } + } + + /** + * Decisive only on hosts with 6 processors or more: below that the count is 2, which a constant would give too. + */ @Test public void selectorThreadsAreChosenFromTheProcessorsWhenNothingIsSet() throws Exception { @@ -223,27 +249,275 @@ public void serverShutdownEndsTheConnectionsOfAStoppedHandler() throws Exception { acceptedConnection(handler); stop(handler); - final Collection drains = (Collection) field(handler, "drains"); - assertEquals(drains.size(), 1, "no transport is left serving the connection of the stopped handler"); + final ServerShutdownListener drain = onlyDrain(handler); + assertTrue(((Collection) field(DirectoryServer.getInstance(), "shutdownListeners")).contains(drain), + "the drain is not registered as a shutdown listener"); + + drain.processServerShutdown(STOP_REASON); + + assertNoticeOfDisconnection(client, STOP_REASON); + } + finally + { + stop(handler); + } + assertThreadsStop(threadPrefix); + } + + /** + * A notice of disconnection queued behind data the client has not read yet still reaches the client: the transport + * is not shut down before the client has read up to it. + */ + @Test + public void serverShutdownDeliversANoticeQueuedBehindUnreadData() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, true, true, true, 2)); + try (Socket client = new Socket()) + { + // Keeps what the client side holds small, so that the backlog is quick to read once the client reads. + client.setReceiveBufferSize(8 * 1024); + client.connect(new InetSocketAddress("127.0.0.1", port)); + final TCPNIOConnection accepted = acceptedConnection(handler); + final TCPNIOTransport transport = transport(handler); + // Written below the LDAP filters, as raw bytes the client will skip, until the socket buffers are full and some + // are left waiting in the write queue. + final byte[] chunk = new byte[UNREAD_CHUNK]; + int unread = 0; + while (accepted.getAsyncWriteQueue().spaceInBytes() == 0) + { + assertTrue(unread < MAX_UNREAD_BYTES, "the socket buffers took " + unread + " bytes"); + transport.getAsyncQueueIO().getWriter().write(accepted, Buffers.wrap(transport.getMemoryManager(), chunk), null); + unread += chunk.length; + // Let the selector move what it can into the socket. + Thread.sleep(50); + } + stop(handler); + final ServerShutdownListener drain = onlyDrain(handler); - ((ServerShutdownListener) drains.iterator().next()).processServerShutdown(STOP_REASON); + final Thread shutdown = new Thread(() -> drain.processServerShutdown(STOP_REASON)); + shutdown.start(); + // Let the notice join the queue while the client reads nothing. + Thread.sleep(300); client.setSoTimeout((int) TIMEOUT_MS); - final ByteArrayOutputStream received = new ByteArrayOutputStream(); final InputStream in = client.getInputStream(); - final byte[] buffer = new byte[256]; - for (int read; (read = in.read(buffer)) != -1;) + final byte[] buffer = new byte[64 * 1024]; + for (int left = unread; left > 0;) { - received.write(buffer, 0, read); + final int read = in.read(buffer, 0, Math.min(buffer.length, left)); + assertTrue(read > 0, "the connection ended with " + left + " unread bytes still queued"); + left -= read; } - assertTrue(new String(received.toByteArray(), StandardCharsets.ISO_8859_1).contains(NOTICE_OF_DISCONNECTION_OID), - "the connection was closed without a notice of disconnection"); + assertNoticeOfDisconnection(client, STOP_REASON); + shutdown.join(TIMEOUT_MS); + assertFalse(shutdown.isAlive(), "the drain is still waiting for a closed connection"); } finally { stop(handler); } - assertThreadsStop(threadPrefix); + } + + /** + * The server shutting down, here for an in-core restart, ends the connections of a handler that is still listening + * with a notice of disconnection. The handler stops its listener on its own thread, after the server has moved on. + */ + @Test + public void serverRestartEndsTheConnectionsOfAListeningHandler() throws Exception + { + try (Socket client = new Socket("127.0.0.1", TestCaseUtils.getServerLdapPort())) + { + awaitServerConnection(client.getLocalPort()); + + TestCaseUtils.restartServer(); + + assertNoticeOfDisconnection(client, INFO_CONNHANDLER_CLOSED_BY_SHUTDOWN.get()); + } + } + + /** + * A handler that stops listening and starts again, as it does when an SSL change it cannot use is applied and then + * undone, runs two transports. The drained one ends with its own connections, whatever the other one serves. + */ + @Test + public void drainEndsWithTheConnectionsOfItsOwnTransport() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, true, true, true, 2)); + try (Socket first = new Socket("127.0.0.1", port)) + { + acceptedConnection(handler); + final TCPNIOTransport drained = listenAgain(handler); + final TCPNIOTransport live = transport(handler); + try (Socket second = new Socket("127.0.0.1", port)) + { + awaitClientConnections(handler, 2); + + first.close(); + + awaitStopped(drained); + assertFalse(live.isStopped(), "the transport of the listening handler stopped"); + } + } + finally + { + stop(handler); + } + } + + /** The server shutting down ends, through a drain, only the connections of the drained transport. */ + @Test + public void drainEndsAtServerShutdownOnlyTheConnectionsOfItsOwnTransport() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, true, true, true, 2)); + try (Socket first = new Socket("127.0.0.1", port)) + { + acceptedConnection(handler); + final TCPNIOTransport drained = listenAgain(handler); + try (Socket second = new Socket("127.0.0.1", port)) + { + awaitClientConnections(handler, 2); + + onlyDrain(handler).processServerShutdown(STOP_REASON); + + assertNoticeOfDisconnection(first, STOP_REASON); + awaitStopped(drained); + second.setSoTimeout((int) KEEPS_OPEN_MS); + try + { + final int read = second.getInputStream().read(); + fail("the connection of the listening transport was " + (read == -1 ? "closed" : "written to")); + } + catch (SocketTimeoutException expected) + { + // still open + } + } + } + finally + { + stop(handler); + } + } + + /** A listener that cannot bind leaves no transport behind: the transport it started for it is shut down. */ + @Test + public void failedListenLeavesNoSelectorThreads() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = new LDAPConnectionHandler2(); + // Both initialization and the configuration check verify the port first: take it only afterwards. + handler.initializeConnectionHandler(DirectoryServer.getInstance().getServerContext(), + configuration(port, true, true, false, 2)); + try (ServerSocket taken = new ServerSocket()) + { + taken.bind(new InetSocketAddress("127.0.0.1", port)); + handler.start(); + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + while ((Boolean) field(handler, "enabled")) + { + assertTrue(System.currentTimeMillis() < deadline, "the handler did not give up listening"); + Thread.sleep(50); + } + assertThreadsStop(handler.getConnectionHandlerName() + " Request Handler"); + } + finally + { + stop(handler); + } + } + + /** + * Makes the handler stop listening and start again in the same instance, and returns the transport it stopped with. + */ + private static TCPNIOTransport listenAgain(LDAPConnectionHandler2 handler) throws Exception + { + final TCPNIOTransport stopped = transport(handler); + // What the handler does to itself when it cannot use an SSL change: its configuration still enables it. + setField(handler, "enabled", false); + awaitDrains(handler, 1); + setField(handler, "enabled", true); + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + // The transport starts before the listener binds: only the listener tells that the port accepts connections. + while (field(handler, "listener") == null) + { + assertTrue(System.currentTimeMillis() < deadline, "the handler did not listen again"); + Thread.sleep(50); + } + return stopped; + } + + private static ServerShutdownListener onlyDrain(LDAPConnectionHandler2 handler) throws Exception + { + final Collection drains = (Collection) field(handler, "drains"); + assertEquals(drains.size(), 1, "no transport is left serving the connections of the stopped listener"); + return (ServerShutdownListener) drains.iterator().next(); + } + + private static void awaitDrains(LDAPConnectionHandler2 handler, int expected) throws Exception + { + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + while (((Collection) field(handler, "drains")).size() != expected) + { + assertTrue(System.currentTimeMillis() < deadline, "the handler did not stop listening"); + Thread.sleep(50); + } + } + + private static void awaitClientConnections(LDAPConnectionHandler2 handler, int expected) throws Exception + { + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + while (handler.getClientConnections().size() != expected) + { + assertTrue(System.currentTimeMillis() < deadline, "the handler did not accept the connections"); + Thread.sleep(50); + } + } + + /** Waits for a connection handler of the server to accept the connection from the given client port. */ + private static void awaitServerConnection(int clientPort) throws Exception + { + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + while (true) + { + for (ConnectionHandler connectionHandler : DirectoryServer.getConnectionHandlers()) + { + for (ClientConnection connection : connectionHandler.getClientConnections()) + { + if (connection.getClientPort() == clientPort) + { + return; + } + } + } + assertTrue(System.currentTimeMillis() < deadline, "the server did not accept the connection"); + Thread.sleep(50); + } + } + + private static void awaitStopped(TCPNIOTransport transport) throws InterruptedException + { + final long deadline = System.currentTimeMillis() + TIMEOUT_MS; + while (!transport.isStopped()) + { + assertTrue(System.currentTimeMillis() < deadline, "the drained transport is still running"); + Thread.sleep(50); + } + } + + /** Reads the next message from the client, checks it is the expected notice of disconnection, then the end. */ + private static void assertNoticeOfDisconnection(Socket client, LocalizableMessage reason) throws Exception + { + client.setSoTimeout((int) TIMEOUT_MS); + final LDAPMessage message = new LDAPReader(client).readMessage(); + assertNotNull(message, "the connection was closed without a notice of disconnection"); + final ExtendedResponseProtocolOp notice = message.getExtendedResponseProtocolOp(); + assertEquals(notice.getOID(), NOTICE_OF_DISCONNECTION_OID); + assertEquals(notice.getResultCode(), LDAPResultCode.UNAVAILABLE, "result code"); + assertEquals(String.valueOf(notice.getErrorMessage()), reason.toString(), "diagnostic message"); + assertEquals(client.getInputStream().read(), -1, "the connection stayed open after the notice"); } /** Returns the prefix of the names of the selector threads of the handler, after checking that some run. */ @@ -268,6 +542,7 @@ private static void assertThreadsStop(String prefix) throws InterruptedException private static void assertSelectorThreads(LDAPConnectionHandlerCfg config, int expected) throws Exception { final LDAPConnectionHandler2 handler = start(config); + final String threadPrefix = selectorThreadPrefix(handler); try { final TCPNIOTransport transport = transport(handler); @@ -278,6 +553,7 @@ private static void assertSelectorThreads(LDAPConnectionHandlerCfg config, int e { stop(handler); } + assertThreadsStop(threadPrefix); } private static void restore(String selectors) @@ -338,6 +614,13 @@ private static Object field(Object target, String... names) throws Exception return value; } + private static void setField(Object target, String name, Object value) throws Exception + { + final Field field = target.getClass().getDeclaredField(name); + field.setAccessible(true); + field.set(target, value); + } + private static LDAPConnectionHandlerCfg configuration(int port, boolean keepAlive, boolean noDelay, boolean reuseAddress, Integer numRequestHandlers) throws Exception { From 5aa486c5cf470944740d451254b1aa0db6d59bde Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Wed, 30 Sep 2026 18:38:03 +0300 Subject: [PATCH 3/4] [#1119] Check the server of the handler when a drain starts, and bound how long a closed drain runs Review round 2 of #1125: - A handler remembers the server instance it was initialized for, and a TransportDrain checks whether that instance is shutting down, instead of a flag recorded at finalization. A handler disabled or deleted just before an in-core restart may stop its listener after shutDown notified its listeners, or after the next instance became current. Its drain then ends the connections with the shutdown notice and registers with no server. - A drain whose last connection closed while it registered deregisters itself, so a finished drain no longer stays a shutdown listener. - Tests: the in-core restart case reads the notice while the server restarts, since macOS drops unread data of a loopback connection closed for longer than net.inet.tcp.fin_timeout. New case: a drain of a handler whose server is shutting down while another instance is current. The queued-notice case checks that the drain ends within a second once its connection is closed. --- .../reactive/LDAPConnectionHandler2.java | 40 ++++++----- ...APConnectionHandler2TransportTestCase.java | 69 ++++++++++++++++++- 2 files changed, 89 insertions(+), 20 deletions(-) diff --git a/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java b/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java index a1989df53b..5ac4ba1900 100644 --- a/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java +++ b/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java @@ -18,6 +18,7 @@ package org.forgerock.opendj.reactive; import static java.util.Collections.*; +import static org.opends.messages.CoreMessages.INFO_CONNHANDLER_CLOSED_BY_SHUTDOWN; import static org.opends.messages.ProtocolMessages.*; import static org.opends.server.loggers.AccessLogger.logConnect; import static org.opends.server.util.ServerConstants.*; @@ -202,9 +203,22 @@ private TransportDrain(TCPNIOTransport drained, Collection acc private void start() { drains.add(this); - registered = true; - DirectoryServer.registerShutdownListener(this); - connectionClosed(); + if (!server.isShuttingDown()) { + // Registers with the server of this handler: another one is current only once this one shut down. + registered = true; + DirectoryServer.registerShutdownListener(this); + if (finished.get()) { + // The last connection closed while this drain registered, and finished it before it was registered. + DirectoryServer.deregisterShutdownListener(this); + } + } + if (server.isShuttingDown()) { + // The server began shutting down before this drain registered, and may have notified its listeners + // already: end the connections as it would have. Their handler is not finalized again by then. + processServerShutdown(INFO_CONNHANDLER_CLOSED_BY_SHUTDOWN.get()); + } else { + connectionClosed(); + } } /** Shuts the transport down if no connection it has accepted is left. */ @@ -327,14 +341,12 @@ private void shutdownTransport() { /** Indicates whether the Directory Server is in the process of shutting down. */ private volatile boolean shutdownRequested; - /** The reason given to the client connections still open when the server shuts down. */ - private volatile LocalizableMessage closeReason; - /** - * Whether the server was shutting down when this handler was finalized. Read then rather than when the listener - * stops, up to a second later: after an in-core restart the server instance is by then the next one. + * The server this handler was initialized for. The listener stops on the thread of the handler, up to a second + * after the handler is finalized: after an in-core restart the current server instance is by then the next one, + * so only this one tells whether the server of the connections is shutting down. */ - private volatile boolean serverShuttingDown; + private DirectoryServer server; /* Internal LDAP connection handler state */ @@ -508,9 +520,6 @@ private void configureSSL(LDAPConnectionHandlerCfg config) throws DirectoryExcep @Override public void finalizeConnectionHandler(LocalizableMessage finalizeReason) { - // DirectoryServer.shutDown sets its flag before it finalizes the connection handlers, on this thread. - serverShuttingDown = DirectoryServer.getInstance().isShuttingDown(); - closeReason = finalizeReason; shutdownRequested = true; currentConfig.removeLDAPChangeListener(this); @@ -669,6 +678,7 @@ public void initializeConnectionHandler(ServerContext serverContext, LDAPConnect friendlyName = config.name(); } + server = DirectoryServer.getInstance(); // Save this configuration for future reference. currentConfig = config; enabled = config.isEnabled(); @@ -849,11 +859,7 @@ void stopListener() { transport = null; transportConnections = null; transportOpenConnections = null; - if (serverShuttingDown) { - drain.processServerShutdown(closeReason); - } else { - drain.start(); - } + drain.start(); } } diff --git a/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java index 848f1aecf0..0650ed8b52 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java @@ -19,6 +19,7 @@ import static org.testng.Assert.*; import java.io.InputStream; +import java.lang.reflect.Constructor; import java.lang.reflect.Field; import java.net.InetSocketAddress; import java.net.ServerSocket; @@ -29,6 +30,10 @@ import java.util.ArrayList; import java.util.Collection; import java.util.List; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import org.forgerock.i18n.LocalizableMessage; import org.forgerock.opendj.reactive.LDAPConnectionHandler2; @@ -74,6 +79,11 @@ public class LDAPConnectionHandler2TransportTestCase extends DirectoryServerTest private static final int UNREAD_CHUNK = 64 * 1024; /** Far more than any socket buffers hold. */ private static final int MAX_UNREAD_BYTES = 64 * 1024 * 1024; + /** + * How soon a drain ends once its last connection is closed: well below the 2 s it waits at most for connections to + * close, which the notice queued behind unread data takes about 1.3 s to reach. + */ + private static final long CLOSED_DRAIN_ENDS_MS = 1000; @BeforeClass public void setUp() throws Exception @@ -310,7 +320,8 @@ public void serverShutdownDeliversANoticeQueuedBehindUnreadData() throws Excepti left -= read; } assertNoticeOfDisconnection(client, STOP_REASON); - shutdown.join(TIMEOUT_MS); + // The client has read everything and the server has closed: the drain must not wait out its bound. + shutdown.join(CLOSED_DRAIN_ENDS_MS); assertFalse(shutdown.isAlive(), "the drain is still waiting for a closed connection"); } finally @@ -321,7 +332,9 @@ public void serverShutdownDeliversANoticeQueuedBehindUnreadData() throws Excepti /** * The server shutting down, here for an in-core restart, ends the connections of a handler that is still listening - * with a notice of disconnection. The handler stops its listener on its own thread, after the server has moved on. + * with a notice of disconnection. The handler stops its listener on its own thread, which may wake before or after + * the next server instance is current: {@link #drainEndsTheConnectionsWhenTheServerOfItsHandlerIsShuttingDown()} + * covers the latter. */ @Test public void serverRestartEndsTheConnectionsOfAListeningHandler() throws Exception @@ -329,11 +342,61 @@ public void serverRestartEndsTheConnectionsOfAListeningHandler() throws Exceptio try (Socket client = new Socket("127.0.0.1", TestCaseUtils.getServerLdapPort())) { awaitServerConnection(client.getLocalPort()); + // Read while the server restarts: macOS resets a closed loopback connection after net.inet.tcp.fin_timeout + // (60 s), and drops what the client has not read yet. + final ExecutorService reader = Executors.newSingleThreadExecutor(); + try + { + final Future notice = reader.submit(() -> { + assertNoticeOfDisconnection(client, INFO_CONNHANDLER_CLOSED_BY_SHUTDOWN.get()); + return null; + }); + TestCaseUtils.restartServer(); + notice.get(TIMEOUT_MS, TimeUnit.MILLISECONDS); + } + finally + { + reader.shutdownNow(); + } + } + } + + /** + * A handler finalized while its server runs, as disabling or deleting it does, may stop its listener after that + * server has begun shutting down, and even after the next server instance is current. The drain then ends the + * connections with a notice of disconnection as the server shutdown would have, and registers with no server. + */ + @Test + public void drainEndsTheConnectionsWhenTheServerOfItsHandlerIsShuttingDown() throws Exception + { + final int port = TestCaseUtils.findFreePort(); + final LDAPConnectionHandler2 handler = start(configuration(port, true, true, true, 2)); + final String threadPrefix = selectorThreadPrefix(handler); + try (Socket client = new Socket("127.0.0.1", port)) + { + acceptedConnection(handler); + // The server of the handler is shutting down, and the current instance is another one that is not. + final Constructor newServer = DirectoryServer.class.getDeclaredConstructor(); + newServer.setAccessible(true); + final DirectoryServer previous = newServer.newInstance(); + setField(previous, "shuttingDown", true); + setField(handler, "server", previous); - TestCaseUtils.restartServer(); + stop(handler); assertNoticeOfDisconnection(client, INFO_CONNHANDLER_CLOSED_BY_SHUTDOWN.get()); + assertTrue(((Collection) field(handler, "drains")).isEmpty(), "the drain is still serving the connections"); + for (Object listener : (Collection) field(DirectoryServer.getInstance(), "shutdownListeners")) + { + assertNotEquals(((ServerShutdownListener) listener).getShutdownListenerName(), + "Transport drain of " + handler.getConnectionHandlerName(), "the drain registered with the current server"); + } } + finally + { + stop(handler); + } + assertThreadsStop(threadPrefix); } /** From a96d7969db4396823a4eae6463dff050818b9751 Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Wed, 30 Sep 2026 23:58:54 +0300 Subject: [PATCH 4/4] [#1119] Give the cases without SO_REUSEADDR a port a plain bind accepts The two cases whose handler has allow-tcp-reuse-address: false failed on the ubuntu/JDK 11 leg of both CI runs of #1125: the port check of initializeConnectionHandler could not bind 127.0.0.1:6552x. The port came from TestCaseUtils.findFreePort(), which checks with SO_REUSEADDR only, and every test class counts ports down from 65535 in a JVM of its own, so a port can still carry a connection an earlier class left in TIME_WAIT. Such a port accepts a bind with SO_REUSEADDR and refuses one without it. These cases now take a port that a plain bind on 127.0.0.1 accepts, and skip the others. --- ...APConnectionHandler2TransportTestCase.java | 34 +++++++++++++++++-- 1 file changed, 32 insertions(+), 2 deletions(-) diff --git a/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java index 0650ed8b52..84fe0f9710 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java @@ -18,9 +18,11 @@ import static org.opends.messages.CoreMessages.INFO_CONNHANDLER_CLOSED_BY_SHUTDOWN; import static org.testng.Assert.*; +import java.io.IOException; import java.io.InputStream; import java.lang.reflect.Constructor; import java.lang.reflect.Field; +import java.net.BindException; import java.net.InetSocketAddress; import java.net.ServerSocket; import java.net.Socket; @@ -141,7 +143,7 @@ public Object[][] reuseAddress() public void listenSocketUsesTheConfiguredReuseAddress(boolean reuseAddress) throws Exception { final LDAPConnectionHandler2 handler = - start(configuration(TestCaseUtils.findFreePort(), true, true, reuseAddress, 2)); + start(configuration(freePort(reuseAddress), true, true, reuseAddress, 2)); try { final Collection serverConnections = (Collection) field(field(handler, "listener"), "impl", "serverConnections"); @@ -469,7 +471,7 @@ public void drainEndsAtServerShutdownOnlyTheConnectionsOfItsOwnTransport() throw @Test public void failedListenLeavesNoSelectorThreads() throws Exception { - final int port = TestCaseUtils.findFreePort(); + final int port = freePort(false); final LDAPConnectionHandler2 handler = new LDAPConnectionHandler2(); // Both initialization and the configuration check verify the port first: take it only afterwards. handler.initializeConnectionHandler(DirectoryServer.getInstance().getServerContext(), @@ -492,6 +494,34 @@ public void failedListenLeavesNoSelectorThreads() throws Exception } } + /** + * Returns a free port that a listen socket with the given SO_REUSEADDR setting can bind to on 127.0.0.1. + * {@link TestCaseUtils#findFreePort()} checks its ports with SO_REUSEADDR only, and every test class counts them down + * from the same number in a JVM of its own: a port can still carry a connection a previous class left in TIME_WAIT, + * which refuses only a socket without SO_REUSEADDR. + */ + private static int freePort(boolean reuseAddress) throws IOException + { + while (true) + { + final int port = TestCaseUtils.findFreePort(); + if (reuseAddress) + { + return port; + } + try (ServerSocket probe = new ServerSocket()) + { + probe.setReuseAddress(false); + probe.bind(new InetSocketAddress("127.0.0.1", port)); + return port; + } + catch (BindException inUse) + { + // Try the next one: findFreePort() hands out each port once, and throws when none is left. + } + } + } + /** * Makes the handler stop listening and start again in the same instance, and returns the transport it stopped with. */