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

Valery Kharseko
10 hours ago bed2b17a444ee960268da0c05d1ca5365bdf2d98
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
@@ -106,7 +106,8 @@
  private static final int LISTEN_PORT_PROBE_TIMEOUT_MS = 200;
  private volatile ServerSocket listenSocket;
  private Thread listenThread;
  /** Volatile like its socket above: a port change reads it from the configuration thread. */
  private volatile Thread listenThread;
  private Thread connectThread;
  /** The current configuration of this replication server. */
@@ -130,8 +131,6 @@
  private boolean externalChangelogRegistered;
  private final AtomicBoolean shutdown = new AtomicBoolean();
  /** Written by the thread applying a configuration change, read by the listen thread. */
  private volatile boolean stopListen;
  private final ReplSessionSecurity replSessionSecurity;
  private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
@@ -155,6 +154,15 @@
   */
  static final AtomicInteger listenPortBindFailures = new AtomicInteger();
  /**
   * Number of listen port changes whose wait for the previous listen thread was interrupted.
   * <p>
   * This is required for unit testing: a wait which is interrupted and a wait which is over
   * before it starts leave this replication server in the same state, so nothing else tells
   * a test that it exercised the interruption instead of passing over it.
   */
  static final AtomicInteger interruptedListenThreadStops = new AtomicInteger();
  /** Monitors for synchronizing domain creation with the connect thread. */
  private final Object domainTicketLock = new Object();
  private final Object connectThreadLock = new Object();
