From 9ff409beae28518a8846ac7423eb34f5faec6ee2 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Mon, 14 Sep 2026 07:17:19 +0000
Subject: [PATCH] [#963] Read the changelog again for a following replica the domain is ahead of (#964)
---
opendj-server-legacy/src/test/java/org/opends/server/replication/server/MissedUpdateRecoveryTest.java | 830 +++++++++++++++++++++++++++++++++++++++++++++++++++
opendj-server-legacy/src/messages/org/opends/messages/replication.properties | 3
opendj-server-legacy/src/main/java/org/opends/server/replication/server/DataServerHandler.java | 14
opendj-server-legacy/src/main/java/org/opends/server/replication/server/MessageHandler.java | 107 ++++++
4 files changed, 954 insertions(+), 0 deletions(-)
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/DataServerHandler.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/DataServerHandler.java
index a68f2d2..c13aea8 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/DataServerHandler.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/DataServerHandler.java
@@ -282,6 +282,20 @@
return status;
}
+ /**
+ * {@inheritDoc}
+ * <p>
+ * A directory server is handed every update the domain receives, from whichever server, unless
+ * it has a bad generation id or is being initialized - see
+ * {@code ReplicationServerDomain.isUpdateMsgFiltered()}.
+ */
+ @Override
+ boolean isFedByTheDomain()
+ {
+ final ServerStatus dsStatus = getStatus();
+ return dsStatus != BAD_GEN_ID_STATUS && dsStatus != FULL_UPDATE_STATUS;
+ }
+
@Override
public boolean isDataServer()
{
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/MessageHandler.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/MessageHandler.java
index 9106352..4bebd3e 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/MessageHandler.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/MessageHandler.java
@@ -81,6 +81,19 @@
/** Specifies whether the consumer is following the producer (is not late). */
private boolean following;
/**
+ * The state of the domain at the previous wait of {@link #getNextMessage()} on an empty queue,
+ * or {@code null} when the queue was not empty then. A change it holds has had a whole wait to
+ * reach the queue, see {@link #stopFollowingWhenChangesAreMissing()}. Guarded by
+ * {@link #msgQueue}.
+ */
+ private ServerState domainStateAtPreviousWait;
+ /**
+ * The state of the domain the changelog was last read again for by
+ * {@link #stopFollowingWhenChangesAreMissing()}, or {@code null} when it never was. Guarded by
+ * {@link #msgQueue}.
+ */
+ private ServerState lastMissingChangesDomainState;
+ /**
* Specifies whether the last update message returned by
* {@link #getNextMessage()} was re-read from the changelog DB (catch-up
* path) rather than taken from the in-memory {@link #msgQueue}. Only ever
@@ -233,6 +246,25 @@
}
/**
+ * The state of the domain at the previous wait of {@link #getNextMessage()} on an empty queue,
+ * or {@code null} when the queue was not empty then - see {@link #domainStateAtPreviousWait}.
+ * <p>
+ * Package private for testing: a test which holds a change on its way to the queue hands it
+ * over once a wait has seen the domain hold it, so that
+ * {@link #stopFollowingWhenChangesAreMissing()} has run with the change on its way - by
+ * construction rather than by wall clock.
+ *
+ * @return the state of the domain at the previous wait on an empty queue, or {@code null}
+ */
+ ServerState getDomainStateAtPreviousWait()
+ {
+ synchronized (msgQueue)
+ {
+ return domainStateAtPreviousWait;
+ }
+ }
+
+ /**
* Indicates whether the last update message returned by
* {@code getNextMessage()} was re-read from the changelog DB (catch-up
* path) rather than taken from the in-memory queue.
@@ -392,11 +424,17 @@
{
return null;
}
+ stopFollowingWhenChangesAreMissing();
}
} catch (InterruptedException e)
{
return null;
}
+ if (msgQueue.isEmpty())
+ {
+ // this handler is going back to the changelog, see stopFollowingWhenChangesAreMissing()
+ continue;
+ }
UpdateMsg msg = msgQueue.removeFirst();
if (updateServerState(msg))
{
@@ -420,6 +458,75 @@
}
/**
+ * Leaves the in-memory queue path when the domain holds changes this handler was never given.
+ * <p>
+ * A handler which is following is served by {@link #add(UpdateMsg)} alone, so an empty queue
+ * means the consumer has everything the domain received. A change which reached the changelog
+ * without reaching that queue breaks this: {@link #fillLateQueue()} is the only reader of the
+ * changelog and is not called again while the handler follows, so the change would never be
+ * sent and the consumer would be considered up to date forever - see issue #963. Going back to
+ * the catch-up path reads the changelog again and delivers it.
+ * <p>
+ * A change is missing when the domain held it at the previous wait and it is still neither in
+ * the queue nor in the state of this handler. The state of the domain is advanced by
+ * {@code ReplicationServerDomain.publishUpdateMsg()} slightly before {@code addUpdate()} queues
+ * the change, so what the domain received since the previous wait may simply be on its way;
+ * comparing with the state of the domain as it was one wait ago leaves such a change alone even
+ * when the domain has been ahead of this handler for longer, because of a gap the changelog was
+ * already read for.
+ * <p>
+ * The changelog is read again once per advance of the state of the domain: a gap the re-read
+ * does not close is reported once, and read for again when the domain receives something new.
+ * <p>
+ * Must be called while holding the {@link #msgQueue} monitor.
+ */
+ private void stopFollowingWhenChangesAreMissing()
+ {
+ if (!msgQueue.isEmpty() || !following || !isFedByTheDomain())
+ {
+ domainStateAtPreviousWait = null;
+ return;
+ }
+ final ServerState missingSince = domainStateAtPreviousWait;
+ final ServerState domainState = replicationServerDomain.getLatestServerState();
+ domainStateAtPreviousWait = domainState;
+ if (missingSince == null || serverState.cover(missingSince))
+ {
+ return;
+ }
+ if (lastMissingChangesDomainState != null && lastMissingChangesDomainState.cover(missingSince))
+ {
+ // the changelog has already been read again for these changes and gave nothing: reading it
+ // once more would give nothing either until the domain receives something new
+ return;
+ }
+ lastMissingChangesDomainState = domainState;
+ following = false;
+ logger.warn(WARN_CHANGELOG_READ_AGAIN_FOR_MISSING_CHANGES, replicationServer.getServerId(),
+ baseDN, getMonitorInstanceName(), serverState, domainState);
+ }
+
+ /**
+ * Whether the domain hands every update it receives to this handler.
+ * <p>
+ * Only for such a handler does the state of the domain being ahead of its own mean that a
+ * change was missed, see {@link #stopFollowingWhenChangesAreMissing()}. The state of any other
+ * handler is legitimately behind the state of the domain, and the changelog must not be read on
+ * its behalf: a peer replication server is handed only the changes of the directory servers
+ * connected to this one, so what a third replication server relayed never reaches its queue
+ * nor its state, and reading the changelog again would send it what it already holds; a
+ * directory server the domain filters out would have the writer drop every change read while
+ * its state moved past it.
+ *
+ * @return {@code true} when {@code ReplicationServerDomain.put()} queues every update it
+ * receives for this handler
+ */
+ boolean isFedByTheDomain()
+ {
+ return false;
+ }
+
+ /**
* Fills the late queue with the most recent changes, accepting only the
* messages from provided replica ids.
*/
diff --git a/opendj-server-legacy/src/messages/org/opends/messages/replication.properties b/opendj-server-legacy/src/messages/org/opends/messages/replication.properties
index d5a7e69..fe6c496 100644
--- a/opendj-server-legacy/src/messages/org/opends/messages/replication.properties
+++ b/opendj-server-legacy/src/messages/org/opends/messages/replication.properties
@@ -676,6 +676,9 @@
for the replay of one of its changes to finish. A change which reaches the backend from now on \
is not recorded in the ServerState being saved, so the replication server sends it again and it \
is replayed a second time
+WARN_CHANGELOG_READ_AGAIN_FOR_MISSING_CHANGES_321=Replication server RS(%d) holds changes for domain \
+ "%s" which never reached the message queue of %s: its state %s is behind the state %s of the \
+ domain although it is being served from that queue. Reading the changelog again to send them
ERR_REPLICATION_DOMAIN_CONFIG_CHANGE_FAILED_326=Could not apply a configuration change to the \
replication domain on "%s": %s
NOTE_REPLICATION_DOMAIN_SESSION_NOT_RESTARTED_327=The configuration change was applied to the \
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/MissedUpdateRecoveryTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/MissedUpdateRecoveryTest.java
new file mode 100644
index 0000000..ba74239
--- /dev/null
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/MissedUpdateRecoveryTest.java
@@ -0,0 +1,830 @@
+/*
+ * 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.messages.ReplicationMessages.WARN_CHANGELOG_READ_AGAIN_FOR_MISSING_CHANGES;
+import static org.opends.server.TestCaseUtils.TEST_ROOT_DN_STRING;
+
+import java.net.SocketTimeoutException;
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.SortedSet;
+import java.util.TreeSet;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.Callable;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+
+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.ServerState;
+import org.opends.server.replication.common.ServerStatus;
+import org.opends.server.replication.protocol.DeleteMsg;
+import org.opends.server.replication.protocol.ReplicationMsg;
+import org.opends.server.replication.protocol.UpdateMsg;
+import org.opends.server.replication.service.ReplicationBroker;
+import org.opends.server.util.TestTimer;
+import org.opends.server.util.TimeThread;
+import org.testng.annotations.Test;
+
+/**
+ * A replica which is served from the in-memory message queue of its handler - it is "following" -
+ * is given every change newer than its state by {@code ReplicationServerDomain.put()}. When that
+ * does not happen, the handler must read the changelog again instead of waiting forever for a
+ * delivery which is not coming: a change missed once would otherwise never be sent again, and the
+ * replication server would consider the replica up to date - see issue #963.
+ * <p>
+ * The handler notices it is behind on the ticks of its 500 ms wait on an empty queue, and only
+ * concludes anything about a change it has seen the domain hold over a whole tick: the state of
+ * the domain is advanced slightly before the change is queued, and a change seen on one tick only
+ * may simply be on its way. The watch windows below are sized in ticks accordingly, and a test
+ * which holds a change on its way hands it to the queue once a tick has seen it in the changelog:
+ * the check has then run with the change on its way, by construction rather than by wall clock.
+ */
+@SuppressWarnings("javadoc")
+public class MissedUpdateRecoveryTest extends ReplicationTestCase
+{
+ private static final int RS_ID = 105;
+ private static final int PEER_RS_ID = 109;
+ private static final int THIRD_RS_ID = 110;
+ private static final int DS_ID = 106;
+ /** The replica whose generation id does not match the one of the domain. */
+ private static final int BAD_GENID_DS_ID = 107;
+ /** The replica the changes of this test come from. */
+ private static final int PUBLISHER_DS_ID = 108;
+ /**
+ * A second source of changes, for the tests which hold a change of one replica missing while
+ * the changes of another are delivered: a state holds one CSN per replica, so a newer change of
+ * the same replica would cover the missing one.
+ */
+ private static final int OTHER_PUBLISHER_DS_ID = 111;
+ private static final int WINDOW_SIZE = 100;
+ private static final int SOCKET_TIMEOUT_MS = 30000;
+ /**
+ * How long a change published straight into the changelog is given to reach the replica. The
+ * handler notices it is behind on the next tick of its 500 ms wait on an empty queue.
+ */
+ private static final long DELIVERY_TIMEOUT_MS = 10000;
+ /** How long a handler which must not read the changelog again is watched for: five ticks. */
+ private static final long NO_DELIVERY_WATCH_MS = 2500;
+ /** How many changes are handed to the queue late by the tests which do so. */
+ private static final int CHANGES_ON_THEIR_WAY = 5;
+
+ /**
+ * The regression this test pins: the change is in the changelog and the state of the replica is
+ * behind it, so the replica must be sent the change, whatever kept it out of the message queue
+ * of its handler. The warning which reports it names the state of the replica and the state of
+ * the domain, and is written once.
+ */
+ @Test
+ public void aFollowingReplicaIsSentAChangeItsQueueNeverReceived() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ ReplicationServer replicationServer = null;
+ ReplicationBroker broker = null;
+ Receiver receiver = null;
+ try
+ {
+ final int replicationPort = TestCaseUtils.findFreePort();
+ replicationServer = newReplicationServer("missedUpdateRecoveryDb", replicationPort);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ final ReplicationServerDomain domain =
+ replicationServer.getReplicationServerDomain(baseDN, true);
+ final DataServerHandler handler = waitForFollowingDirectoryServer(domain, DS_ID);
+ // a change of the replica itself, so that its state is not empty in the warning
+ final CSN ownCSN = new CSN(TimeThread.getTime(), 1, DS_ID);
+ broker.publish(newDeleteMsg(baseDN, "cn=own", ownCSN));
+ waitForServerState(handler, ownCSN);
+ receiver = new Receiver(broker);
+
+ final CSN csn = new CSN(TimeThread.getTime(), 1, PUBLISHER_DS_ID);
+ publishToChangelogOnly(replicationServer, baseDN, csn);
+
+ final UpdateMsg receivedMsg = receiver.next(DELIVERY_TIMEOUT_MS);
+ assertThat(receivedMsg)
+ .as("the replica was never sent the change the changelog holds for it").isNotNull();
+ assertThat(receivedMsg.getCSN()).isEqualTo(csn);
+
+ final List<String> warnings = warningsFor(handler, csn);
+ assertThat(new HashSet<>(warnings))
+ .as("the changelog is read again once for a change which is missing once").hasSize(1);
+ assertThat(warnings.get(0))
+ .as("the warning names the state of the replica and the state of the domain")
+ .contains(ownCSN.toString())
+ .contains(csn.toString());
+ assertThat(handler.isFollowing())
+ .as("the replica is served from the queue again once it has been sent the change")
+ .isTrue();
+ }
+ finally
+ {
+ stop(broker);
+ stop(receiver);
+ removeQuietly(replicationServer);
+ }
+ }
+
+ /**
+ * The domain does not hand its updates to a replica whose generation id does not match, and the
+ * changelog must not be read on behalf of such a replica either: the writer would drop every
+ * change it read while the state of the handler moved past it, and the changes would then be
+ * missing from the delivery which follows the reinitialization of the replica.
+ */
+ @Test
+ public void aReplicaTheDomainDoesNotFeedIsLeftAlone() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ ReplicationServer replicationServer = null;
+ ReplicationBroker broker = null;
+ ReplicationBroker badGenIdBroker = null;
+ try
+ {
+ final int replicationPort = TestCaseUtils.findFreePort();
+ replicationServer = newReplicationServer("missedUpdateBadGenIdDb", replicationPort);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ // a change the domain holds and must not hand to the bad-generation-id replica opened below
+ broker.publish(newDeleteMsg(baseDN, "cn=first", new CSN(TimeThread.getTime(), 1, DS_ID)));
+ badGenIdBroker = openReplicationSession(
+ baseDN, BAD_GENID_DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID + 1);
+ final ReplicationServerDomain domain =
+ replicationServer.getReplicationServerDomain(baseDN, true);
+ waitForFollowingDirectoryServer(domain, DS_ID);
+ final DataServerHandler badGenIdHandler = waitForFollowingDirectoryServer(domain, BAD_GENID_DS_ID);
+ assertThat(badGenIdHandler.getStatus())
+ .as("the second replica was expected to be refused the changes of the domain")
+ .isEqualTo(ServerStatus.BAD_GEN_ID_STATUS);
+
+ final CSN csn = new CSN(TimeThread.getTime(), 1, PUBLISHER_DS_ID);
+ publishToChangelogOnly(replicationServer, baseDN, csn);
+
+ assertLeftAlone(badGenIdHandler, csn);
+ }
+ finally
+ {
+ stop(badGenIdBroker);
+ stop(broker);
+ removeQuietly(replicationServer);
+ }
+ }
+
+ /**
+ * The mirror of {@link #aReplicaTheDomainDoesNotFeedIsLeftAlone()} for the other status the
+ * domain filters out: a replica which is being initialized is not sent the changes of the
+ * domain, and the changelog must not be read on its behalf either.
+ */
+ @Test
+ public void aReplicaBeingInitializedIsLeftAlone() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ ReplicationServer replicationServer = null;
+ ReplicationBroker broker = null;
+ try
+ {
+ final int replicationPort = TestCaseUtils.findFreePort();
+ replicationServer = newReplicationServer("missedUpdateFullUpdateDb", replicationPort);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ final ReplicationServerDomain domain =
+ replicationServer.getReplicationServerDomain(baseDN, true);
+ final DataServerHandler handler = waitForFollowingDirectoryServer(domain, DS_ID);
+ broker.signalStatusChange(ServerStatus.FULL_UPDATE_STATUS);
+ waitForStatus(handler, ServerStatus.FULL_UPDATE_STATUS);
+
+ final CSN csn = new CSN(TimeThread.getTime(), 1, PUBLISHER_DS_ID);
+ publishToChangelogOnly(replicationServer, baseDN, csn);
+
+ assertLeftAlone(handler, csn);
+ }
+ finally
+ {
+ stop(broker);
+ removeQuietly(replicationServer);
+ }
+ }
+
+ /**
+ * The state of the domain is advanced slightly before the change is queued, so a change seen
+ * on one tick of the wait only is on its way and not missing: the changelog is not read again
+ * for it, and it reaches the replica once, from the queue.
+ */
+ @Test
+ public void aChangeOnItsWayToTheQueueIsNotReadAgain() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ ReplicationServer replicationServer = null;
+ ReplicationBroker broker = null;
+ Receiver receiver = null;
+ try
+ {
+ final int replicationPort = TestCaseUtils.findFreePort();
+ replicationServer = newReplicationServer("missedUpdateOnItsWayDb", replicationPort);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ final ReplicationServerDomain domain =
+ replicationServer.getReplicationServerDomain(baseDN, true);
+ final DataServerHandler handler = waitForFollowingDirectoryServer(domain, DS_ID);
+ receiver = new Receiver(broker);
+
+ final List<CSN> csns = new ArrayList<>();
+ for (int i = 1; i <= CHANGES_ON_THEIR_WAY; i++)
+ {
+ final CSN csn = new CSN(TimeThread.getTime(), i, PUBLISHER_DS_ID);
+ csns.add(csn);
+ publishThenQueueLate(replicationServer, handler, baseDN, csn);
+ final UpdateMsg receivedMsg = receiver.next(DELIVERY_TIMEOUT_MS);
+ assertThat(receivedMsg).as("the replica was not sent change " + i).isNotNull();
+ assertThat(receivedMsg.getCSN()).isEqualTo(csn);
+ }
+
+ assertThat(receiver.next(NO_DELIVERY_WATCH_MS))
+ .as("the replica was sent a change a second time").isNull();
+ for (CSN csn : csns)
+ {
+ assertThat(warningsFor(handler, csn))
+ .as("the changelog was read again for a change which was on its way").isEmpty();
+ }
+ assertThat(handler.isFollowing()).isTrue();
+ }
+ finally
+ {
+ stop(broker);
+ stop(receiver);
+ removeQuietly(replicationServer);
+ }
+ }
+
+ /**
+ * The ordinary path: a change the domain hands to the queue is sent from the queue, and the
+ * replica stays served from it - nothing is read again and nothing is reported.
+ */
+ @Test
+ public void aChangeTheDomainQueuesKeepsTheReplicaFollowing() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ ReplicationServer replicationServer = null;
+ ReplicationBroker broker = null;
+ ReplicationBroker publisher = null;
+ Receiver receiver = null;
+ try
+ {
+ final int replicationPort = TestCaseUtils.findFreePort();
+ replicationServer = newReplicationServer("missedUpdateQueuedDb", replicationPort);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ publisher = openReplicationSession(
+ baseDN, PUBLISHER_DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ final ReplicationServerDomain domain =
+ replicationServer.getReplicationServerDomain(baseDN, true);
+ final DataServerHandler handler = waitForFollowingDirectoryServer(domain, DS_ID);
+ receiver = new Receiver(broker);
+
+ final CSN csn = new CSN(TimeThread.getTime(), 1, PUBLISHER_DS_ID);
+ publisher.publish(newDeleteMsg(baseDN, "cn=queued", csn));
+
+ final UpdateMsg receivedMsg = receiver.next(DELIVERY_TIMEOUT_MS);
+ assertThat(receivedMsg).as("the replica was not sent the change").isNotNull();
+ assertThat(receivedMsg.getCSN()).isEqualTo(csn);
+ assertThat(receiver.next(NO_DELIVERY_WATCH_MS))
+ .as("the replica was sent the change a second time").isNull();
+ assertThat(warningsFor(handler, csn))
+ .as("the changelog was read again for a change the queue delivered").isEmpty();
+ assertThat(handler.isFollowing()).isTrue();
+ }
+ finally
+ {
+ stop(broker, publisher);
+ stop(receiver);
+ removeQuietly(replicationServer);
+ }
+ }
+
+ /**
+ * The changelog is read again once per advance of the state of the domain. A gap it was already
+ * read for is not read for again while the domain holds nothing new, whatever the state of the
+ * replica says - and it is read for again as soon as the domain receives something.
+ * <p>
+ * The gap is reopened by hand: a changelog which holds a change and cannot yield it cannot be
+ * built from outside, so what is pinned here is that the handler does not open a cursor twice
+ * for one state of the domain.
+ */
+ @Test
+ public void aGapAlreadyReadForIsNotReadAgainUntilTheDomainAdvances() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ ReplicationServer replicationServer = null;
+ ReplicationBroker broker = null;
+ Receiver receiver = null;
+ try
+ {
+ final int replicationPort = TestCaseUtils.findFreePort();
+ replicationServer = newReplicationServer("missedUpdateThrottleDb", replicationPort);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ final ReplicationServerDomain domain =
+ replicationServer.getReplicationServerDomain(baseDN, true);
+ final DataServerHandler handler = waitForFollowingDirectoryServer(domain, DS_ID);
+ receiver = new Receiver(broker);
+
+ final CSN missed = new CSN(TimeThread.getTime(), 1, PUBLISHER_DS_ID);
+ publishToChangelogOnly(replicationServer, baseDN, missed);
+ assertThat(receiver.next(DELIVERY_TIMEOUT_MS)).isNotNull();
+ waitForFollowing(handler);
+ final int warningsBefore = warningsFor(handler, missed).size();
+
+ // the same gap again, with the state of the domain as it was when the changelog was read
+ reopenGap(handler, missed);
+ assertThat(receiver.next(NO_DELIVERY_WATCH_MS))
+ .as("the changelog was read again for a state it was already read for").isNull();
+ assertThat(warningsFor(handler, missed)).hasSize(warningsBefore);
+ assertThat(handler.isFollowing()).isTrue();
+
+ final CSN next = new CSN(TimeThread.getTime(), 2, PUBLISHER_DS_ID);
+ publishToChangelogOnly(replicationServer, baseDN, next);
+ final List<CSN> received = new ArrayList<>();
+ for (int i = 0; i < 2; i++)
+ {
+ final UpdateMsg receivedMsg = receiver.next(DELIVERY_TIMEOUT_MS);
+ assertThat(receivedMsg)
+ .as("the changelog was not read again once the domain received something new")
+ .isNotNull();
+ received.add(receivedMsg.getCSN());
+ }
+ assertThat(received).containsExactly(missed, next);
+ }
+ finally
+ {
+ stop(broker);
+ stop(receiver);
+ removeQuietly(replicationServer);
+ }
+ }
+
+ /**
+ * Next to a gap the changelog was already read for, a change on its way to the queue is still
+ * on its way: the state of the domain is ahead on every tick because of the gap, and that must
+ * not turn the tick which sees the new change into the one which reports it. The report which
+ * follows the advance of the domain comes once the change has been queued and sent, and names
+ * a state of the replica which holds it.
+ */
+ @Test
+ public void aChangeOnItsWayIsNotReportedMissingNextToAGapAlreadyReadFor() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ ReplicationServer replicationServer = null;
+ ReplicationBroker broker = null;
+ Receiver receiver = null;
+ try
+ {
+ final int replicationPort = TestCaseUtils.findFreePort();
+ replicationServer = newReplicationServer("missedUpdateGapOnItsWayDb", replicationPort);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, replicationPort, SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ final ReplicationServerDomain domain =
+ replicationServer.getReplicationServerDomain(baseDN, true);
+ final DataServerHandler handler = waitForFollowingDirectoryServer(domain, DS_ID);
+ receiver = new Receiver(broker);
+
+ final CSN missed = new CSN(TimeThread.getTime(), 1, PUBLISHER_DS_ID);
+ publishToChangelogOnly(replicationServer, baseDN, missed);
+ assertThat(receiver.next(DELIVERY_TIMEOUT_MS)).isNotNull();
+
+ for (int i = 1; i <= CHANGES_ON_THEIR_WAY; i++)
+ {
+ waitForFollowing(handler);
+ reopenGap(handler, missed);
+ // seen over a tick or two: the gap is throttled, nothing is read again for it
+ Thread.sleep(NO_DELIVERY_WATCH_MS);
+
+ final CSN csn = new CSN(TimeThread.getTime(), i, OTHER_PUBLISHER_DS_ID);
+ publishThenQueueLate(replicationServer, handler, baseDN, csn);
+ final List<CSN> received = new ArrayList<>();
+ for (int j = 0; j < 2; j++)
+ {
+ final UpdateMsg receivedMsg = receiver.next(DELIVERY_TIMEOUT_MS);
+ assertThat(receivedMsg).as("delivery " + j + " of round " + i).isNotNull();
+ received.add(receivedMsg.getCSN());
+ }
+ assertThat(received)
+ .as("round " + i + ": the queued change is sent from the queue, the gap is read again after")
+ .containsExactly(csn, missed);
+ final List<String> warnings = warningsFor(handler, csn);
+ assertThat(warnings).as("round " + i).isNotEmpty();
+ for (String warning : warnings)
+ {
+ assertThat(stateOfTheReplicaIn(warning))
+ .as("round " + i + ": a change on its way was reported missing")
+ .contains(csn.toString());
+ }
+ }
+ }
+ finally
+ {
+ stop(broker);
+ stop(receiver);
+ removeQuietly(replicationServer);
+ }
+ }
+
+ /**
+ * A peer replication server is handed only the changes of the directory servers connected to
+ * this one - what a third replication server relayed never reaches its handler, nor its state,
+ * so the domain is always ahead of that handler. That is not a miss: in a mesh of three, the
+ * handler of one peer must not be sent again what the other peer already sent both.
+ */
+ @Test
+ public void aPeerReplicationServerIsLeftAlone() throws Exception
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ final List<ReplicationServer> replicationServers = new ArrayList<>();
+ ReplicationBroker broker = null;
+ try
+ {
+ final int[] ports = { TestCaseUtils.findFreePort(), TestCaseUtils.findFreePort(), TestCaseUtils.findFreePort() };
+ // the replica writes through the third replication server: the two others relay nothing
+ // of what it writes to each other. It connects first, so that the domain has a generation
+ // id by the time the two others are told about it - a replication server presenting none
+ // is not relayed anything
+ final ReplicationServer thirdReplicationServer =
+ newReplicationServer("missedUpdateMeshDb3", ports[2], THIRD_RS_ID, ports[0], ports[1]);
+ replicationServers.add(thirdReplicationServer);
+ broker = openReplicationSession(
+ baseDN, DS_ID, WINDOW_SIZE, ports[2], SOCKET_TIMEOUT_MS, EMPTY_DN_GENID);
+ waitForGenerationId(thirdReplicationServer, baseDN);
+ replicationServers.add(newReplicationServer("missedUpdateMeshDb1", ports[0], RS_ID, ports[1], ports[2]));
+ replicationServers.add(newReplicationServer("missedUpdateMeshDb2", ports[1], PEER_RS_ID, ports[0], ports[2]));
+ final ReplicationServerDomain domain = waitForMesh(replicationServers.get(1), baseDN);
+ final ReplicationServerHandler peerHandler = domain.getConnectedRSs().get(PEER_RS_ID);
+ waitForFollowing(peerHandler);
+
+ final CSN csn = new CSN(TimeThread.getTime(), 1, DS_ID);
+ broker.publish(newDeleteMsg(baseDN, "cn=relayed", csn));
+ waitForDomainState(domain, csn);
+
+ Thread.sleep(NO_DELIVERY_WATCH_MS);
+ assertThat(warningsFor(peerHandler, csn))
+ .as("the changelog was read again for a peer replication server").isEmpty();
+ assertThat(peerHandler.isFollowing()).isTrue();
+ }
+ finally
+ {
+ stop(broker);
+ for (ReplicationServer replicationServer : replicationServers)
+ {
+ removeQuietly(replicationServer);
+ }
+ }
+ }
+
+ /**
+ * Writes a change into the changelog the way {@code ReplicationServerDomain.put()} does, minus
+ * the copy it hands to the message queue of every connected handler. This is the state the
+ * replication server is left in by a change which reached the changelog but not the queue of a
+ * handler.
+ */
+ private void publishToChangelogOnly(ReplicationServer replicationServer, DN baseDN, CSN csn)
+ throws Exception
+ {
+ publishToChangelogOnly(replicationServer, baseDN, newDeleteMsg(baseDN, "cn=missed", csn));
+ }
+
+ private void publishToChangelogOnly(ReplicationServer replicationServer, DN baseDN, UpdateMsg msg)
+ throws Exception
+ {
+ replicationServer.getChangelogDB().getReplicationDomainDB().publishUpdateMsg(baseDN, msg);
+ }
+
+ /**
+ * Writes a change into the changelog and hands it to the queue of the handler once a tick has
+ * seen the domain hold it: the change is on its way for exactly one check, and that check
+ * compared it with a state of the domain seen at a tick before it was published - the previous
+ * tick is waited for first, since a delivery from the queue leaves the handler with no state to
+ * compare with until it has waited on an empty queue again.
+ */
+ private void publishThenQueueLate(
+ ReplicationServer replicationServer, MessageHandler handler, DN baseDN, CSN csn) throws Exception
+ {
+ final UpdateMsg msg = newDeleteMsg(baseDN, "cn=late", csn);
+ waitForATickOnAnEmptyQueue(handler);
+ publishToChangelogOnly(replicationServer, baseDN, msg);
+ waitForATickWhichSaw(handler, csn);
+ handler.add(msg);
+ }
+
+ /**
+ * Takes a change the handler was sent back out of its state, so that the domain is ahead of the
+ * handler by that change again while holding nothing it did not hold when the changelog was
+ * last read for it.
+ */
+ private void reopenGap(MessageHandler handler, CSN csn)
+ {
+ assertThat(handler.getServerState().removeCSN(csn))
+ .as("the handler was expected to hold " + csn).isTrue();
+ }
+
+ private DeleteMsg newDeleteMsg(DN baseDN, String rdn, CSN csn) throws Exception
+ {
+ return new DeleteMsg(DN.valueOf(rdn + "," + baseDN), csn, "uuid");
+ }
+
+ /**
+ * The handler of a replica the domain does not feed keeps its state where it is and reads
+ * nothing: a change it read would be dropped by the writer while its state moved past it.
+ */
+ private void assertLeftAlone(DataServerHandler handler, CSN csn) throws Exception
+ {
+ final long deadline = System.currentTimeMillis() + NO_DELIVERY_WATCH_MS;
+ while (System.currentTimeMillis() < deadline)
+ {
+ assertThat(handler.getServerState().cover(csn))
+ .as("the change was recorded as sent to a replica which is not being sent anything")
+ .isFalse();
+ Thread.sleep(100);
+ }
+ assertThat(warningsFor(handler, csn))
+ .as("the changelog was read again for a replica the domain does not feed").isEmpty();
+ }
+
+ /**
+ * The records of {@code WARN_CHANGELOG_READ_AGAIN_FOR_MISSING_CHANGES} in the error log which
+ * name the given handler and the given change, without their timestamp. The test harness
+ * registers two error log publishers on the same writer, so every record is there twice:
+ * compare sizes with each other, and count distinct records to count warnings.
+ */
+ private static List<String> warningsFor(MessageHandler handler, CSN csn)
+ {
+ final String msgId = "msgID=" + WARN_CHANGELOG_READ_AGAIN_FOR_MISSING_CHANGES.ordinal() + " ";
+ final String handlerName = handler.getMonitorInstanceName();
+ final List<String> warnings = new ArrayList<>();
+ for (String record : TestCaseUtils.ERROR_TEXT_WRITER.getMessages())
+ {
+ if (record.contains(msgId) && record.contains(handlerName) && record.contains(csn.toString()))
+ {
+ warnings.add(record.substring(record.indexOf(" msg=")));
+ }
+ }
+ return warnings;
+ }
+
+ /** The state of the replica as the warning names it: what it says comes before the domain. */
+ private static String stateOfTheReplicaIn(String warning)
+ {
+ final int domainState = warning.indexOf(" is behind the state ");
+ assertThat(domainState).as("the warning names both states: " + warning).isPositive();
+ return warning.substring(0, domainState);
+ }
+
+ private DataServerHandler waitForFollowingDirectoryServer(
+ final ReplicationServerDomain domain, final int serverId) throws Exception
+ {
+ return timer().repeatUntilSuccess(new Callable<DataServerHandler>()
+ {
+ @Override
+ public DataServerHandler call() throws Exception
+ {
+ final DataServerHandler handler = domain.getConnectedDSs().get(serverId);
+ assertThat(handler).as("the directory server never connected").isNotNull();
+ assertThat(handler.isFollowing())
+ .as("the directory server never caught up with the changelog").isTrue();
+ return handler;
+ }
+ });
+ }
+
+ private void waitForFollowing(final MessageHandler handler) throws Exception
+ {
+ timer().repeatUntilSuccess(new Callable<Void>()
+ {
+ @Override
+ public Void call() throws Exception
+ {
+ assertThat(handler.isFollowing()).as("the handler never went back to its queue").isTrue();
+ return null;
+ }
+ });
+ }
+
+ /** Waits for a tick of the handler on an empty queue: the next check has a state to compare with. */
+ private void waitForATickOnAnEmptyQueue(final MessageHandler handler) throws Exception
+ {
+ timer().repeatUntilSuccess(new Callable<Void>()
+ {
+ @Override
+ public Void call() throws Exception
+ {
+ assertThat(handler.getDomainStateAtPreviousWait())
+ .as("the handler never waited on an empty queue").isNotNull();
+ return null;
+ }
+ });
+ }
+
+ /** Waits for a tick of the handler which saw the domain hold the given change. */
+ private void waitForATickWhichSaw(final MessageHandler handler, final CSN csn) throws Exception
+ {
+ timer().repeatUntilSuccess(new Callable<Void>()
+ {
+ @Override
+ public Void call() throws Exception
+ {
+ final ServerState seen = handler.getDomainStateAtPreviousWait();
+ assertThat(seen != null && seen.cover(csn))
+ .as("no wait of the handler saw the domain hold " + csn).isTrue();
+ return null;
+ }
+ });
+ }
+
+ private void waitForServerState(final MessageHandler handler, final CSN csn) throws Exception
+ {
+ timer().repeatUntilSuccess(new Callable<Void>()
+ {
+ @Override
+ public Void call() throws Exception
+ {
+ assertThat(handler.getServerState().cover(csn)).as("the handler never saw " + csn).isTrue();
+ return null;
+ }
+ });
+ }
+
+ private void waitForStatus(final DataServerHandler handler, final ServerStatus status) throws Exception
+ {
+ timer().repeatUntilSuccess(new Callable<Void>()
+ {
+ @Override
+ public Void call() throws Exception
+ {
+ assertThat(handler.getStatus()).isEqualTo(status);
+ return null;
+ }
+ });
+ }
+
+ private void waitForDomainState(final ReplicationServerDomain domain, final CSN csn) throws Exception
+ {
+ timer().repeatUntilSuccess(new Callable<Void>()
+ {
+ @Override
+ public Void call() throws Exception
+ {
+ assertThat(domain.getLatestServerState().cover(csn))
+ .as("the change never reached the changelog of RS(" + domain.getLocalRSServerId() + ")").isTrue();
+ return null;
+ }
+ });
+ }
+
+ private void waitForGenerationId(final ReplicationServer replicationServer, final DN baseDN) throws Exception
+ {
+ timer().repeatUntilSuccess(new Callable<Void>()
+ {
+ @Override
+ public Void call() throws Exception
+ {
+ final ReplicationServerDomain domain = replicationServer.getReplicationServerDomain(baseDN, false);
+ assertThat(domain).as("the replica never connected").isNotNull();
+ assertThat(domain.getGenerationId()).as("the domain never took a generation id").isPositive();
+ return null;
+ }
+ });
+ }
+
+ /** Waits for the domain of the replication server to be connected to the two other ones. */
+ private ReplicationServerDomain waitForMesh(final ReplicationServer replicationServer, final DN baseDN)
+ throws Exception
+ {
+ return timer().repeatUntilSuccess(new Callable<ReplicationServerDomain>()
+ {
+ @Override
+ public ReplicationServerDomain call() throws Exception
+ {
+ final ReplicationServerDomain domain = replicationServer.getReplicationServerDomain(baseDN, false);
+ assertThat(domain).as("the domain never reached this replication server").isNotNull();
+ assertThat(domain.getConnectedRSs().keySet())
+ .as("the replication servers never all connected to each other")
+ .containsOnly(PEER_RS_ID, THIRD_RS_ID);
+ return domain;
+ }
+ });
+ }
+
+ private TestTimer timer()
+ {
+ return new TestTimer.Builder()
+ .maxSleep(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS)
+ .sleepTimes(10, TimeUnit.MILLISECONDS)
+ .toTimer();
+ }
+
+ private ReplicationServer newReplicationServer(String dbDirName, int replicationPort)
+ throws Exception
+ {
+ return newReplicationServer(dbDirName, replicationPort, RS_ID);
+ }
+
+ private ReplicationServer newReplicationServer(
+ String dbDirName, int replicationPort, int serverId, int... peerPorts) throws Exception
+ {
+ final SortedSet<String> peers = new TreeSet<>();
+ for (int peerPort : peerPorts)
+ {
+ peers.add("localhost:" + peerPort);
+ }
+ return new ReplicationServer(new ReplServerFakeConfiguration(
+ replicationPort, dbDirName, 0, serverId, 0, WINDOW_SIZE, peers));
+ }
+
+ private void removeQuietly(ReplicationServer replicationServer)
+ {
+ try
+ {
+ remove(replicationServer);
+ }
+ catch (Exception ignored)
+ {
+ // the test has already reported what matters
+ }
+ }
+
+ private static void stop(Receiver receiver)
+ {
+ if (receiver != null)
+ {
+ receiver.stop();
+ }
+ }
+
+ /**
+ * Collects the update messages a broker is sent, on a thread of its own. The thread leaves when
+ * the broker is stopped - {@code receive()} returns null from then on, without blocking - or
+ * when it is interrupted.
+ */
+ private static final class Receiver implements Runnable
+ {
+ private final ReplicationBroker broker;
+ private final BlockingQueue<UpdateMsg> received = new LinkedBlockingQueue<>();
+ private final Thread thread;
+
+ Receiver(ReplicationBroker broker)
+ {
+ this.broker = broker;
+ this.thread = new Thread(this, "MissedUpdateRecoveryTest receiver for DS(" + broker.getServerId() + ")");
+ this.thread.setDaemon(true);
+ this.thread.start();
+ }
+
+ @Override
+ public void run()
+ {
+ while (!Thread.currentThread().isInterrupted())
+ {
+ try
+ {
+ final ReplicationMsg msg = broker.receive();
+ if (msg == null)
+ {
+ return; // the broker was stopped
+ }
+ if (msg instanceof UpdateMsg)
+ {
+ received.add((UpdateMsg) msg);
+ }
+ }
+ catch (SocketTimeoutException ignored)
+ {
+ // nothing was sent to this replica within the socket timeout, keep reading
+ }
+ }
+ }
+
+ /** The next update message the broker was sent, or null when none came within the timeout. */
+ UpdateMsg next(long timeoutInMillis) throws InterruptedException
+ {
+ return received.poll(timeoutInMillis, TimeUnit.MILLISECONDS);
+ }
+
+ void stop()
+ {
+ thread.interrupt();
+ }
+ }
+}
--
Gitblit v1.10.0