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/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java |  278 ++++++++++++++++++++++++++++++++++++++++++++++++++++++-
 1 files changed, 273 insertions(+), 5 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java b/opendj-server-legacy/src/main/java/org/forgerock/opendj/reactive/LDAPConnectionHandler2.java
index 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);
     }
 

--
Gitblit v1.10.0