/*
|
* 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.Collections;
|
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;
|
/** A peer replication server the collocated one relays the message to. */
|
private static final int RS_ID = 11;
|
private static final int OTHER_RS_ID = 12;
|
|
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, RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/**
|
* The announcement is made before the message is published, so the broker may still refuse
|
* it - no usable session, or stopped in between. What was announced and never written must
|
* not hold the shutdown back: nobody will forward it.
|
*/
|
@Test
|
public void canShutdownOnceTheReplicaOfflineMsgIsWithdrawn() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgNotSent(baseDN1, offlineCSN);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/**
|
* A withdrawal takes back the very announcement it names, and not whatever the replica
|
* announced last: a stale one, of a message the replica has since announced again, is ignored,
|
* and the newer announcement is still owed its forward. Ignored, and not taken for the
|
* withdrawal of the newer one: that would put the stale announcement back in front, where the
|
* forward of its message - which says nothing about the newer one - would end the wait.
|
*/
|
@Test
|
public void theWithdrawalOfAnEarlierMessageLeavesANewerOneAlone() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN refusedByTheBroker = newCSN(SERVER_ID, 1);
|
final CSN sentByTheShutdown = newCSN(SERVER_ID, 2);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, refusedByTheBroker);
|
shutdownSync.replicaOfflineMsgSent(baseDN1, sentByTheShutdown);
|
shutdownSync.replicaOfflineMsgNotSent(baseDN1, refusedByTheBroker);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
|
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, refusedByTheBroker, RS_ID);
|
assertThat(shutdownSync.canShutdown(baseDN1))
|
.as("the stale withdrawal was ignored, not turned into a restore")
|
.isFalse();
|
}
|
|
/**
|
* A replica announces itself offline on every disableService(), and a later announcement
|
* takes the place of the earlier one. When the broker then refuses the later message, the
|
* earlier one - which did go out, and which a peer still has to forward - must get its wait
|
* back: withdrawing the later announcement must not take the earlier one with it.
|
*/
|
@Test
|
public void theWithdrawalOfALaterMessageGivesTheEarlierOneItsWaitBack() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN sentByTheShutdown = newCSN(SERVER_ID, 1);
|
final CSN refusedByTheBroker = newCSN(SERVER_ID, 2);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, sentByTheShutdown);
|
shutdownSync.replicaOfflineMsgSent(baseDN1, refusedByTheBroker);
|
shutdownSync.replicaOfflineMsgNotSent(baseDN1, refusedByTheBroker);
|
|
assertThat(shutdownSync.canShutdown(baseDN1))
|
.as("the earlier message went out and nobody has forwarded it yet")
|
.isFalse();
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheShutdown, RS_ID);
|
assertThat(shutdownSync.canShutdown(baseDN1))
|
.as("the forward of the earlier message ends the wait")
|
.isTrue();
|
}
|
|
/**
|
* What the withdrawal gives back is the very announcement which was displaced, with the peers
|
* its message was queued for: the forward of one of them does not end a wait which is for
|
* several, as it would for a message which was announced again from scratch.
|
*/
|
@Test
|
public void theRestoredAnnouncementIsStillOwedTheForwardsItWasQueuedFor() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN sentByTheShutdown = newCSN(SERVER_ID, 1);
|
final CSN refusedByTheBroker = newCSN(SERVER_ID, 2);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, sentByTheShutdown);
|
shutdownSync.replicaOfflineMsgDispatched(
|
baseDN1, sentByTheShutdown, asList(RS_ID, OTHER_RS_ID));
|
shutdownSync.replicaOfflineMsgSent(baseDN1, refusedByTheBroker);
|
shutdownSync.replicaOfflineMsgNotSent(baseDN1, refusedByTheBroker);
|
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheShutdown, RS_ID);
|
assertThat(shutdownSync.canShutdown(baseDN1))
|
.as("the restored message is still owed the other peer's forward")
|
.isFalse();
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheShutdown, OTHER_RS_ID);
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/**
|
* The restored announcement keeps its own clock: the shutdown waits out what is left of the
|
* earlier message's grace period, not a new one counted from the withdrawal.
|
*/
|
@Test
|
public void theRestoredAnnouncementKeepsWhatIsLeftOfItsOwnGracePeriod() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(GRACE_PERIOD);
|
final CSN sentByTheShutdown = newCSN(SERVER_ID, 1);
|
final CSN refusedByTheBroker = newCSN(SERVER_ID, 2);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, sentByTheShutdown);
|
Thread.sleep(GRACE_PERIOD - 200);
|
shutdownSync.replicaOfflineMsgSent(baseDN1, refusedByTheBroker);
|
shutdownSync.replicaOfflineMsgNotSent(baseDN1, refusedByTheBroker);
|
// past the end of the earlier message's grace period, well short of a whole new one
|
Thread.sleep(250);
|
|
assertThat(shutdownSync.canShutdown(baseDN1))
|
.as("the wait started over at the withdrawal")
|
.isTrue();
|
}
|
|
/**
|
* While the announcement of a refused message stands in the place of the earlier one, what is
|
* reported about the earlier message is not seen by it: a forward reported in that window is
|
* lost, and the restored announcement waits out what is left of its own grace period. The
|
* window is the one refused publish; this pins the trade-off, so that a change to it is made
|
* knowingly.
|
*/
|
@Test
|
public void aForwardReportedWhileARefusedAnnouncementStoodIsNotSeen() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN sentByTheShutdown = newCSN(SERVER_ID, 1);
|
final CSN refusedByTheBroker = newCSN(SERVER_ID, 2);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, sentByTheShutdown);
|
shutdownSync.replicaOfflineMsgDispatched(baseDN1, sentByTheShutdown, asList(RS_ID));
|
shutdownSync.replicaOfflineMsgSent(baseDN1, refusedByTheBroker);
|
// not seen: the announcement of the refused message stands in front
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheShutdown, RS_ID);
|
shutdownSync.replicaOfflineMsgNotSent(baseDN1, refusedByTheBroker);
|
|
assertThat(shutdownSync.canShutdown(baseDN1))
|
.as("a forward reported while the refused announcement stood is not seen")
|
.isFalse();
|
}
|
|
@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, RS_ID);
|
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, RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
|
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheShutdown, RS_ID);
|
|
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), RS_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, RS_ID);
|
|
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, RS_ID);
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, ofTheOtherReplica, RS_ID);
|
|
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);
|
}
|
|
/**
|
* The collocated replication server queues the message for every peer it relays to, and each
|
* of them is served by its own writer: the forward of one peer says nothing about the others,
|
* whose queue the shutdown is about to clear.
|
*/
|
@Test
|
public void theForwardOfOnePeerDoesNotEndTheWaitOfTheOthers() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgDispatched(baseDN1, offlineCSN, asList(RS_ID, OTHER_RS_ID));
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, offlineCSN, RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
|
}
|
|
@Test
|
public void canShutdownOnceEveryPeerTheMessageWasQueuedForForwardedIt() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgDispatched(baseDN1, offlineCSN, asList(RS_ID, OTHER_RS_ID));
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, offlineCSN, RS_ID);
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, offlineCSN, OTHER_RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/**
|
* With no peer to relay the message to - none connected, or none sharing the generation id of
|
* the domain - there is nothing to wait for.
|
*/
|
@Test
|
public void canShutdownWhenTheMessageWasQueuedForNoPeer() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgDispatched(
|
baseDN1, offlineCSN, Collections.<Integer> emptyList());
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/**
|
* A peer which is no longer connected cannot forward anything, so the shutdown must not spend
|
* the rest of its window waiting for it.
|
*/
|
@Test
|
public void aPeerWhichStoppedIsNoLongerWaitedFor() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgDispatched(baseDN1, offlineCSN, asList(RS_ID, OTHER_RS_ID));
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, offlineCSN, RS_ID);
|
shutdownSync.replicaOfflineMsgNotForwarded(baseDN1, OTHER_RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/** The peers which are still connected keep their part of the grace period. */
|
@Test
|
public void aPeerWhichStoppedDoesNotEndTheWaitOfTheOthers() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgDispatched(baseDN1, offlineCSN, asList(RS_ID, OTHER_RS_ID));
|
shutdownSync.replicaOfflineMsgNotForwarded(baseDN1, OTHER_RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
|
}
|
|
/**
|
* A peer which connected after the message was queued was never given it, so what it forwards
|
* is a message of its own catch-up and says nothing about the peers which still owe theirs.
|
*/
|
@Test
|
public void theForwardOfAPeerTheMessageWasNotQueuedForDoesNotEndTheWait() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgDispatched(baseDN1, offlineCSN, asList(RS_ID));
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, offlineCSN, OTHER_RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
|
}
|
|
/**
|
* The peers are recorded by the collocated replication server when it queues the message for
|
* them, which the message of a replica connected to a remote replication server never reaches.
|
* With no peer recorded the wait keeps the behaviour it had before they were tracked: the
|
* first forward ends it.
|
*/
|
@Test
|
public void theFirstForwardEndsTheWaitWhenNoPeerWasRecorded() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, offlineCSN, RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/**
|
* A peer going away says nothing about a message it was never given, which is the opposite of
|
* what a forward says: with no peer recorded the first forward ends the wait, and a give-up
|
* must leave it running. Otherwise any peer disconnecting would release a message the
|
* collocated replication server has not queued for anybody yet - the very bug the recipients
|
* were introduced to close, in a new shape.
|
*/
|
@Test
|
public void aPeerStoppingBeforeTheMessageIsQueuedDoesNotEndTheWait() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
|
shutdownSync.replicaOfflineMsgSent(baseDN1, newCSN(SERVER_ID));
|
shutdownSync.replicaOfflineMsgNotForwarded(baseDN1, RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isFalse();
|
}
|
|
/**
|
* A replica announces itself offline on every disableService(), so the peers recorded for an
|
* earlier announcement say nothing about the one the shutdown is waiting for.
|
*/
|
@Test
|
public void thePeersOfAnEarlierAnnouncementAreNotTakenForThoseOfThisOne() 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, sentByTheShutdown);
|
shutdownSync.replicaOfflineMsgDispatched(
|
baseDN1, queuedByAnEarlierImport, asList(RS_ID, OTHER_RS_ID));
|
shutdownSync.replicaOfflineMsgForwarded(baseDN1, sentByTheShutdown, RS_ID);
|
|
assertThat(shutdownSync.canShutdown(baseDN1)).isTrue();
|
}
|
|
/**
|
* The peer which goes away must wake the shutdown up, and not leave it waiting for a forward
|
* nobody can report any more.
|
*/
|
@Test
|
public void theWaitEndsWhenTheLastPeerExpectedToForwardStops() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
shutdownSync.replicaOfflineMsgDispatched(baseDN1, offlineCSN, asList(RS_ID));
|
final Thread peerStopper = newPeerStopperThread(shutdownSync, RS_ID);
|
|
final long startTime = System.nanoTime();
|
peerStopper.start();
|
shutdownSync.awaitReplicaOfflineMsgsForwarded(
|
asList(baseDN1), shutdownSync.newShutdownDeadline());
|
final long elapsed = millisSince(startTime);
|
peerStopper.join();
|
|
assertThat(elapsed).isGreaterThanOrEqualTo(FORWARD_DELAY);
|
assertThat(elapsed)
|
.as("the peer going away did not wake the wait up")
|
.isLessThan(LONG_GRACE_PERIOD);
|
}
|
|
/**
|
* A withdrawal must wake the shutdown up as a forward does, and not leave it waiting for the
|
* forward of a message which never left.
|
*/
|
@Test
|
public void theWaitEndsWhenTheMessageIsWithdrawn() throws Exception
|
{
|
final DSRSShutdownSync shutdownSync = new DSRSShutdownSync(LONG_GRACE_PERIOD);
|
final CSN offlineCSN = newCSN(SERVER_ID);
|
shutdownSync.replicaOfflineMsgSent(baseDN1, offlineCSN);
|
final Thread withdrawer = newWithdrawerThread(shutdownSync, offlineCSN);
|
|
final long startTime = System.nanoTime();
|
withdrawer.start();
|
shutdownSync.awaitReplicaOfflineMsgsForwarded(
|
asList(baseDN1), shutdownSync.newShutdownDeadline());
|
final long elapsed = millisSince(startTime);
|
withdrawer.join();
|
|
assertThat(elapsed).isGreaterThanOrEqualTo(FORWARD_DELAY);
|
assertThat(elapsed)
|
.as("the withdrawal did not wake the wait up")
|
.isLessThan(LONG_GRACE_PERIOD);
|
assertThat(shutdownSync.canShutdown(baseDN1))
|
.as("the withdrawn message holds nothing back")
|
.isTrue();
|
}
|
|
/** Withdraws the announcement, as the broker refusing the message during the wait does. */
|
private Thread newWithdrawerThread(final DSRSShutdownSync shutdownSync, final CSN offlineCSN)
|
{
|
return new Thread(new Runnable()
|
{
|
@Override
|
public void run()
|
{
|
try
|
{
|
Thread.sleep(FORWARD_DELAY);
|
shutdownSync.replicaOfflineMsgNotSent(baseDN1, offlineCSN);
|
}
|
catch (InterruptedException e)
|
{
|
Thread.currentThread().interrupt();
|
}
|
}
|
});
|
}
|
|
/** Stops the peer the message was queued for, as a disconnection during the wait does. */
|
private Thread newPeerStopperThread(final DSRSShutdownSync shutdownSync, final int peerId)
|
{
|
return new Thread(new Runnable()
|
{
|
@Override
|
public void run()
|
{
|
try
|
{
|
Thread.sleep(FORWARD_DELAY);
|
shutdownSync.replicaOfflineMsgNotForwarded(baseDN1, peerId);
|
}
|
catch (InterruptedException e)
|
{
|
Thread.currentThread().interrupt();
|
}
|
}
|
});
|
}
|
|
/** 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, RS_ID);
|
Thread.sleep(FORWARD_DELAY);
|
shutdownSync.replicaOfflineMsgForwarded(baseDN2, ofTheSecondDomain, RS_ID);
|
}
|
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);
|
}
|
}
|