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