From 0ad962f75f40b0d90fb958e3ca9689a788d2ae0e Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Thu, 01 Oct 2026 09:32:06 +0000
Subject: [PATCH] [#1119] Serve each LDAP connection handler with a transport of its own, built from its configuration (#1125)
---
opendj-server-legacy/pom.xml | 5
opendj-server-legacy/src/test/java/org/opends/server/protocols/ldap/LDAPConnectionHandler2TransportTestCase.java | 767 ++++++++++++++++++++++++++++++++++++++
opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListenerTestCase.java | 40 ++
opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/GrizzlyLDAPListener.java | 14
opendj-maven-plugin/src/main/resources/config/xml/org/forgerock/opendj/server/config/LDAPConnectionHandlerConfiguration.xml | 71 +++
opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java | 278 +++++++++++++
6 files changed, 1,165 insertions(+), 10 deletions(-)
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 19e36fe..127283d 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<TCPNIOTransport> GRIZZLY_TRANSPORT = Option.of(TCPNIOTransport.class, null);
+
private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
private final ReferenceCountedObject<TCPNIOTransport>.Reference transport;
private final Collection<TCPNIOServerConnection> serverConnections;
@@ -66,7 +74,7 @@
* @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 @@
final Function<LDAPClientContext,
ReactiveHandler<LDAPClientContext, LdapRequestEnvelope, Stream<Response>>,
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 89b7787..fae5bbd 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;
@@ -281,6 +284,43 @@
}
/**
+ * 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
* connection.
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 a21dbba..0e596cc 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 @@
<adm:property-reference name="listen-port" />
<adm:property-reference name="use-ssl" />
<adm:property-reference name="ssl-cert-nickname" />
- <adm:property-reference name="use-tcp-keep-alive" />
- <adm:property-reference name="use-tcp-no-delay" />
+ <adm:property name="use-tcp-keep-alive" advanced="true">
+ <adm:synopsis>
+ Indicates whether the
+ <adm:user-friendly-name />
+ should use TCP keep-alive.
+ </adm:synopsis>
+ <adm:description>
+ 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.
+ </adm:description>
+ <adm:requires-admin-action>
+ <adm:component-restart />
+ </adm:requires-admin-action>
+ <adm:default-behavior>
+ <adm:defined>
+ <adm:value>true</adm:value>
+ </adm:defined>
+ </adm:default-behavior>
+ <adm:syntax>
+ <adm:boolean />
+ </adm:syntax>
+ <adm:profile name="ldap">
+ <ldap:attribute>
+ <ldap:name>ds-cfg-use-tcp-keep-alive</ldap:name>
+ </ldap:attribute>
+ </adm:profile>
+ </adm:property>
+ <adm:property name="use-tcp-no-delay" advanced="true">
+ <adm:synopsis>
+ Indicates whether the
+ <adm:user-friendly-name />
+ should use TCP no-delay.
+ </adm:synopsis>
+ <adm:description>
+ 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.
+ </adm:description>
+ <adm:requires-admin-action>
+ <adm:component-restart />
+ </adm:requires-admin-action>
+ <adm:default-behavior>
+ <adm:defined>
+ <adm:value>true</adm:value>
+ </adm:defined>
+ </adm:default-behavior>
+ <adm:syntax>
+ <adm:boolean />
+ </adm:syntax>
+ <adm:profile name="ldap">
+ <ldap:attribute>
+ <ldap:name>ds-cfg-use-tcp-no-delay</ldap:name>
+ </ldap:attribute>
+ </adm:profile>
+ </adm:property>
<adm:property-reference name="allow-tcp-reuse-address" />
<adm:property name="key-manager-provider">
<adm:synopsis>
@@ -338,6 +402,9 @@
each client connection and used to buffer LDAP response messages data
when writing.
</adm:description>
+ <adm:requires-admin-action>
+ <adm:component-restart />
+ </adm:requires-admin-action>
<adm:default-behavior>
<adm:defined>
<adm:value>4096 bytes</adm:value>
diff --git a/opendj-server-legacy/pom.xml b/opendj-server-legacy/pom.xml
index c883bc7..81ba8ee 100644
--- a/opendj-server-legacy/pom.xml
+++ b/opendj-server-legacy/pom.xml
@@ -97,6 +97,11 @@
<dependency>
<groupId>org.openidentityplatform.opendj</groupId>
+ <artifactId>opendj-grizzly</artifactId>
+ </dependency>
+
+ <dependency>
+ <groupId>org.openidentityplatform.opendj</groupId>
<artifactId>opendj-ldap-toolkit</artifactId>
<version>${project.version}</version>
</dependency>
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 5f40bd1..5ac4ba1 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 @@
}
}
+ /**
+ * 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<Connection<?>> 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<ClientConnection> 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<ClientConnection> 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 @@
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<ClientConnection> 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<TransportDrain> drains = new CopyOnWriteArrayList<>();
+
/** The current configuration state. */
private LDAPConnectionHandlerCfg currentConfig;
@@ -152,9 +326,28 @@
/** 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 @@
// * 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 @@
friendlyName = config.name();
}
+ server = DirectoryServer.getInstance();
// Save this configuration for future reference.
currentConfig = config;
enabled = config.isEnabled();
@@ -502,6 +697,10 @@
// 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 @@
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<LocalizableMessage> unacceptableReasons) {
@@ -640,9 +854,58 @@
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<ClientConnection> 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<ClientConnection> accepted = ConcurrentHashMap.newKeySet();
+ transport = newTransport(openConnections);
+ transportConnections = accepted;
+ transportOpenConnections = openConnections;
listener = new LDAPListener(
listenAddresses,
new Function<LDAPClientContext,
@@ -653,23 +916,24 @@
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<LDAPClientContext, LdapRequestEnvelope, Stream<Response>>() {
@@ -682,7 +946,11 @@
}
}, 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 0000000..84fe0f9
--- /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<DirectoryServer> 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<String> 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<String> threadsNamed(String prefix)
+ {
+ final List<String> 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<String> 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");
+ }
+}
--
Gitblit v1.10.0