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..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.*; @@ -41,9 +42,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 +58,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 +73,14 @@ 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; +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 +137,155 @@ 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); + } + + /** + * 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 + * 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; + /** 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, Collection accepted, OpenConnections open) { + this.drained = drained; + this.accepted = accepted; + this.open = open; + } + + private void start() { + drains.add(this); + 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. */ + private void connectionClosed() { + 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(); + } + } + + @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 : accepted) { + clientConnection.disconnect(DisconnectReason.SERVER_SHUTDOWN, true, reason); + } + shutdownTransport(); + } + } + + private boolean finish() { + if (!finished.compareAndSet(false, true)) { + return false; + } + drains.remove(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) { + 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"; + + /** 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"; @@ -135,6 +294,21 @@ 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 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<>(); + /** The current configuration state. */ private LDAPConnectionHandlerCfg currentConfig; @@ -152,9 +326,28 @@ 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 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 DirectoryServer server; + /* Internal LDAP connection handler state */ /** Indicates whether this connection handler is enabled. */ @@ -280,7 +473,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. @@ -484,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(); @@ -502,6 +697,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 +760,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 +854,58 @@ void stopListener() { listener = null; logger.info(NOTE_CONNHANDLER_STOPPED_LISTENING, handlerName); } + if (transport != null) { + final TransportDrain drain = new TransportDrain(transport, transportConnections, transportOpenConnections); + transport = null; + transportConnections = null; + transportOpenConnections = null; + drain.start(); + } + } + + private void connectionClosed(ClientConnection connection, Collection accepted) { + accepted.remove(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(OpenConnections openConnections) 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); + // 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) { + newTransport.shutdownNow(); + throw e; + } + return newTransport; } private void startListener() throws IOException { + 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) { - clientConnections.remove(conn); + connectionClosed(conn, accepted); } @Override public void handleConnectionDisconnected(final LDAPClientContext context, final ResultCode resultCode, String diagnosticMessage) { - clientConnections.remove(conn); + connectionClosed(conn, accepted); } @Override public void handleConnectionClosed(final LDAPClientContext context, final UnbindRequest unbindRequest) { - clientConnections.remove(conn); + connectionClosed(conn, accepted); } }); return new ReactiveHandler>() { @@ -682,7 +946,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..84fe0f9710 --- /dev/null +++ b/opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java @@ -0,0 +1,767 @@ +/* + * 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.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; +import java.net.SocketTimeoutException; +import java.nio.channels.ServerSocketChannel; +import java.nio.channels.SocketChannel; +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; +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; +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; + /** 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; + /** + * 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 + { + 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(freePort(reuseAddress), 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 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 + { + 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 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); + + 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 InputStream in = client.getInputStream(); + final byte[] buffer = new byte[64 * 1024]; + for (int left = unread; left > 0;) + { + 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; + } + assertNoticeOfDisconnection(client, STOP_REASON); + // 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 + { + stop(handler); + } + } + + /** + * 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, which may wake before or after + * the next server instance is current: {@link #drainEndsTheConnectionsWhenTheServerOfItsHandlerIsShuttingDown()} + * covers the latter. + */ + @Test + public void serverRestartEndsTheConnectionsOfAListeningHandler() throws Exception + { + 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); + + 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); + } + + /** + * 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 = 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(), + 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); + } + } + + /** + * 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. + */ + 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. */ + 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); + final String threadPrefix = selectorThreadPrefix(handler); + try + { + final TCPNIOTransport transport = transport(handler); + assertEquals(transport.getSelectorRunnersCount(), expected, "selector runners"); + assertEquals(transport.getKernelThreadPoolConfig().getMaxPoolSize(), expected, "selector threads"); + } + finally + { + stop(handler); + } + assertThreadsStop(threadPrefix); + } + + 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 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 + { + 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"); + } +}