@@ -260,15 +268,22 @@
   * This thread accept incoming connections on the replication server
   * ports from other replication servers or from LDAP servers
   * and spawn further thread responsible for handling those connections
   * <p>
   * The socket is the one this thread was created for, not the one this replication server
   * currently listens on: closing it is what stops this thread, and it is then the only
   * thread it stops, even when another one is already listening on another port.
   *
   * @param socket
   *          the bound socket this thread accepts connections on
   */
  void runListen()
  void runListen(ServerSocket socket)
  {
    logger.info(NOTE_REPLICATION_SERVER_LISTENING,
        getServerId(),
        listenSocket.getInetAddress().getHostAddress(),
        listenSocket.getLocalPort());
        socket.getInetAddress().getHostAddress(),
        socket.getLocalPort());
    while (!shutdown.get() && !stopListen)
    while (!shutdown.get() && !socket.isClosed())
    {
      // Wait on the replicationServer port.
      // Read incoming messages and create LDAP or ReplicationServer listener
@@ -279,7 +294,7 @@
        Socket newSocket = null;
        try
        {
          newSocket = listenSocket.accept();
          newSocket = socket.accept();
          newSocket.setTcpNoDelay(true);
          newSocket.setKeepAlive(true);
          int timeoutMS = MultimasterReplication.getConnectionTimeoutMS();
@@ -486,8 +501,8 @@
   * Initialization function for the replicationServer.
   *
   * @throws ConfigException
   *           when the replication server cannot be started, in particular when its listen
   *           port cannot be bound.
   *           when the replication server cannot be started, in particular when its changelog
   *           cannot be read or when its listen port cannot be bound.
   */
  private void initialize() throws ConfigException
  {
@@ -495,9 +510,13 @@
    try
    {
      // Assigned before the changelog is opened: the monitor instance name of the domains it
      // restores, and of their changelogs, embeds it, and a provider registered under a name
      // which later changes can never be deregistered again.
      setServerURL();
      this.changelogDB.initializeDB();
      setServerURL();
      // Assigned before the threads are created, so that a failure below still releases it.
      listenSocket = bindListenPort(getReplicationPort());
@@ -520,6 +539,12 @@
      {
        logger.trace("RS " + getMonitorInstanceName() + " successfully initialized");
      }
    } catch (ChangelogException e)
    {
      // A replication server which cannot read its changelog is as dead as one which cannot
      // bind its listen port (issue #802). The message already names the changelog directory.
      logger.traceException(e);
      throw new ConfigException(e.getMessageObject(), e);
    } catch (UnknownHostException e)
    {
      // Not logged here: the caller reports the ConfigException, logging it once.
@@ -527,8 +552,8 @@
      throw new ConfigException(ERR_UNKNOWN_HOSTNAME.get(), e);
    } catch (IOException e)
    {
      // A replication server whose listen port is not bound is dead: every consumer would
      // otherwise only learn about it as a "connection refused" somewhere else.
      // Every consumer would otherwise only learn about it as a "connection refused"
      // somewhere else (issue #792).
      logger.traceException(e);
      throw new ConfigException(bindFailureMessage(getReplicationPort(), e), e);
    }
@@ -582,28 +607,29 @@
  }
  /**
   * Stops the listen thread and releases the listen port.
   * Stops the listen thread of the provided listen socket and releases that port.
   * <p>
   * Both are the ones the thread was started on rather than the current ones, so this stops
   * that thread only, even when another one is already listening on another port.
   *
   * @param socket
   *          the listen socket to close, which is what stops its thread
   * @param thread
   *          the listen thread of that socket, {@code null} when it was never started
   * @throws InterruptedException
   *           if this thread is interrupted while waiting for the listen thread to stop
   */
  private void stopListenThread() throws InterruptedException
  private void stopListenThread(ServerSocket socket, Thread thread) throws InterruptedException
  {
    stopListen = true;
    close(listenSocket);
    if (listenThread != null)
    close(socket);
    if (thread != null)
    {
      listenThread.join();
      listenThread = null;
      thread.join();
    }
  }
  /**
   * Starts a listen thread on the provided listen socket.
   * <p>
   * {@code stopListen} is only cleared here, i.e. once the listen port is bound, so a
   * failure to bind leaves this replication server consistently stopped rather than with a
   * listen thread which would spin on a closed socket.
   *
   * @param boundListenSocket
   *          the bound socket the listen thread will accept connections on
@@ -611,30 +637,40 @@
  private void startListenThread(ServerSocket boundListenSocket)
  {
    listenSocket = boundListenSocket;
    stopListen = false;
    listenThread = new ReplicationServerListenThread(this);
    listenThread = new ReplicationServerListenThread(this, boundListenSocket);
    listenThread.start();
  }
  /**
   * Switches the listen port to the one of the provided configuration.
   * <p>
   * The new port is bound while the current one is still open and serving, so a failure
   * leaves this replication server listening on its current port, with its current
   * configuration: there is nothing to roll back, and no window during which this
   * replication server advertises a port that nothing listens to.
   * The new port is bound, and its listen thread started, while the current one is still
   * open and serving. A failure therefore leaves this replication server listening on its
   * current port, with its current configuration: there is nothing to roll back, and there
   * is no window during which this replication server listens on no port at all.
   * <p>
   * The trade is a window during which both ports accept, so a peer which connects to the
   * previous port just before it is released gets a session which outlives the change. That
   * is the deliberate inverse of a window during which nothing listens at all.
   *
   * @param newConfig
   *          the configuration being applied, whose listen port differs from the current one
   * @param ccr
   *          the result of the configuration change, to which a failure is added
   * @param listenThreadStopInterrupted
   *          set when the wait for the previous listen thread was interrupted, instead of
   *          restoring the interrupt status here: the caller restores it once the rest of
   *          the change, some of which is interruptible, has run
   * @return {@code true} when this replication server listens on the new port, in which
   *         case {@code newConfig} has become its configuration
   */
  private boolean switchListenPort(ReplicationServerCfg newConfig, ConfigChangeResult ccr)
  private boolean switchListenPort(ReplicationServerCfg newConfig, ConfigChangeResult ccr,
      AtomicBoolean listenThreadStopInterrupted)
  {
    final ReplicationServerCfg previousConfig = this.config;
    final String previousServerURL = serverURL;
    final ServerSocket previousListenSocket = listenSocket;
    final Thread previousListenThread = listenThread;
    final int newPort = newConfig.getReplicationPort();
    ServerSocket newListenSocket = null;
    try
@@ -645,12 +681,28 @@
      this.config = newConfig;
      setServerURL();
      stopListenThread();
      // The new port is served before the current one is released, so that this replication
      // server is never left with no listener at all, whatever happens next.
      startListenThread(newListenSocket);
      newListenSocket = null;
      // In step with getReplicationPort(), which answers the new port from here on: the
      // wait below blocks for as long as the previous thread takes to serve its current
      // connection, and localPorts must not trail it for that whole window.
      localPorts.remove(previousConfig.getReplicationPort());
      localPorts.add(newPort);
      try
      {
        stopListenThread(previousListenSocket, previousListenThread);
      }
      catch (InterruptedException e)
      {
        // The previous port is already released and its thread stops on its own as soon as
        // it wakes up on its closed socket: only the wait for it was cut short.
        interruptedListenThreadStops.incrementAndGet();
        listenThreadStopInterrupted.set(true);
        logger.traceException(e);
      }
      return true;
    }
    catch (UnknownHostException e)
@@ -666,15 +718,6 @@
      ccr.setResultCode(ResultCode.OPERATIONS_ERROR);
      ccr.addMessage(bindFailureMessage(newPort, e));
    }
    catch (InterruptedException e)
    {
      // The previous listen thread may still be running, so do not hand it a new socket:
      // stopListen is still set, which makes that thread stop as soon as it wakes up.
      Thread.currentThread().interrupt();
      logger.traceException(e);
      ccr.setResultCode(ResultCode.OPERATIONS_ERROR);
      ccr.addMessage(ERR_COULD_NOT_STOP_LISTEN_THREAD.get(getExceptionMessage(e)));
    }
    // The failure is reported through the ConfigChangeResult, which the configuration
    // handler logs: nothing of the new configuration was applied.
    this.config = previousConfig;
@@ -754,6 +797,22 @@
    {
      listenThread.interrupt();
    }
    // Opening the changelog restores one domain per domain it holds, and each of them starts
    // its threads and registers its monitor provider: a failure after that point, such as a
    // listen port which cannot be bound, would otherwise leave them behind. Shut them down
    // before the changelog they write to, and one unchecked exception at a time: the changelog
    // this one is built on is known to be broken, and what follows still has to run.
    for (ReplicationServerDomain domain : getReplicationServerDomains())
    {
      try
      {
        domain.shutdown();
      }
      catch (RuntimeException ignored)
      {
        logger.traceException(ignored);
      }
    }
    shutdownExternalChangelog();
    if (this.changelogDB != null)
    {
@@ -1195,8 +1254,9 @@
    // done first, and the new port is bound before the current one is released, so that a
    // change which cannot be applied leaves this replication server as it was, instead of
    // half configured and, worse, without any listener.
    final AtomicBoolean listenThreadStopInterrupted = new AtomicBoolean();
    if (configuration.getReplicationPort() != oldConfig.getReplicationPort()
        && !switchListenPort(configuration, ccr))
        && !switchListenPort(configuration, ccr, listenThreadStopInterrupted))
    {
      return ccr;
    }
@@ -1262,6 +1322,15 @@
    {
      ccr.setAdminActionRequired(true);
    }
    // The interrupt which cut short the wait for the previous listen thread, deferred by
    // switchListenPort(): restored only now, because the steps above include interruptible
    // ones — stopping the handlers of removed replication servers locks interruptibly —
    // which an interrupt status left set would have failed while the change reports SUCCESS.
    if (listenThreadStopInterrupted.get())
    {
      Thread.currentThread().interrupt();
    }
    return ccr;
  }