| | |
| | | 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.*; |
| | |
| | | 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; |
| | |
| | | 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; |
| | |
| | | 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; |
| | |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * 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"; |
| | | |
| | |
| | | |
| | | 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; |
| | | |
| | |
| | | /** 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. */ |
| | |
| | | // * 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. |
| | |
| | | friendlyName = config.name(); |
| | | } |
| | | |
| | | server = DirectoryServer.getInstance(); |
| | | // Save this configuration for future reference. |
| | | currentConfig = config; |
| | | enabled = config.isEnabled(); |
| | |
| | | // 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())); |
| | |
| | | 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) { |
| | |
| | | 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, |
| | |
| | | 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>>() { |
| | |
| | | } |
| | | }, 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); |
| | | } |
| | | |