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

Valery Kharseko
21 hours ago 129fc4e34608d8212e62435f07283132e76ee1d6
[#983] Interrupt the listen thread after the ReplicaOfflineMsgs are forwarded (#987)
2 files modified
356 ■■■■ changed files
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java 26 ●●●● patch | view | raw | blame | history
opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java 330 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
@@ -1660,12 +1660,9 @@
      connectThread.interrupt();
    }
    // shutdown the listener thread
    // Stop accepting connections. Closing the socket is what ends the accept() of the listen
    // thread, and the loop of that thread already stops on its closed socket.
    close(listenSocket);
    if (listenThread != null)
    {
      listenThread.interrupt();
    }
    /*
     * Let the ReplicaOfflineMsgs a collocated DS sent be forwarded while every handler is still
@@ -1680,6 +1677,25 @@
     */
    awaitReplicaOfflineMsgsForwarded();
    /*
     * Only now interrupt the listen thread: the handshake of an incoming connection runs in it,
     * and an interrupt sent before the wait tears down a peer replication server whose handshake
     * is in flight - one of the very servers the message has to be forwarded to. Whether that
     * peer is already registered in the domain and waiting on the startup of its session, or
     * still owes its TopologyMsg - a receive no interrupt breaks, so the flag survives until the
     * handshake has registered it and reaches the startup of its session - the abort unregisters
     * it, and it is never told that the replica went offline.
     * <p>
     * The interrupt still precedes the shutdown of the domains, so a handshake which has not
     * finished by then is still aborted before its reader and its writer are started. Without it
     * such a handshake would register its peer after the domains were stopped and serve it: a
     * writer parked on a cursor over the closed changelog, a heartbeat to a server which is gone.
     */
    if (listenThread != null)
    {
      listenThread.interrupt();
    }
    for (ReplicationServerDomain domain : getReplicationServerDomains())
    {
      domain.shutdown();
opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java
@@ -33,6 +33,7 @@
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
@@ -50,6 +51,7 @@
import org.opends.server.replication.protocol.ReplicaOfflineMsg;
import org.opends.server.replication.protocol.ReplicationMsg;
import org.opends.server.replication.protocol.Session;
import org.opends.server.replication.protocol.StopMsg;
import org.opends.server.replication.protocol.TopologyMsg;
import org.opends.server.replication.protocol.WindowMsg;
import org.opends.server.replication.service.DSRSShutdownSync;
@@ -66,9 +68,12 @@
 * <p>
 * Most tests drive {@link DSRSShutdownSync} directly rather than through a collocated directory
 * server: the contract they pin is when the shutdown of the replication server waits, and how
 * long. {@link #thePeerReceivesTheReplicaOfflineMsgBeforeTheShutdownReturns()} and
 * long. {@link #thePeerReceivesTheReplicaOfflineMsgBeforeTheShutdownReturns()},
 * {@link #thePeerWhoseHandshakeIsInFlightIsStillToldTheReplicaWentOffline()} and
 * {@link #theShutdownWaitsForEveryPeerToBeToldTheReplicaWentOffline()} pin the outcome those
 * waits exist for, on peers connected through the real handshake.
 * waits exist for, on peers connected through the real handshake, and
 * {@link #thePeerWhoseHandshakeCompletesAfterTheShutdownIsStoppedNotServed()} pins what the
 * shutdown still does to a handshake which outlasts the wait.
 */
@SuppressWarnings("javadoc")
public class ReplicationServerShutdownSyncTest extends ReplicationTestCase
@@ -84,6 +89,13 @@
  private static final int HELD_BACK_RS_ID = 95;
  /** The peer replication server whose handshake is aborted while the message is pushed. */
  private static final int ABORTED_RS_ID = 96;
  /** The peer replication server whose handshake is still in flight when the shutdown starts. */
  private static final int HANDSHAKING_RS_ID = 97;
  /**
   * A replica which announces itself offline and whose message nobody can forward: it never
   * published one, so the shutdown spends its whole grace period waiting for it.
   */
  private static final int UNREACHABLE_DS_ID = 98;
  /** Send window a peer advertises when nothing has to hold its writer back. */
  private static final int PEER_WINDOW = 100;
  /**
@@ -193,7 +205,8 @@
      replicationServer =
          newReplicationServer(shutdownSync, "shutdownSyncDeliveryDb", 8226, replicationPort);
      broker = openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      peer = new FakePeerReplicationServer(replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID);
      peer = FakePeerReplicationServer.connected(
          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID);
      final ReplicationServerDomain domain =
          replicationServer.getReplicationServerDomain(baseDN, true);
@@ -232,6 +245,178 @@
  }
  /**
   * A peer replication server whose handshake is in flight when the shutdown starts must live
   * long enough to be told: it is one of the servers the ReplicaOfflineMsg is forwarded to, and
   * the handshake runs in the very listen thread the shutdown interrupts.
   */
  @Test
  public void thePeerWhoseHandshakeIsInFlightIsStillToldTheReplicaWentOffline() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final CountDownLatch shutdownIsWaiting = new CountDownLatch(1);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync()
    {
      @Override
      public void awaitReplicaOfflineMsgsForwarded(Collection<DN> baseDNs, long deadline)
      {
        shutdownIsWaiting.countDown();
        super.awaitReplicaOfflineMsgsForwarded(baseDNs, deadline);
      }
    };
    ReplicationServer replicationServer = null;
    ReplicationBroker broker = null;
    FakePeerReplicationServer connectedPeer = null;
    FakePeerReplicationServer handshakingPeer = null;
    Thread shutdown = null;
    try
    {
      final int replicationPort = TestCaseUtils.findFreePort();
      replicationServer =
          newReplicationServer(shutdownSync, "shutdownSyncHandshakeDb", 8234, replicationPort);
      broker = openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      connectedPeer = FakePeerReplicationServer.connected(
          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID);
      final ReplicationServerDomain domain =
          replicationServer.getReplicationServerDomain(baseDN, true);
      waitForConnectedReplicationServer(domain, REMOTE_RS_ID);
      /*
       * The second peer stops between the two phases of its handshake, which leaves the listen
       * thread blocked in the receive of the TopologyMsg it owes, the peer not yet registered on
       * the domain - the state in which the shutdown used to send the interrupt which, once the
       * handshake got to the startup of its session, tore the peer down. A message nobody can
       * forward holds the shutdown in its wait meanwhile, so what the connected peer does with
       * the message published below cannot end the wait before this one has been served.
       * <p>
       * Two clocks run from here, both far longer than what follows costs on the loopback: the
       * grace period of that announcement, which is all that holds the wait, and the connection
       * timeout of the replication server, which bounds its receive of the TopologyMsg. The
       * handshake, the registration and the forward must land inside the first, else the domains
       * are stopped and the peer reports the very message of the regression - hence the
       * assertion on the elapsed time below; the TopologyMsg must arrive inside the second, else
       * the replication server gives the peer up itself and completeHandshake() fails on its
       * own account.
       */
      handshakingPeer = FakePeerReplicationServer.handshaking(
          replicationPort, HANDSHAKING_RS_ID, baseDN, EMPTY_DN_GENID);
      final CSN unreachableCSN = newOfflineCSN(UNREACHABLE_DS_ID);
      final long startTime = System.nanoTime();
      shutdownSync.replicaOfflineMsgSent(baseDN, unreachableCSN);
      shutdown = newShutdownThread(replicationServer);
      shutdown.start();
      assertThat(shutdownIsWaiting.await(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS))
          .as("the shutdown never reached its wait for the ReplicaOfflineMsgs").isTrue();
      handshakingPeer.completeHandshake();
      final Future<ReplicaOfflineMsg> received = handshakingPeer.receive(ReplicaOfflineMsg.class);
      waitForConnectedReplicationServer(domain, HANDSHAKING_RS_ID);
      final CSN offlineCSN = newOfflineCSN();
      shutdownSync.replicaOfflineMsgSent(baseDN, offlineCSN);
      broker.publish(new ReplicaOfflineMsg(offlineCSN));
      final ReplicaOfflineMsg forwarded = received.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS);
      final long elapsed = elapsedMillis(startTime);
      assertThat(forwarded)
          .as("the peer which was handshaking when the shutdown started was never told that the "
              + "replica went offline, its read ended with: %s", handshakingPeer.failure())
          .isNotNull();
      assertThat(forwarded.getCSN().getServerId()).isEqualTo(LOCAL_DS_ID);
      assertThat(elapsed)
          .as("the handshake and the forward must land inside the grace period which holds the "
              + "wait, or the pass says nothing about the order of the shutdown")
          .isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
      /*
       * The announcement has done its job. An entry which was never queued for anybody is
       * released by the first forward reported for it, so end the wait now rather than sit out
       * the rest of its grace period in the join below.
       */
      shutdownSync.replicaOfflineMsgForwarded(baseDN, unreachableCSN, HANDSHAKING_RS_ID);
    }
    finally
    {
      joinQuietly(shutdown);
      closeQuietly(handshakingPeer);
      closeQuietly(connectedPeer);
      stop(broker);
      removeQuietly(replicationServer);
    }
  }
  /**
   * The other half of the order of the shutdown: the interrupt still precedes the shutdown of
   * the domains, so a handshake which completes only once the shutdown is over is stopped rather
   * than served. Without the interrupt such a peer registers after the domains were stopped and
   * nothing ever stops it: its session, reader, writer and heartbeat are started, and its writer
   * parks forever on a cursor over the changelog the shutdown has closed - five threads and a
   * peer heartbeating a replication server which no longer exists, on a server which was merely
   * removed from the configuration.
   */
  @Test
  public void thePeerWhoseHandshakeCompletesAfterTheShutdownIsStoppedNotServed() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    ReplicationServer replicationServer = null;
    FakePeerReplicationServer handshakingPeer = null;
    try
    {
      final int replicationPort = TestCaseUtils.findFreePort();
      replicationServer = newReplicationServer(
          new DSRSShutdownSync(), "shutdownSyncLateHandshakeDb", 8235, replicationPort);
      final ReplicationServerDomain domain =
          replicationServer.getReplicationServerDomain(baseDN, true);
      handshakingPeer = FakePeerReplicationServer.handshaking(
          replicationPort, HANDSHAKING_RS_ID, baseDN, EMPTY_DN_GENID);
      /*
       * Nothing is pending, so the wait returns at once and the interrupt lands on the handshake
       * in flight - on a receive it does not break, where the flag survives until the startup of
       * the session. The handshake is completed once the shutdown has returned, inside the
       * connection timeout which bounds that receive: past it the replication server gives the
       * peer up on its own, without a StopMsg, and completeHandshake() fails on a closed session
       * rather than letting the case pass for the wrong reason.
       */
      replicationServer.shutdown();
      handshakingPeer.completeHandshake();
      final Future<StopMsg> stopped = handshakingPeer.receive(StopMsg.class);
      StopMsg stopMsg = null;
      try
      {
        stopMsg = stopped.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS);
      }
      catch (TimeoutException e)
      {
        // A served peer keeps receiving heartbeats: its read neither ends nor yields a StopMsg.
      }
      assertThat(stopMsg)
          .as("the peer whose handshake completed after the shutdown was served instead of "
              + "stopped: no StopMsg within %s ms, its read ended with: %s",
              SOCKET_TIMEOUT_MS, handshakingPeer.failure())
          .isNotNull();
      // abortStart() sends the StopMsg when it closes the session, and unregisters the peer after
      newConnectionTimer().repeatUntilSuccess(new TestTimer.CallableVoid()
      {
        @Override
        public void call() throws Exception
        {
          assertThat(domain.getConnectedRSs())
              .as("the peer whose handshake was aborted by the shutdown stayed registered")
              .doesNotContainKey(HANDSHAKING_RS_ID);
        }
      });
    }
    finally
    {
      closeQuietly(handshakingPeer);
      removeQuietly(replicationServer);
    }
  }
  /**
   * Only a peer replication server learning about the offline replica ends the wait.
   * ReplicationServerDomain.put() never queues a ReplicaOfflineMsg for a directory server, but
   * the changelog cursor of a directory server which is catching up synthesizes one from the
@@ -316,9 +501,9 @@
          newReplicationServer(shutdownSync, "shutdownSyncEveryPeerDb", 8229, replicationPort);
      broker =
          openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      peer = new FakePeerReplicationServer(
      peer = FakePeerReplicationServer.connected(
          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID, PEER_WINDOW);
      heldBackPeer = new FakePeerReplicationServer(
      heldBackPeer = FakePeerReplicationServer.connected(
          replicationPort, HELD_BACK_RS_ID, baseDN, EMPTY_DN_GENID, HELD_BACK_PEER_WINDOW);
      final ReplicationServerDomain domain =
@@ -412,7 +597,7 @@
          newReplicationServer(shutdownSync, "shutdownSyncDroppedMsgDb", 8230, replicationPort);
      broker =
          openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      peer = new FakePeerReplicationServer(
      peer = FakePeerReplicationServer.connected(
          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID, HELD_BACK_PEER_WINDOW);
      final ReplicationServerDomain domain =
@@ -493,7 +678,8 @@
          shutdownSync, "shutdownSyncAbortedHandshakeDb", 8231, replicationPort);
      broker =
          openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      peer = new FakePeerReplicationServer(replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID);
      peer = FakePeerReplicationServer.connected(
          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID);
      final ReplicationServerDomain domain =
          replicationServer.getReplicationServerDomain(baseDN, true);
@@ -569,9 +755,9 @@
          shutdownSync, "shutdownSyncDisconnectedPeerDb", 8232, replicationPort);
      broker =
          openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      peer = new FakePeerReplicationServer(
      peer = FakePeerReplicationServer.connected(
          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID, PEER_WINDOW);
      heldBackPeer = new FakePeerReplicationServer(
      heldBackPeer = FakePeerReplicationServer.connected(
          replicationPort, HELD_BACK_RS_ID, baseDN, EMPTY_DN_GENID, HELD_BACK_PEER_WINDOW);
      final ReplicationServerDomain domain =
@@ -643,7 +829,8 @@
          shutdownSync, "shutdownSyncOtherGenerationIdDb", 8233, replicationPort);
      broker =
          openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      peer = new FakePeerReplicationServer(replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID);
      peer = FakePeerReplicationServer.connected(
          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID);
      final ReplicationServerDomain domain =
          replicationServer.getReplicationServerDomain(baseDN, true);
@@ -900,12 +1087,15 @@
   * The registration is not the end of the handshake: {@code startFromRemoteRS()} puts the
   * handler in {@code connectedRSs} before it calls {@code finalizeStart()}, which is what
   * starts the reader and the writer, and the whole handshake runs in the listen thread of the
   * replication server. {@link ReplicationServer#shutdown()} interrupts that thread before it
   * waits for the ReplicaOfflineMsgs, so a shutdown triggered while the handshake is still in
   * {@code Session.waitForStartup()} aborts it: the session is closed and the handler
   * replication server. {@link ReplicationServer#shutdown()} interrupts that thread once it is
   * done waiting for the ReplicaOfflineMsgs, so a handshake which is still in
   * {@code Session.waitForStartup()} by then is aborted: the session is closed and the handler
   * unregistered, and the peer is gone before the message could be queued for it, let alone
   * forwarded. Issue #821 recorded that same window, from the dead handler it used to leave
   * behind.
   * behind; {@link #thePeerWhoseHandshakeIsInFlightIsStillToldTheReplicaWentOffline()} pins the
   * part of it the wait now covers, and
   * {@link #thePeerWhoseHandshakeCompletesAfterTheShutdownIsStoppedNotServed()} the part the
   * shutdown still aborts.
   * <p>
   * The listen thread serves one handshake at a time, so a peer which is past that window also
   * puts every connection accepted before it - the collocated directory server of these tests
@@ -970,6 +1160,18 @@
        .toTimer();
  }
  private Thread newShutdownThread(final ReplicationServer replicationServer)
  {
    return new Thread(new Runnable()
    {
      @Override
      public void run()
      {
        replicationServer.shutdown();
      }
    }, "ReplicationServerShutdownSyncTest shutdown");
  }
  private Thread newForwarderThread(final DSRSShutdownSync shutdownSync, final DN baseDN,
      final CSN offlineCSN)
  {
@@ -1054,7 +1256,13 @@
  /** The CSN of a message the collocated replica announces, as PendingChanges generates it. */
  private static CSN newOfflineCSN()
  {
    return new CSNGenerator(LOCAL_DS_ID, 0).newCSN();
    return newOfflineCSN(LOCAL_DS_ID);
  }
  /** The CSN of a message the provided replica announces. */
  private static CSN newOfflineCSN(int serverId)
  {
    return new CSNGenerator(serverId, 0).newCSN();
  }
  private static boolean sleepQuietly(long millis)
@@ -1269,12 +1477,25 @@
  }
  /**
   * A peer replication server which connects to the replication server under test and completes
   * the handshake, so that the handler it leaves behind on the domain has a real writer and can
   * actually forward what the domain pushes to it.
   * A peer replication server which connects to the replication server under test.
   * {@link #connected(int, int, DN, long)} completes the handshake, so that the handler it leaves
   * behind on the domain has a real writer and can actually forward what the domain pushes to
   * it.
   * <p>
   * {@link #handshaking(int, int, DN, long)} stops between the two phases of the handshake
   * instead, which leaves the listen thread of the replication server blocked in the receive of
   * the TopologyMsg this peer owes it, before the handshake has registered the peer on the
   * domain. That pause is bounded by the socket timeout of the handshake, so
   * {@link #completeHandshake()} must follow shortly: a peer which stays silent for longer is
   * given up on by the replication server itself.
   */
  private static final class FakePeerReplicationServer
  {
    private static final byte GROUP_ID = 1;
    private final int serverId;
    private final long generationId;
    private final String serverURL;
    private final Session session;
    private final ExecutorService reader = Executors.newSingleThreadExecutor();
    /**
@@ -1285,42 +1506,73 @@
     */
    private final AtomicReference<Exception> failure = new AtomicReference<>();
    FakePeerReplicationServer(int replicationPort, int serverId, DN baseDN, long generationId)
        throws Exception
    /** A peer which has completed its handshake and is served by a writer of the domain. */
    static FakePeerReplicationServer connected(
        int replicationPort, int serverId, DN baseDN, long generationId) throws Exception
    {
      this(replicationPort, serverId, baseDN, generationId, PEER_WINDOW);
      return connected(replicationPort, serverId, baseDN, generationId, PEER_WINDOW);
    }
    FakePeerReplicationServer(int replicationPort, int serverId, DN baseDN, long generationId,
        int windowSize) throws Exception
    /** A connected peer which advertises the given send window to the replication server. */
    static FakePeerReplicationServer connected(int replicationPort, int serverId, DN baseDN,
        long generationId, int windowSize) throws Exception
    {
      final FakePeerReplicationServer peer = new FakePeerReplicationServer(
          replicationPort, serverId, baseDN, generationId, windowSize);
      boolean handshaken = false;
      try
      {
        peer.completeHandshake();
        handshaken = true;
      }
      finally
      {
        if (!handshaken)
        {
          // The caller has no handle on this peer yet, so nothing else would close it.
          peer.close();
        }
      }
      return peer;
    }
    /** A peer whose handshake stops after its first phase, before it sends its TopologyMsg. */
    static FakePeerReplicationServer handshaking(
        int replicationPort, int serverId, DN baseDN, long generationId) throws Exception
    {
      return new FakePeerReplicationServer(
          replicationPort, serverId, baseDN, generationId, PEER_WINDOW);
    }
    private FakePeerReplicationServer(int replicationPort, int serverId, DN baseDN,
        long generationId, int windowSize) throws Exception
    {
      this.serverId = serverId;
      this.generationId = generationId;
      final Socket socket = new Socket();
      Session newSession = null;
      boolean handshaken = false;
      String newServerURL = null;
      boolean started = false;
      try
      {
        socket.setTcpNoDelay(true);
        socket.connect(new InetSocketAddress("127.0.0.1", replicationPort), SOCKET_TIMEOUT_MS);
        newSession = getReplSessionSecurity().createClientSession(socket, SOCKET_TIMEOUT_MS);
        final String serverURL = "127.0.0.1:" + socket.getLocalPort();
        final byte groupId = (byte) 1;
        newSession.publish(new ReplServerStartMsg(serverId, serverURL, baseDN, windowSize,
            new ServerState(), generationId, false, groupId, 5000));
        newServerURL = "127.0.0.1:" + socket.getLocalPort();
        newSession.publish(new ReplServerStartMsg(serverId, newServerURL, baseDN, windowSize,
            new ServerState(), generationId, false, GROUP_ID, 5000));
        final ReplServerStartMsg inStartMsg =
            waitForSpecificMsg(newSession, ReplServerStartMsg.class);
        if (!inStartMsg.getSSLEncryption())
        {
          newSession.stopEncryption();
        }
        newSession.publish(new TopologyMsg(null,
            newArrayList(new RSInfo(serverId, serverURL, generationId, groupId, 1))));
        waitForSpecificMsg(newSession, TopologyMsg.class);
        handshaken = true;
        started = true;
      }
      finally
      {
        if (!handshaken)
        if (!started)
        {
          // The caller has no handle on this peer yet, so nothing else would close it.
          reader.shutdownNow();
@@ -1334,10 +1586,22 @@
          }
        }
      }
      serverURL = newServerURL;
      session = newSession;
    }
    /**
     * Runs the second phase of the handshake, the one which registers this peer on the domain of
     * the replication server.
     */
    void completeHandshake() throws Exception
    {
      session.publish(new TopologyMsg(null,
          newArrayList(new RSInfo(serverId, serverURL, generationId, GROUP_ID, 1))));
      waitForSpecificMsg(session, TopologyMsg.class);
    }
    /**
     * Returns the first message of the given type this peer receives, or null if its session
     * ends first.
     */