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