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); } }