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

Valery Kharseko
11 hours ago 995d54a8963d87335196b86d8d4fc1b252a38ee4
[#900] Time the ReplicaOfflineMsg grace period per replica, and spend it where the message can still be forwarded (#919)
6 files modified
2 files added
1421 ■■■■■ changed files
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java 4 ●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/PendingChanges.java 6 ●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java 44 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerDomain.java 10 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java 30 ●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/service/DSRSShutdownSync.java 241 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java 793 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/test/java/org/opends/server/replication/service/DSRSShutdownSyncTest.java 293 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
@@ -1975,8 +1975,8 @@
  @Override
  public void publishReplicaOfflineMsg()
  {
    pendingChanges.putReplicaOfflineMsg();
    dsrsShutdownSync.replicaOfflineMsgSent(getBaseDN());
    final CSN offlineCSN = pendingChanges.putReplicaOfflineMsg();
    dsrsShutdownSync.replicaOfflineMsgSent(getBaseDN(), offlineCSN);
  }
  /**
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/PendingChanges.java
@@ -13,6 +13,7 @@
 *
 * Copyright 2009 Sun Microsystems, Inc.
 * Portions Copyright 2011-2015 ForgeRock AS.
 * Portions Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.plugin;
@@ -123,8 +124,10 @@
  /**
   * Add a replica offline message to the pending list.
   *
   * @return the CSN of the message which was added
   */
  public synchronized void putReplicaOfflineMsg()
  public synchronized CSN putReplicaOfflineMsg()
  {
    final CSN offlineCSN = csnGenerator.newCSN();
    final PendingChange pendingChange =
@@ -133,6 +136,7 @@
    pendingChanges.put(offlineCSN, pendingChange);
    pushCommittedChanges();
    return offlineCSN;
  }
  /**
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
@@ -178,6 +178,11 @@
  /**
   * Creates a new Replication server using the provided configuration entry.
   * <p>
   * The synchronization object this creates is its own, so the resulting server does not
   * synchronize its shutdown with a collocated directory server. A server which has to must be
   * built with {@link #ReplicationServer(ReplicationServerCfg, DSRSShutdownSync)}, passing the
   * instance the directory server side records its ReplicaOfflineMsgs on.
   *
   * @param cfg The configuration of this replication server.
   * @throws ConfigException When Configuration is invalid.
@@ -802,6 +807,9 @@
    // 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.
    // Nothing in an instance which never finished coming up can forward a pending
    // ReplicaOfflineMsg, so this path does not wait for one: it would only delay the failure
    // which is being reported by a grace period which cannot pay off.
    for (ReplicationServerDomain domain : getReplicationServerDomains())
    {
      try
@@ -1181,7 +1189,19 @@
      listenThread.interrupt();
    }
    // shutdown all the replication domains
    /*
     * Let the ReplicaOfflineMsgs a collocated DS sent be forwarded while every handler is still
     * up, and only then stop the domains: shutting a domain down deactivates the consumer of its
     * handlers, clears their message queue and closes their session - see OPENDJ-1453. All the
     * domains wait together and share one deadline, so the shutdown is bounded by one grace
     * period and the wait of one domain does not spend the grace period of the next.
     * <p>
     * This also runs before the assured timer of any domain is cancelled, so an assured update
     * still waiting for acks keeps timing out during the wait instead of holding its sender
     * until the sessions are closed.
     */
    awaitReplicaOfflineMsgsForwarded();
    for (ReplicationServerDomain domain : getReplicationServerDomains())
    {
      domain.shutdown();
@@ -1203,6 +1223,28 @@
  }
  /**
   * Waits for the ReplicaOfflineMsg of every domain which has a replication server to forward it
   * to. With no such server connected there is nobody to forward the message to, and waiting
   * would only delay the shutdown by the whole grace period.
   */
  private void awaitReplicaOfflineMsgsForwarded()
  {
    final List<DN> domainsToWaitFor = new ArrayList<>();
    for (ReplicationServerDomain domain : getReplicationServerDomains())
    {
      if (!domain.getConnectedRSs().isEmpty())
      {
        domainsToWaitFor.add(domain.getBaseDN());
      }
    }
    if (!domainsToWaitFor.isEmpty())
    {
      dsrsShutdownSync.awaitReplicaOfflineMsgsForwarded(
          domainsToWaitFor, dsrsShutdownSync.newShutdownDeadline());
    }
  }
  /**
   * Retrieves the time after which changes must be deleted from the
   * persistent storage (in milliseconds).
   *
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerDomain.java
@@ -1661,7 +1661,15 @@
    return monitorMsg;
  }
  /** Shutdown this ReplicationServerDomain. */
  /**
   * Shutdown this ReplicationServerDomain.
   * <p>
   * A ReplicaOfflineMsg which a collocated DS sent and which is still to be forwarded must be
   * waited for before this runs: stopping the server handlers deactivates their consumer, clears
   * their message queue and closes their session, after which the message can no longer be sent
   * - see OPENDJ-1453. ReplicationServer.shutdown() waits for the messages of all of its domains
   * before it stops any of them.
   */
  public void shutdown()
  {
    DirectoryServer.deregisterMonitorProvider(this);
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java
@@ -13,6 +13,7 @@
 *
 * Copyright 2006-2009 Sun Microsystems, Inc.
 * Portions Copyright 2011-2015 ForgeRock AS.
 * Portions Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.server;
@@ -87,9 +88,14 @@
    LocalizableMessage errMessage = null;
    try
    {
      boolean shutdown = false;
      while (!shutdown
          || !dsrsShutdownSync.canShutdown(replicationServerDomain.getBaseDN()))
      /*
       * Looping here to wait for a pending ReplicaOfflineMsg would achieve nothing: this writer
       * only stops once its handler has been shut down, which deactivates the consumer, clears
       * the message queue and closes the session. The shutdown of the domain waits for the
       * message to be forwarded before it stops the handlers - see
       * ReplicationServerDomain.shutdown() and OPENDJ-1453.
       */
      while (true)
      {
        final UpdateMsg updateMsg = this.handler.take();
        if (updateMsg == null)
@@ -97,15 +103,25 @@
          // this connection is closing
          errMessage = LocalizableMessage.raw(
           "Connection closure: null update returned by domain.");
          shutdown = true;
          break;
        }
        else if (!isUpdateMsgFiltered(updateMsg))
        if (!isUpdateMsgFiltered(updateMsg))
        {
          // Publish the update to the remote server using a protocol version it supports
          session.publish(updateMsg);
          if (updateMsg instanceof ReplicaOfflineMsg)
          /*
           * Only the forward to a peer RS ends the wait of the shutdown: what the grace period
           * buys is the rest of the topology learning that the replica went offline.
           * ReplicationServerDomain.put() never queues this message for a directory server - its
           * isUpdateMsgFiltered() drops it there - but a directory server which is catching up
           * reads its updates from the changelog, where ReplicaCursor synthesizes a
           * ReplicaOfflineMsg from the offline CSN of the replica. Publishing that one says
           * nothing about the peer RSs the shutdown is waiting for.
           */
          if (updateMsg instanceof ReplicaOfflineMsg && !handler.isDataServer())
          {
            dsrsShutdownSync.replicaOfflineMsgForwarded(replicationServerDomain.getBaseDN());
            dsrsShutdownSync.replicaOfflineMsgForwarded(
                replicationServerDomain.getBaseDN(), updateMsg.getCSN());
          }
        }
      }
opendj-server-legacy/src/main/java/org/opends/server/replication/service/DSRSShutdownSync.java
@@ -12,13 +12,19 @@
 * information: "Portions Copyright [year] [name of copyright owner]".
 *
 * Copyright 2014-2016 ForgeRock AS.
 * Portions Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.service;
import java.util.concurrent.ConcurrentSkipListSet;
import java.util.concurrent.atomic.AtomicLong;
import static java.util.concurrent.TimeUnit.MILLISECONDS;
import static java.util.concurrent.TimeUnit.NANOSECONDS;
import java.util.Collection;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import org.forgerock.opendj.ldap.DN;
import org.opends.server.replication.common.CSN;
/**
 * Class useful for the case where DS/RS instances are collocated inside the
@@ -27,24 +33,73 @@
 * More specifically, it ensures a ReplicaOfflineMsg sent by the DS is
 * relayed/forwarded by the collocated RS to the other RSs in the topology
 * before the whole process shuts down.
 * <p>
 * The state is kept per domain and per instance: the collocated DS and RS
 * sides coordinate through the single instance MultimasterReplication hands
 * to both of them.
 *
 * @since OPENDJ-1453
 */
public class DSRSShutdownSync
{
  private static final ConcurrentSkipListSet<DN> replicaOfflineMsgs = new ConcurrentSkipListSet<>();
  private static AtomicLong stopInstanceTimestamp = new AtomicLong();
  /**
   * How long a ReplicaOfflineMsg may hold back the shutdown of the collocated
   * RS, in milliseconds, counted from the moment the message was sent.
   */
  public static final long REPLICA_OFFLINE_GRACE_PERIOD = 5000;
  private final long gracePeriod;
  /**
   * The ReplicaOfflineMsg which has not been forwarded yet, per domain and per
   * replica of that domain.
   * <p>
   * It is kept per domain because a domain sends this message whenever its
   * replication service is disabled - an online import, a restore, a
   * configuration change - and not only when the process shuts down. A single
   * entry for the whole process would be the one of the first such message and
   * would leave no grace period at all to the shutdown this class exists for.
   * <p>
   * It is kept per replica because the collocated RS relays the message of
   * every replica connected to it, and the forward of another replica's
   * message says nothing about this one.
   */
  private final ConcurrentMap<DN, ConcurrentMap<Integer, PendingOfflineMsg>> replicaOfflineMsgs =
      new ConcurrentHashMap<>();
  /** Monitor notified whenever a ReplicaOfflineMsg has been forwarded. */
  private final Object forwardedMonitor = new Object();
  /** Creates a synchronization object using the default grace period. */
  public DSRSShutdownSync()
  {
    this(REPLICA_OFFLINE_GRACE_PERIOD);
  }
  /**
   * Creates a synchronization object using the provided grace period.
   *
   * @param gracePeriod
   *          how long a ReplicaOfflineMsg may hold back the shutdown, in milliseconds
   */
  DSRSShutdownSync(long gracePeriod)
  {
    this.gracePeriod = gracePeriod;
  }
  /**
   * Message has been sent.
   *
   * @param baseDN
   *          the domain for which the message has been sent
   * @param offlineCSN
   *          the CSN of the message, which identifies both the replica which announced itself
   *          offline and the announcement being waited for
   */
  public void replicaOfflineMsgSent(DN baseDN)
  public void replicaOfflineMsgSent(DN baseDN, CSN offlineCSN)
  {
    stopInstanceTimestamp.compareAndSet(0, System.currentTimeMillis());
    replicaOfflineMsgs.add(baseDN);
    replicaOfflineMsgs
        .computeIfAbsent(baseDN, dn -> new ConcurrentHashMap<Integer, PendingOfflineMsg>())
        .put(offlineCSN.getServerId(), new PendingOfflineMsg(offlineCSN, System.nanoTime()));
  }
  /**
@@ -52,23 +107,179 @@
   *
   * @param baseDN
   *          the domain for which the message has been sent
   * @param forwardedCSN
   *          the CSN of the forwarded message
   */
  public void replicaOfflineMsgForwarded(DN baseDN)
  public void replicaOfflineMsgForwarded(DN baseDN, CSN forwardedCSN)
  {
    replicaOfflineMsgs.remove(baseDN);
    final ConcurrentMap<Integer, PendingOfflineMsg> msgs = replicaOfflineMsgs.get(baseDN);
    if (msgs != null)
    {
      final int serverId = forwardedCSN.getServerId();
      final PendingOfflineMsg pending = msgs.get(serverId);
      /*
       * A replica announces itself offline on every disableService(), so the message which is
       * forwarded now may be an older one - queued behind a backlog since an earlier import, or
       * synthesized from the offline CSN of the changelog for a server which is catching up.
       * Such a forward says nothing about the announcement the shutdown is waiting for, and must
       * not consume its grace period.
       */
      if (pending != null && pending.csn.isOlderThanOrEqualTo(forwardedCSN))
      {
        msgs.remove(serverId, pending);
      }
    }
    synchronized (forwardedMonitor)
    {
      forwardedMonitor.notifyAll();
    }
  }
  /**
   * Whether a ReplicationServer ServerReader or ServerWriter can proceed with
   * shutdown.
   * Whether the shutdown of a domain can proceed, i.e. its ReplicaOfflineMsg
   * has been forwarded or its grace period has expired.
   * <p>
   * The shutdown itself blocks on {@link #awaitReplicaOfflineMsgsForwarded(Collection, long)}
   * rather than polling this; it is the same state, observable without waiting for it.
   *
   * @param baseDN
   *          the baseDN of the ServerReader or ServerWriter .
   * @return true if the caller can shutdown, false otherwise
   *          the baseDN of the domain being shut down
   * @return true if the shutdown of this domain need not wait any longer, i.e. its message was
   *         forwarded or its grace period has expired, false otherwise
   */
  public boolean canShutdown(DN baseDN)
  {
    return !replicaOfflineMsgs.contains(baseDN)
        || System.currentTimeMillis() - stopInstanceTimestamp.get() > 5000;
    return remainingGracePeriod(baseDN) <= 0;
  }
  /**
   * Returns the time by which every wait of one shutdown must be over.
   * <p>
   * A process shuts its domains down one after the other and each of them may have a message
   * pending, so a deadline computed once and shared by all of them keeps the whole shutdown
   * bounded by one grace period instead of one per domain.
   *
   * @return the point in time, on the {@link System#nanoTime()} clock, by which the waits must
   *         be over
   */
  public long newShutdownDeadline()
  {
    return System.nanoTime() + MILLISECONDS.toNanos(gracePeriod);
  }
  /**
   * Waits for the ReplicaOfflineMsg of every provided domain to be forwarded, or for their grace
   * periods or the provided deadline to expire.
   * <p>
   * This must be called before the server handlers of those domains are stopped: stopping them
   * deactivates their consumer, clears their message queue and closes their session, after which
   * the message can no longer be forwarded.
   * <p>
   * All the domains of one shutdown wait together rather than one after the other, so that the
   * shutdown is bounded by one grace period without the wait of one domain spending the grace
   * period of the next.
   *
   * @param baseDNs
   *          the baseDNs of the domains whose messages must be forwarded
   * @param deadline
   *          the point in time, on the {@link System#nanoTime()} clock, by which this wait must
   *          be over whatever the domains announce in the meantime - see
   *          {@link #newShutdownDeadline()}. A deadline which is not in the future returns
   *          without waiting at all, for a caller which has nothing to wait for.
   */
  public void awaitReplicaOfflineMsgsForwarded(Collection<DN> baseDNs, long deadline)
  {
    if (deadline - System.nanoTime() <= 0)
    {
      return;
    }
    synchronized (forwardedMonitor)
    {
      while (true)
      {
        final long timeout = Math.min(remainingGracePeriod(baseDNs),
            NANOSECONDS.toMillis(deadline - System.nanoTime()));
        if (timeout <= 0)
        {
          return;
        }
        try
        {
          forwardedMonitor.wait(timeout);
        }
        catch (InterruptedException e)
        {
          /*
           * Give up waiting. The interrupt is deliberately not restored: what follows this call is
           * the rest of the shutdown - joining the reader and writer thread of every handler, then
           * closing the changelog DB - and an interrupt flag would make all of it give up too.
           */
          return;
        }
      }
    }
  }
  /**
   * Returns the time left, in milliseconds, to forward the ReplicaOfflineMsg of the replica of
   * the provided domains which has the longest to wait, zero or less if none of them has a
   * message pending.
   */
  private long remainingGracePeriod(Collection<DN> baseDNs)
  {
    long remaining = 0;
    for (DN baseDN : baseDNs)
    {
      remaining = Math.max(remaining, remainingGracePeriod(baseDN));
    }
    return remaining;
  }
  /**
   * Returns the time left, in milliseconds, to forward the ReplicaOfflineMsg of the replica of
   * this domain which has the longest to wait, zero or less if no message of this domain is
   * pending.
   */
  private long remainingGracePeriod(DN baseDN)
  {
    final ConcurrentMap<Integer, PendingOfflineMsg> msgs = replicaOfflineMsgs.get(baseDN);
    if (msgs == null)
    {
      return 0;
    }
    final long now = System.nanoTime();
    long remaining = 0;
    for (PendingOfflineMsg pending : msgs.values())
    {
      remaining = Math.max(remaining, gracePeriod - NANOSECONDS.toMillis(now - pending.sentTime));
    }
    return remaining;
  }
  /**
   * A ReplicaOfflineMsg a replica announced and which has not been forwarded yet.
   * <p>
   * This deliberately does not override {@code equals}: the two-argument
   * {@link ConcurrentMap#remove(Object, Object)} of the forward guard must match the very
   * announcement it read, not another one which happens to carry the same values.
   */
  private static final class PendingOfflineMsg
  {
    /** The CSN of the message, so that the forward of an older one is not taken for this one. */
    private final CSN csn;
    /** When the message was announced, on the {@link System#nanoTime()} clock. */
    private final long sentTime;
    private PendingOfflineMsg(CSN csn, long sentTime)
    {
      this.csn = csn;
      this.sentTime = sentTime;
    }
    @Override
    public String toString()
    {
      return "PendingOfflineMsg(" + csn + ")";
    }
  }
}
opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java
New file
@@ -0,0 +1,793 @@
/*
 * The contents of this file are subject to the terms of the Common Development and
 * Distribution License (the License). You may not use this file except in compliance with the
 * License.
 *
 * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the
 * specific language governing permission and limitations under the License.
 *
 * When distributing Covered Software, include this CDDL Header Notice in each file and include
 * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL
 * Header, with the fields enclosed by brackets [] replaced by your own identifying
 * information: "Portions copyright [year] [name of copyright owner]".
 *
 * Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.server;
import static org.assertj.core.api.Assertions.assertThat;
import static org.opends.server.TestCaseUtils.TEST_ROOT_DN_STRING;
import static org.opends.server.util.CollectionUtils.newArrayList;
import java.net.InetSocketAddress;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.TreeSet;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import org.forgerock.opendj.ldap.DN;
import org.opends.server.TestCaseUtils;
import org.opends.server.replication.ReplicationTestCase;
import org.opends.server.replication.common.CSN;
import org.opends.server.replication.common.CSNGenerator;
import org.opends.server.replication.common.RSInfo;
import org.opends.server.replication.common.ServerState;
import org.opends.server.replication.protocol.ReplServerStartMsg;
import org.opends.server.replication.protocol.ReplSessionSecurity;
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.TopologyMsg;
import org.opends.server.replication.service.DSRSShutdownSync;
import org.opends.server.replication.service.ReplicationBroker;
import org.opends.server.util.StaticUtils;
import org.opends.server.util.TestTimer;
import org.testng.annotations.Test;
/**
 * The shutdown of a replication server must let a ReplicaOfflineMsg sent by a collocated
 * directory server be forwarded to the other replication servers of the topology before the
 * server handlers are stopped - stopping them deactivates their consumer, clears their message
 * queue and closes their session, after which the message can no longer be sent.
 * <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()} pins the outcome
 * those waits exist for, on a peer connected through the real handshake.
 */
@SuppressWarnings("javadoc")
public class ReplicationServerShutdownSyncTest extends ReplicationTestCase
{
  private static final int SOCKET_TIMEOUT_MS = 30000;
  /** A session end nobody owns is discarded, so its cleanup waits far less than a live one. */
  private static final int DISCARDED_SESSION_TIMEOUT_MS = 2000;
  private static final int REMOTE_RS_ID = 92;
  private static final int REMOTE_DS_ID = 93;
  /** The collocated replica whose ReplicaOfflineMsg the shutdown waits for. */
  private static final int LOCAL_DS_ID = 94;
  /** Time given to the forwarding thread before it releases the shutdown. */
  private static final long FORWARD_DELAY = 500;
  /** How often the domains of {@link #theGracePeriodIsSharedByAllTheDomainsOfOneShutdown()}
   * announce themselves offline again while the shutdown is waiting for them. */
  private static final long REANNOUNCE_INTERVAL = 200;
  @Test
  public void shutdownWaitsForTheReplicaOfflineMsgToBeForwarded() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    ReplicationServer replicationServer = null;
    try (ServerSocket listen = TestCaseUtils.bindFreePort())
    {
      listen.setSoTimeout(SOCKET_TIMEOUT_MS);
      replicationServer = newReplicationServer(shutdownSync, "shutdownSyncWaitDb", 8221);
      final Session[] sessionPair = connectSessionPair(listen, getReplSessionSecurity());
      try (Session remoteEnd = sessionPair[0];
          Session session = sessionPair[1])
      {
        registerConnectedReplicationServer(replicationServer, baseDN, session);
        final long startTime = System.nanoTime();
        shutdownSync.replicaOfflineMsgSent(baseDN, newOfflineCSN());
        replicationServer.shutdown();
        final long elapsed = elapsedMillis(startTime);
        assertThat(elapsed).isGreaterThanOrEqualTo(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
      }
    }
    finally
    {
      removeQuietly(replicationServer);
    }
  }
  @Test
  public void shutdownResumesAsSoonAsTheReplicaOfflineMsgIsForwarded() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    ReplicationServer replicationServer = null;
    Thread forwarder = null;
    try (ServerSocket listen = TestCaseUtils.bindFreePort())
    {
      listen.setSoTimeout(SOCKET_TIMEOUT_MS);
      replicationServer = newReplicationServer(shutdownSync, "shutdownSyncForwardDb", 8222);
      final Session[] sessionPair = connectSessionPair(listen, getReplSessionSecurity());
      try (Session remoteEnd = sessionPair[0];
          Session session = sessionPair[1])
      {
        registerConnectedReplicationServer(replicationServer, baseDN, session);
        final CSN offlineCSN = newOfflineCSN();
        forwarder = newForwarderThread(shutdownSync, baseDN, offlineCSN);
        shutdownSync.replicaOfflineMsgSent(baseDN, offlineCSN);
        final long startTime = System.nanoTime();
        forwarder.start();
        replicationServer.shutdown();
        final long elapsed = elapsedMillis(startTime);
        assertThat(elapsed).isGreaterThanOrEqualTo(FORWARD_DELAY)
            .isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
      }
    }
    finally
    {
      joinQuietly(forwarder);
      removeQuietly(replicationServer);
    }
  }
  /**
   * The outcome the grace period exists for, end to end: a peer replication server connected
   * through the real handshake has received the ReplicaOfflineMsg of the collocated replica by
   * the time the shutdown returns.
   * <p>
   * The waiting tests above measure durations only, so they stay green if the wait is moved
   * after the handlers are stopped - which reintroduces OPENDJ-1453 and loses the message. This
   * one fails in that case.
   */
  @Test
  public void thePeerReceivesTheReplicaOfflineMsgBeforeTheShutdownReturns() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    ReplicationServer replicationServer = null;
    ReplicationBroker broker = null;
    FakePeerReplicationServer peer = null;
    Thread publisher = null;
    try
    {
      final int replicationPort = TestCaseUtils.findFreePort();
      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);
      final ReplicationServerDomain domain =
          replicationServer.getReplicationServerDomain(baseDN, true);
      waitForConnectedReplicationServer(domain);
      final Future<ReplicaOfflineMsg> received = peer.receiveReplicaOfflineMsg();
      /*
       * The replica announces itself offline once the shutdown of the replication server is
       * already waiting for the message, which is the ordering the grace period exists for.
       */
      final CSN offlineCSN = newOfflineCSN();
      shutdownSync.replicaOfflineMsgSent(baseDN, offlineCSN);
      publisher = newPublisherThread(broker, offlineCSN);
      final long startTime = System.nanoTime();
      publisher.start();
      replicationServer.shutdown();
      final long elapsed = elapsedMillis(startTime);
      final ReplicaOfflineMsg forwarded = received.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS);
      assertThat(forwarded)
          .as("the peer replication server was never told that the replica went offline, its "
              + "read ended with: %s", peer.readerFailure())
          .isNotNull();
      assertThat(forwarded.getCSN().getServerId()).isEqualTo(LOCAL_DS_ID);
      assertThat(elapsed).isGreaterThanOrEqualTo(FORWARD_DELAY)
          .isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
    }
    finally
    {
      joinQuietly(publisher);
      closeQuietly(peer);
      stop(broker);
      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
   * offline CSN of the replica, so the writer serving a directory server can publish it - and
   * the peer replication servers would still know nothing.
   */
  @Test
  public void theForwardToADirectoryServerDoesNotEndTheWait() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    ReplicationServer replicationServer = null;
    ReplicationBroker broker = null;
    try (ServerSocket listen = TestCaseUtils.bindFreePort())
    {
      listen.setSoTimeout(SOCKET_TIMEOUT_MS);
      final int replicationPort = TestCaseUtils.findFreePort();
      replicationServer =
          newReplicationServer(shutdownSync, "shutdownSyncDataServerForwardDb", 8227, replicationPort);
      broker = openReplicationSession(baseDN, REMOTE_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      final Session[] sessionPair = connectSessionPair(listen, getReplSessionSecurity());
      try (Session remoteEnd = sessionPair[0];
          Session session = sessionPair[1])
      {
        // a peer replication server, so that the shutdown does wait for the message: what this
        // test pins is that the directory server receiving it is not what ends that wait
        registerConnectedReplicationServer(replicationServer, baseDN, session);
        final ReplicationServerDomain domain =
            replicationServer.getReplicationServerDomain(baseDN, true);
        final DataServerHandler dsHandler = waitForConnectedDirectoryServer(domain);
        final CSN offlineCSN = newOfflineCSN();
        final long startTime = System.nanoTime();
        shutdownSync.replicaOfflineMsgSent(baseDN, offlineCSN);
        // the very message the shutdown waits for, so only the guard of the writer can save it
        dsHandler.add(new ReplicaOfflineMsg(offlineCSN));
        // the directory server did receive it, so its writer went through the forwarding code
        assertThat(waitForSpecificMsg(broker, ReplicaOfflineMsg.class).getCSN().getServerId())
            .isEqualTo(LOCAL_DS_ID);
        assertThat(elapsedMillis(startTime))
            .as("the fixture must deliver the message well inside the grace period, otherwise "
                + "the wait asserted below cannot be told apart from a slow delivery")
            .isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD / 2);
        replicationServer.shutdown();
        final long elapsed = elapsedMillis(startTime);
        assertThat(elapsed)
            .as("the message published to a directory server ended the wait of the shutdown")
            .isGreaterThanOrEqualTo(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
      }
    }
    finally
    {
      stop(broker);
      removeQuietly(replicationServer);
    }
  }
  /**
   * With no other replication server connected there is nobody to forward the message to, so
   * waiting would only delay the shutdown of a standalone server by the whole grace period.
   */
  @Test
  public void shutdownIsNotDelayedWhenNoOtherReplicationServerCanForwardTheMessage() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    ReplicationServer replicationServer = null;
    try
    {
      replicationServer = newReplicationServer(shutdownSync, "shutdownSyncAloneDb", 8223);
      replicationServer.getReplicationServerDomain(baseDN, true);
      final long startTime = System.nanoTime();
      shutdownSync.replicaOfflineMsgSent(baseDN, newOfflineCSN());
      replicationServer.shutdown();
      final long elapsed = elapsedMillis(startTime);
      assertThat(elapsed).isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
    }
    finally
    {
      removeQuietly(replicationServer);
    }
  }
  /**
   * The writer serving a directory server must not hold back the shutdown either: it used to
   * loop on the pending message until the grace period expired, although its handler had already
   * been shut down - which deactivates its consumer and leaves the loop nothing to take.
   */
  @Test
  public void shutdownIsNotDelayedByTheWriterServingADirectoryServer() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    ReplicationServer replicationServer = null;
    ReplicationBroker broker = null;
    try
    {
      final int replicationPort = TestCaseUtils.findFreePort();
      replicationServer =
          newReplicationServer(shutdownSync, "shutdownSyncDataServerDb", 8225, replicationPort);
      broker = openReplicationSession(baseDN, REMOTE_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
      final long startTime = System.nanoTime();
      shutdownSync.replicaOfflineMsgSent(baseDN, newOfflineCSN());
      replicationServer.shutdown();
      final long elapsed = elapsedMillis(startTime);
      assertThat(elapsed).isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
    }
    finally
    {
      stop(broker);
      removeQuietly(replicationServer);
    }
  }
  @Test
  public void shutdownIsNotDelayedWhenNoReplicaOfflineMsgIsPending() throws Exception
  {
    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    ReplicationServer replicationServer = null;
    try (ServerSocket listen = TestCaseUtils.bindFreePort())
    {
      listen.setSoTimeout(SOCKET_TIMEOUT_MS);
      replicationServer = newReplicationServer(shutdownSync, "shutdownSyncNoMsgDb", 8224);
      final Session[] sessionPair = connectSessionPair(listen, getReplSessionSecurity());
      try (Session remoteEnd = sessionPair[0];
          Session session = sessionPair[1])
      {
        registerConnectedReplicationServer(replicationServer, baseDN, session);
        final long startTime = System.nanoTime();
        replicationServer.shutdown();
        final long elapsed = elapsedMillis(startTime);
        assertThat(elapsed).isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
      }
    }
    finally
    {
      removeQuietly(replicationServer);
    }
  }
  /**
   * The domains of a replication server are shut down one after the other, so the grace period
   * must bound the whole shutdown and not each of its domains: a process with several base DNs
   * would otherwise pay it once per domain.
   * <p>
   * Both domains keep announcing themselves offline while the shutdown is running, so neither
   * wait can be ended by a forward and each of them runs to its bound - one grace period in
   * total if it is shared, one per domain otherwise.
   */
  @Test
  public void theGracePeriodIsSharedByAllTheDomainsOfOneShutdown() throws Exception
  {
    final DN baseDN1 = DN.valueOf(TEST_ROOT_DN_STRING);
    final DN baseDN2 = DN.valueOf("dc=world,dc=company");
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    final AtomicBoolean stopped = new AtomicBoolean();
    ReplicationServer replicationServer = null;
    Thread reAnnouncer = null;
    try (ServerSocket listen1 = TestCaseUtils.bindFreePort();
        ServerSocket listen2 = TestCaseUtils.bindFreePort())
    {
      listen1.setSoTimeout(SOCKET_TIMEOUT_MS);
      listen2.setSoTimeout(SOCKET_TIMEOUT_MS);
      replicationServer = newReplicationServer(shutdownSync, "shutdownSyncSharedDeadlineDb", 8228);
      final Session[] sessionPair1 = connectSessionPair(listen1, getReplSessionSecurity());
      final Session[] sessionPair2 = connectSessionPair(listen2, getReplSessionSecurity());
      try (Session remoteEnd1 = sessionPair1[0];
          Session session1 = sessionPair1[1];
          Session remoteEnd2 = sessionPair2[0];
          Session session2 = sessionPair2[1])
      {
        registerConnectedReplicationServer(replicationServer, baseDN1, session1);
        registerConnectedReplicationServer(replicationServer, baseDN2, session2);
        /*
         * Announce both domains offline here rather than leaving it to the thread below: the
         * wait of the shutdown must be armed whatever that thread has had time to run.
         */
        final CSNGenerator csns = new CSNGenerator(LOCAL_DS_ID, 0);
        shutdownSync.replicaOfflineMsgSent(baseDN1, csns.newCSN());
        shutdownSync.replicaOfflineMsgSent(baseDN2, csns.newCSN());
        reAnnouncer = newReAnnouncerThread(shutdownSync, baseDN1, baseDN2, stopped);
        reAnnouncer.start();
        final long startTime = System.nanoTime();
        replicationServer.shutdown();
        final long elapsed = elapsedMillis(startTime);
        assertThat(elapsed).isGreaterThanOrEqualTo(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
        assertThat(elapsed)
            .as("each domain waited its own grace period instead of sharing one deadline")
            .isLessThan(2 * DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
      }
    }
    finally
    {
      stopped.set(true);
      joinQuietly(reAnnouncer);
      removeQuietly(replicationServer);
    }
  }
  private ReplicationServer newReplicationServer(DSRSShutdownSync shutdownSync, String dbDirName,
      int serverId) throws Exception
  {
    return newReplicationServer(shutdownSync, dbDirName, serverId, TestCaseUtils.findFreePort());
  }
  private ReplicationServer newReplicationServer(DSRSShutdownSync shutdownSync, String dbDirName,
      int serverId, int replicationPort) throws Exception
  {
    return new ReplicationServer(new ReplServerFakeConfiguration(
        replicationPort, dbDirName, 0, serverId, 0, 100, new TreeSet<String>()), shutdownSync);
  }
  /**
   * Registers a peer replication server on the domain as the handshake does, but without the
   * protocol exchange: the handler this leaves behind has no writer, which is enough for the
   * tests which only need a domain with a connected peer.
   */
  private void registerConnectedReplicationServer(
      ReplicationServer replicationServer, DN baseDN, Session session) throws Exception
  {
    final ReplicationServerDomain domain = replicationServer.getReplicationServerDomain(baseDN, true);
    final ReplicationServerHandler rsHandler =
        new ReplicationServerHandler(session, 100, replicationServer, 100);
    rsHandler.serverId = REMOTE_RS_ID;
    rsHandler.serverURL = "127.0.0.1:1636";
    rsHandler.setBaseDNAndDomain(baseDN, false);
    domain.lock();
    try
    {
      domain.register(rsHandler);
    }
    finally
    {
      domain.release();
    }
  }
  private void waitForConnectedReplicationServer(final ReplicationServerDomain domain)
      throws Exception
  {
    newConnectionTimer().repeatUntilSuccess(new TestTimer.CallableVoid()
    {
      @Override
      public void call() throws Exception
      {
        assertThat(domain.getConnectedRSs())
            .as("the peer replication server never connected").containsKey(REMOTE_RS_ID);
      }
    });
  }
  private DataServerHandler waitForConnectedDirectoryServer(final ReplicationServerDomain domain)
      throws Exception
  {
    return newConnectionTimer().repeatUntilSuccess(new Callable<DataServerHandler>()
    {
      @Override
      public DataServerHandler call() throws Exception
      {
        final DataServerHandler dsHandler = domain.getConnectedDSs().get(REMOTE_DS_ID);
        assertThat(dsHandler).as("the directory server never connected").isNotNull();
        return dsHandler;
      }
    });
  }
  private static TestTimer newConnectionTimer()
  {
    return new TestTimer.Builder()
        .maxSleep(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS)
        .sleepTimes(10, TimeUnit.MILLISECONDS)
        .toTimer();
  }
  private Thread newForwarderThread(final DSRSShutdownSync shutdownSync, final DN baseDN,
      final CSN offlineCSN)
  {
    return new Thread(new Runnable()
    {
      @Override
      public void run()
      {
        if (!sleepQuietly(FORWARD_DELAY))
        {
          return;
        }
        shutdownSync.replicaOfflineMsgForwarded(baseDN, offlineCSN);
      }
    });
  }
  private Thread newPublisherThread(final ReplicationBroker broker, final CSN offlineCSN)
  {
    return new Thread(new Runnable()
    {
      @Override
      public void run()
      {
        if (!sleepQuietly(FORWARD_DELAY))
        {
          return;
        }
        broker.publish(new ReplicaOfflineMsg(offlineCSN));
      }
    });
  }
  private Thread newReAnnouncerThread(final DSRSShutdownSync shutdownSync, final DN baseDN1,
      final DN baseDN2, final AtomicBoolean stopped)
  {
    return new Thread(new Runnable()
    {
      @Override
      public void run()
      {
        final CSNGenerator csns = new CSNGenerator(LOCAL_DS_ID, 0);
        while (!stopped.get())
        {
          shutdownSync.replicaOfflineMsgSent(baseDN1, csns.newCSN());
          shutdownSync.replicaOfflineMsgSent(baseDN2, csns.newCSN());
          if (!sleepQuietly(REANNOUNCE_INTERVAL))
          {
            return;
          }
        }
      }
    });
  }
  /** Milliseconds elapsed since a {@link System#nanoTime()} reading, the clock the waits use. */
  private static long elapsedMillis(long startTime)
  {
    return TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startTime);
  }
  /** 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();
  }
  private static boolean sleepQuietly(long millis)
  {
    try
    {
      Thread.sleep(millis);
      return true;
    }
    catch (InterruptedException e)
    {
      Thread.currentThread().interrupt();
      return false;
    }
  }
  /** Teardown must never mask the primary assertion failure. */
  private void removeQuietly(ReplicationServer replicationServer)
  {
    try
    {
      remove(replicationServer);
    }
    catch (Exception ignored)
    {
    }
  }
  private void joinQuietly(Thread thread)
  {
    if (thread != null)
    {
      try
      {
        thread.join(SOCKET_TIMEOUT_MS);
      }
      catch (InterruptedException e)
      {
        Thread.currentThread().interrupt();
      }
    }
  }
  private void closeQuietly(FakePeerReplicationServer peer)
  {
    if (peer != null)
    {
      peer.close();
    }
  }
  /**
   * Establishes a connected session pair over the given listen socket, as a remote server
   * connecting to the RS would. The TLS negotiation performed by the session factories needs both
   * ends handshaking at the same time, so the client end runs on its own thread.
   *
   * @return the two sessions: the remote (client) end first, then the local (server) end to hand
   *         to the handler under test
   */
  private Session[] connectSessionPair(ServerSocket listenSocket, final ReplSessionSecurity security)
      throws Exception
  {
    final Socket clientSocket = new Socket("127.0.0.1", listenSocket.getLocalPort());
    clientSocket.setTcpNoDelay(true);
    final ExecutorService executor = Executors.newSingleThreadExecutor();
    Future<Session> clientEnd = null;
    Socket serverSocket = null;
    Session serverEnd = null;
    boolean connected = false;
    try
    {
      clientEnd = executor.submit(new Callable<Session>()
      {
        @Override
        public Session call() throws Exception
        {
          return security.createClientSession(clientSocket, SOCKET_TIMEOUT_MS);
        }
      });
      serverSocket = listenSocket.accept();
      serverSocket.setTcpNoDelay(true);
      serverEnd = security.createServerSession(serverSocket, SOCKET_TIMEOUT_MS);
      assertThat(serverEnd).as("could not create a session for the handler under test").isNotNull();
      final Session[] sessionPair =
          new Session[] { clientEnd.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS), serverEnd };
      connected = true;
      return sessionPair;
    }
    finally
    {
      if (!connected)
      {
        // Nobody owns either end yet: close whatever they managed to create.
        closeClientEndQuietly(clientEnd, clientSocket);
        closeServerEndQuietly(serverEnd, serverSocket);
      }
      executor.shutdown();
    }
  }
  private void closeServerEndQuietly(Session serverEnd, Socket serverSocket)
  {
    if (serverEnd != null)
    {
      serverEnd.close();
    }
    else
    {
      StaticUtils.close(serverSocket);
    }
  }
  private void closeClientEndQuietly(Future<Session> clientEnd, Socket clientSocket)
  {
    if (clientEnd != null)
    {
      try
      {
        final Session session = clientEnd.get(DISCARDED_SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
        if (session != null)
        {
          session.close();
        }
      }
      catch (Exception ignored)
      {
        clientEnd.cancel(true);
      }
    }
    StaticUtils.close(clientSocket);
  }
  /**
   * 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.
   */
  private static final class FakePeerReplicationServer
  {
    private final Session session;
    private final ExecutorService reader = Executors.newSingleThreadExecutor();
    /** Why the peer stopped reading, so that a missing message can be told from a failed one. */
    private volatile Exception readerFailure;
    FakePeerReplicationServer(int replicationPort, int serverId, DN baseDN, long generationId)
        throws Exception
    {
      final Socket socket = new Socket();
      Session newSession = null;
      boolean handshaken = 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, 100,
            new ServerState(), generationId, false, groupId, 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;
      }
      finally
      {
        if (!handshaken)
        {
          // The caller has no handle on this peer yet, so nothing else would close it.
          reader.shutdownNow();
          if (newSession != null)
          {
            newSession.close();
          }
          else
          {
            StaticUtils.close(socket);
          }
        }
      }
      session = newSession;
    }
    /** Returns the first ReplicaOfflineMsg this peer receives, or null if its session ends first. */
    Future<ReplicaOfflineMsg> receiveReplicaOfflineMsg()
    {
      return reader.submit(new Callable<ReplicaOfflineMsg>()
      {
        @Override
        public ReplicaOfflineMsg call()
        {
          try
          {
            while (true)
            {
              final ReplicationMsg msg = session.receive();
              if (msg instanceof ReplicaOfflineMsg)
              {
                return (ReplicaOfflineMsg) msg;
              }
            }
          }
          catch (Exception e)
          {
            // The session is closed when the replication server completes its shutdown: whatever
            // has not arrived by then never will.
            readerFailure = e;
            return null;
          }
        }
      });
    }
    /** Returns what ended the read of this peer, null if nothing did. */
    Exception readerFailure()
    {
      return readerFailure;
    }
    void close()
    {
      reader.shutdownNow();
      session.close();
    }
  }
}
opendj-server-legacy/src/test/java/org/opends/server/replication/service/DSRSShutdownSyncTest.java
New file
@@ -0,0 +1,293 @@
/*
 * The contents of this file are subject to the terms of the Common Development and
 * Distribution License (the License). You may not use this file except in compliance with the
 * License.
 *
 * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the
 * specific language governing permission and limitations under the License.
 *
 * When distributing Covered Software, include this CDDL Header Notice in each file and include
 * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL
 * Header, with the fields enclosed by brackets [] replaced by your own identifying
 * information: "Portions copyright [year] [name of copyright owner]".
 *
 * Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.service;
import static java.util.Arrays.asList;
import static org.assertj.core.api.Assertions.assertThat;
import java.util.concurrent.TimeUnit;
import org.forgerock.opendj.ldap.DN;
import org.opends.server.DirectoryServerTestCase;
import org.opends.server.replication.common.CSN;
import org.testng.annotations.BeforeClass;
import org.testng.annotations.Test;
/** Test the {@link DSRSShutdownSync} class. */
@SuppressWarnings("javadoc")
public class DSRSShutdownSyncTest extends DirectoryServerTestCase
{
  /** Short grace period, for the contracts a test has to wait out. */
  private static final long GRACE_PERIOD = 500;
  /**
   * Grace period for the contracts which only read the state: long enough that no scheduling
   * pause between announcing a message and reading the state can expire it.
   */
  private static final long LONG_GRACE_PERIOD = 60000;
  /** Time given to the forwarding thread before it forwards the message of one domain. */
  private static final long FORWARD_DELAY = 200;
  private static final int SERVER_ID = 1;
  private static final int OTHER_SERVER_ID = 2;
  private static DN baseDN1;
  private static DN baseDN2;
  @BeforeClass
  public static void classSetup() throws Exception
  {
    baseDN1 = DN.valueOf("dc=example,dc=com");
    baseDN2 = DN.valueOf("dc=world,dc=company");
  }
  @Test
  public void canShutdownWhenNoReplicaOfflineMsgWasSent() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(GRACE_PERIOD);
    assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
  }
  @Test
  public void cannotShutdownUntilTheReplicaOfflineMsgIsForwarded() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(SERVER_ID));
    assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
  }
  @Test
  public void canShutdownOnceTheReplicaOfflineMsgIsForwarded() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    final CSN offlineCSN = newCSN(SERVER_ID);
    shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, offlineCSN);
    assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
  }
  @Test
  public void canShutdownOnceTheGracePeriodExpired() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(GRACE_PERIOD);
    shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(SERVER_ID));
    Thread.sleep(GRACE_PERIOD + 50);
    assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
  }
  /**
   * A message sent earlier in the life of the process - an online import, a restore, a
   * configuration change - must not consume the grace period of the message sent by the
   * shutdown this class exists for.
   */
  @Test
  public void gracePeriodOfAShutdownIsNotSpentByAnEarlierMessage() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(GRACE_PERIOD);
    // an import disables then re-enables the replication service
    final CSN sentByTheImport = newCSN(SERVER_ID, 1);
    shutdownSync.replicaOfflineMsgSent(baseDN1, sentByTheImport);
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheImport);
    Thread.sleep(GRACE_PERIOD + 50);
    // the shutdown of the process, much later
    shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(SERVER_ID, 2));
    assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
  }
  /**
   * The message of an earlier announcement may still be queued behind a backlog when the
   * shutdown announces the replica offline again. Forwarding that older message says nothing
   * about the one the shutdown is waiting for, so it must not end the wait.
   */
  @Test
  public void aStaleForwardDoesNotConsumeTheGracePeriodOfANewerMessage() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    final CSN queuedByAnEarlierImport = newCSN(SERVER_ID, 1);
    final CSN sentByTheShutdown = newCSN(SERVER_ID, 2);
    shutdownSync.replicaOfflineMsgSent(baseDN1, queuedByAnEarlierImport);
    shutdownSync.replicaOfflineMsgSent(baseDN1, sentByTheShutdown);
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, queuedByAnEarlierImport);
    assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheShutdown);
    assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
  }
  @Test
  public void gracePeriodIsCountedPerDomain() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(GRACE_PERIOD);
    shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(SERVER_ID));
    Thread.sleep(GRACE_PERIOD + 50);
    shutdownSync.replicaOfflineMsgSent(baseDN2, newCSN(SERVER_ID));
    assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
    assertThat(shutdownSync.canShutdown(baseDN2)).isFalse();
  }
  /**
   * A replication server relays the ReplicaOfflineMsg of every replica connected to it, so the
   * forward of another replica's message must not release the shutdown of this one.
   */
  @Test
  public void theForwardOfAnotherReplicasMessageDoesNotEndTheWait() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(SERVER_ID));
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, newCSN(OTHER_SERVER_ID));
    assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
  }
  /** The domain waits for the message of every one of its replicas, not for the first of them. */
  @Test
  public void aReplicaWhichIsStillWaitingHoldsBackTheShutdownOfItsDomain() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    final CSN ofOneReplica = newCSN(SERVER_ID);
    shutdownSync.replicaOfflineMsgSent(baseDN1, ofOneReplica);
    shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(OTHER_SERVER_ID));
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, ofOneReplica);
    assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
  }
  @Test
  public void canShutdownOnceEveryReplicaOfTheDomainIsForwarded() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    final CSN ofOneReplica = newCSN(SERVER_ID);
    final CSN ofTheOtherReplica = newCSN(OTHER_SERVER_ID);
    shutdownSync.replicaOfflineMsgSent(baseDN1, ofOneReplica);
    shutdownSync.replicaOfflineMsgSent(baseDN1, ofTheOtherReplica);
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, ofOneReplica);
    shutdownSync.replicaOfflineMsgForwarded(baseDN1, ofTheOtherReplica);
    assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
  }
  /**
   * The domains of a shutdown wait together: the wait ends when the message of every one of them
   * has been forwarded, not when the first one has. Waiting for them one after the other would
   * leave the domains which come later without a grace period at all, since the wait of the
   * first one spends the deadline they share.
   */
  @Test
  public void oneWaitCoversEveryDomainOfTheShutdown() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    final CSN ofTheFirstDomain = newCSN(SERVER_ID, 1);
    final CSN ofTheSecondDomain = newCSN(SERVER_ID, 2);
    shutdownSync.replicaOfflineMsgSent(baseDN1, ofTheFirstDomain);
    shutdownSync.replicaOfflineMsgSent(baseDN2, ofTheSecondDomain);
    final Thread forwarder =
        newForwarderThread(shutdownSync, ofTheFirstDomain, ofTheSecondDomain);
    final long startTime = System.nanoTime();
    forwarder.start();
    shutdownSync.awaitReplicaOfflineMsgsForwarded(
        asList(baseDN1, baseDN2), shutdownSync.newShutdownDeadline());
    final long elapsed = millisSince(startTime);
    forwarder.join();
    assertThat(elapsed)
        .as("the wait ended on the first domain forwarded, leaving the second one nothing")
        .isGreaterThanOrEqualTo(2 * FORWARD_DELAY);
    assertThat(elapsed).isLessThan(LONG_GRACE_PERIOD);
  }
  /**
   * However long the messages of a shutdown may still hold it back, the deadline the shutdown
   * was given bounds the wait.
   */
  @Test
  public void theWaitIsBoundedByTheDeadlineOfTheShutdown() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
    shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(SERVER_ID));
    shutdownSync.replicaOfflineMsgSent(baseDN2, newCSN(SERVER_ID));
    final long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(GRACE_PERIOD);
    final long startTime = System.nanoTime();
    shutdownSync.awaitReplicaOfflineMsgsForwarded(asList(baseDN1, baseDN2), deadline);
    final long elapsed = millisSince(startTime);
    assertThat(elapsed).isGreaterThanOrEqualTo(GRACE_PERIOD - 50);
    assertThat(elapsed)
        .as("the wait outlived the deadline of the shutdown")
        .isLessThan(2 * GRACE_PERIOD);
  }
  /** Forwards the message of the first domain, then, as long again later, of the second one. */
  private Thread newForwarderThread(final DSRSShutdownSync shutdownSync,
      final CSN ofTheFirstDomain, final CSN ofTheSecondDomain)
  {
    return new Thread(new Runnable()
    {
      @Override
      public void run()
      {
        try
        {
          Thread.sleep(FORWARD_DELAY);
          shutdownSync.replicaOfflineMsgForwarded(baseDN1, ofTheFirstDomain);
          Thread.sleep(FORWARD_DELAY);
          shutdownSync.replicaOfflineMsgForwarded(baseDN2, ofTheSecondDomain);
        }
        catch (InterruptedException e)
        {
          Thread.currentThread().interrupt();
        }
      }
    });
  }
  private static long millisSince(long startTime)
  {
    return TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startTime);
  }
  /**
   * The CSN of a message a replica announced. They are built by hand rather than with a
   * CSNGenerator: this class has no state to share with the server, and a generator would tie
   * the test to the time service the server starts.
   */
  private static CSN newCSN(int serverId)
  {
    return newCSN(serverId, 1);
  }
  private static CSN newCSN(int serverId, int seqNum)
  {
    return new CSN(1, seqNum, serverId);
  }
}