mirror of https://github.com/OpenIdentityPlatform/OpenDJ.git

Valery Kharseko
16 hours ago 0ad962f75f40b0d90fb958e3ca9689a788d2ae0e
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);
    }