From 2d99dd7e2843c998a2d15f072cb827f443a79809 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Wed, 23 Sep 2026 07:39:43 +0000
Subject: [PATCH] [#1014] Report a ReplicaOfflineMsg as forwarded only to a peer which can decode it (#1019)

---
 opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicaOfflineMsg.java               |   30 +++
 opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicationMsg.java                  |   20 ++
 opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java                 |   18 +
 opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java                      |   25 ++
 opendj-server-legacy/src/messages/org/opends/messages/replication.properties                                   |    3 
 opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java |  362 +++++++++++++++++++++++++++++++++++++++
 opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SynchronizationMsgTest.java          |   16 +
 opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerDomain.java           |   21 ++
 8 files changed, 483 insertions(+), 12 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicaOfflineMsg.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicaOfflineMsg.java
index 80dc167..d2768e6 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicaOfflineMsg.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicaOfflineMsg.java
@@ -12,6 +12,7 @@
  * information: "Portions Copyright [year] [name of copyright owner]".
  *
  * Copyright 2014 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
  */
 package org.opends.server.replication.protocol;
 
@@ -81,7 +82,7 @@
   @Override
   public byte[] getBytes(short protocolVersion)
   {
-    if (protocolVersion < ProtocolVersion.REPLICATION_PROTOCOL_V8)
+    if (!isSupportedBy(protocolVersion))
     {
       return null;
     }
@@ -94,6 +95,33 @@
 
   /** {@inheritDoc} */
   @Override
+  public boolean isEncodableFor(short protocolVersion)
+  {
+    return isSupportedBy(protocolVersion);
+  }
+
+  /**
+   * Whether a peer which negotiated the provided replication protocol version
+   * can be sent this message at all.
+   * <p>
+   * The message was introduced by
+   * {@link ProtocolVersion#REPLICATION_PROTOCOL_V8} and has no encoding below
+   * it: an older peer is not told that the replica went offline, and cannot be.
+   * This answers for the message type, for a caller which has no instance at
+   * hand - see {@link #isEncodableFor(short)} for the one which has.
+   *
+   * @param protocolVersion
+   *          The protocol version negotiated with the peer.
+   * @return <code>true</code> if that version carries this message,
+   *         <code>false</code> otherwise.
+   */
+  public static boolean isSupportedBy(short protocolVersion)
+  {
+    return protocolVersion >= ProtocolVersion.REPLICATION_PROTOCOL_V8;
+  }
+
+  /** {@inheritDoc} */
+  @Override
   public int size()
   {
     return bytes(1) + shorts(1) + csns(1);
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicationMsg.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicationMsg.java
index bb0f2bd..ce75fe6 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicationMsg.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/ReplicationMsg.java
@@ -13,6 +13,7 @@
  *
  * Copyright 2006-2010 Sun Microsystems, Inc.
  * Portions copyright 2013-2016 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
  */
 package org.opends.server.replication.protocol;
 
@@ -109,6 +110,25 @@
   public abstract byte[] getBytes(short protocolVersion);
 
   /**
+   * Whether this message has an encoding for a peer which negotiated the
+   * provided replication protocol version.
+   * <p>
+   * A message which has none is dropped by {@link Session#publish(ReplicationMsg)}
+   * - its {@link #getBytes(short)} returns <code>null</code> - with nothing on
+   * the wire and nothing in the log, so a caller which goes on to report what
+   * the peer was told must ask this before it publishes.
+   *
+   * @param protocolVersion
+   *          The protocol version negotiated with the peer.
+   * @return <code>true</code> if this message can be encoded for that version,
+   *         <code>false</code> if it has no encoding for it at all.
+   */
+  public boolean isEncodableFor(short protocolVersion)
+  {
+    return true;
+  }
+
+  /**
    * Generates a ReplicationMsg from its encoded form. This un-serialization is
    * done taking into account the various supported replication protocol
    * versions.
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
index 05232fe..961cd88 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
@@ -1721,21 +1721,25 @@
    * to. With no such server connected there is nobody to forward the message to, and waiting
    * would only delay the shutdown by the whole grace period.
    * <p>
+   * A peer which negotiated a protocol version the message has no encoding for is no such
+   * server either: nothing is ever sent to it and no writer will ever report a forward for it.
+   * <p>
    * The recipients DSRSShutdownSync records when the message is queued are the sharper source of
    * truth and cover the domains this test lets through - a peer which connected after the
-   * message was queued owes nothing, and is not waited for. What this test still covers is the
-   * announcement which was never queued here at all, and for which no recipient can therefore be
-   * recorded: the message of a replica which picked a remote replication server, or the
-   * announcement of issue #918 recorded after its message was relayed. Those wait out the whole
-   * grace period on the first forward, and on a server with no peer connected nothing would ever
-   * report one.
+   * message was queued owes nothing, and is not waited for, and the writer of a peer which
+   * cannot decode the message strikes it off when it drops it. What this test still covers is
+   * the announcement which was never queued here at all, and for which no recipient can
+   * therefore be recorded: the message of a replica which picked a remote replication server, or
+   * the announcement of issue #918 recorded after its message was relayed. Those wait out the
+   * whole grace period on the first forward, and on a server whose only peers cannot decode the
+   * message - or with no peer connected at all - nothing would ever report one.
    */
   private void awaitReplicaOfflineMsgsForwarded()
   {
     final List<DN> domainsToWaitFor = new ArrayList<>();
     for (ReplicationServerDomain domain : getReplicationServerDomains())
     {
-      if (!domain.getConnectedRSs().isEmpty())
+      if (domain.hasPeerWhichCanReceiveReplicaOfflineMsgs())
       {
         domainsToWaitFor.add(domain.getBaseDN());
       }
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerDomain.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerDomain.java
index c402c5b..ad054f8 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerDomain.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerDomain.java
@@ -2379,6 +2379,27 @@
   }
 
   /**
+   * Whether any connected peer replication server can be forwarded a
+   * {@link ReplicaOfflineMsg}: the message has no encoding below replication protocol version 8,
+   * so a peer which negotiated an older one is nobody to wait for - the writer serving it drops
+   * the message instead of publishing it.
+   *
+   * @return {@code true} if at least one connected replication server can receive a
+   *         {@link ReplicaOfflineMsg}, {@code false} otherwise
+   */
+  public boolean hasPeerWhichCanReceiveReplicaOfflineMsgs()
+  {
+    for (ReplicationServerHandler rsHandler : connectedRSs.values())
+    {
+      if (ReplicaOfflineMsg.isSupportedBy(rsHandler.getProtocolVersion()))
+      {
+        return true;
+      }
+    }
+    return false;
+  }
+
+  /**
    * Get the map of connected RSs.
    * @return The map of connected RSs
    */
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java
index 3f2ecdd..d1edac1 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java
@@ -172,6 +172,31 @@
 
   private boolean isUpdateMsgFiltered(UpdateMsg updateMsg)
   {
+    if (!updateMsg.isEncodableFor(handler.getProtocolVersion()))
+    {
+      /*
+       * Session.publish() drops what it cannot encode for its peer and returns as if it had sent
+       * it, so this must be caught here: the drop would otherwise be reported as the forward of
+       * a ReplicaOfflineMsg, which the shutdown takes for this peer having been told - and, for
+       * a message no recipient was recorded for, as the end of its wait for the peers which can
+       * still be told - see OPENDJ-1453 and issue #1014. Dropping it here reports what happened
+       * instead. There is nothing to send to this one: the message is not part of the protocol
+       * version it negotiated, and no version of it will ever reach it.
+       * <p>
+       * The drop is reported rather than traced whoever the consumer is: today the only message
+       * gated on a protocol version is the ReplicaOfflineMsg, which never reaches the writer of
+       * a directory server at all - ReplicationServerDomain.put() does not queue it for one, and
+       * DataServerHandler.updateServerState() drops the copy the changelog cursor of a directory
+       * server which is catching up synthesizes from the offline CSN of the replica (issue
+       * #1029). A message which does get here is one the consumer was to be sent and will not
+       * be, which is what this record says.
+       */
+      logger.warn(WARN_IGNORING_UPDATE_UNSUPPORTED_BY_PEER,
+          handler.getReplicationServerId(), updateMsg.getCSN(), handler.getBaseDN(),
+          handler.getServerId(), session.getReadableRemoteAddress(),
+          handler.getProtocolVersion());
+      return true;
+    }
     if (handler.isDataServer())
     {
       /**
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 594524f..a4eb51f 100644
--- a/opendj-server-legacy/src/messages/org/opends/messages/replication.properties
+++ b/opendj-server-legacy/src/messages/org/opends/messages/replication.properties
@@ -709,3 +709,6 @@
  replication server was stopped before the initialization from server %s completed. The \
  entries which had arrived are imported, and the generation id of the data is computed over \
  them rather than taken from the exporter
+WARN_IGNORING_UPDATE_UNSUPPORTED_BY_PEER_328=Replication server RS(%d) not sending update \
+ %s for domain "%s" to server %d at %s because the replication protocol version %d \
+ negotiated with it has no encoding for this message
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SynchronizationMsgTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SynchronizationMsgTest.java
index bba5b3a..2979910 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SynchronizationMsgTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SynchronizationMsgTest.java
@@ -815,6 +815,22 @@
   }
 
   /**
+   * The ReplicaOfflineMsg was introduced by the 8th version of the replication protocol and has
+   * no encoding below it, where getBytes() returns null. Session.publish() drops such a message
+   * without a trace, so whoever goes on to report that a peer was told must ask first.
+   */
+  @Test
+  public void replicaOfflineMsgHasNoEncodingBelowV8() throws Exception
+  {
+    final ReplicaOfflineMsg msg = new ReplicaOfflineMsg(new CSN(System.currentTimeMillis(), 0, 42));
+
+    assertTrue(msg.isEncodableFor(REPLICATION_PROTOCOL_V8));
+    assertNotNull(msg.getBytes(REPLICATION_PROTOCOL_V8));
+    assertFalse(msg.isEncodableFor(REPLICATION_PROTOCOL_V7));
+    assertNull(msg.getBytes(REPLICATION_PROTOCOL_V7));
+  }
+
+  /**
    * Test that WindowMsg encoding and decoding works
    * by checking that : msg == new WindowMsg(msg.getBytes()).
    */
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java
index 31d948d..c8e503c 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java
@@ -16,6 +16,7 @@
 package org.opends.server.replication.server;
 
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.opends.messages.ReplicationMessages.WARN_IGNORING_UPDATE_UNSUPPORTED_BY_PEER;
 import static org.opends.server.TestCaseUtils.TEST_ROOT_DN_STRING;
 import static org.opends.server.util.CollectionUtils.newArrayList;
 
@@ -23,6 +24,7 @@
 import java.net.InetSocketAddress;
 import java.net.ServerSocket;
 import java.net.Socket;
+import java.util.ArrayList;
 import java.util.Collection;
 import java.util.List;
 import java.util.TreeSet;
@@ -46,6 +48,7 @@
 import org.opends.server.replication.common.RSInfo;
 import org.opends.server.replication.common.ServerState;
 import org.opends.server.replication.protocol.DeleteMsg;
+import org.opends.server.replication.protocol.ProtocolVersion;
 import org.opends.server.replication.protocol.ReplServerStartMsg;
 import org.opends.server.replication.protocol.ReplSessionSecurity;
 import org.opends.server.replication.protocol.ReplicaOfflineMsg;
@@ -58,6 +61,7 @@
 import org.opends.server.replication.service.ReplicationBroker;
 import org.opends.server.util.StaticUtils;
 import org.opends.server.util.TestTimer;
+import org.testng.annotations.DataProvider;
 import org.testng.annotations.Test;
 
 /**
@@ -96,6 +100,15 @@
    * published one, so the shutdown spends its whole grace period waiting for it.
    */
   private static final int UNREACHABLE_DS_ID = 98;
+  /** The peer replication server which speaks a protocol version with no ReplicaOfflineMsg. */
+  private static final int OLD_PEER_RS_ID = 99;
+  /**
+   * Base of the id of the replication server serving a peer of a chosen protocol version: the
+   * version is added to it, so that the two rows of
+   * {@link #onlyAPeerWhichCanDecodeTheMessageIsToldTheReplicaWentOffline(short, boolean)} do not
+   * share a server id.
+   */
+  private static final int PEER_VERSION_RS_ID = 8240;
   /** Send window a peer advertises when nothing has to hold its writer back. */
   private static final int PEER_WINDOW = 100;
   /**
@@ -244,6 +257,155 @@
     }
   }
 
+  @DataProvider
+  public Object[][] peerProtocolVersions()
+  {
+    return new Object[][] {
+      { ProtocolVersion.getCurrentVersion(), true },
+      { ProtocolVersion.REPLICATION_PROTOCOL_V7, false },
+    };
+  }
+
+  /**
+   * A peer replication server which negotiated a protocol version older than the one which
+   * introduced the ReplicaOfflineMsg cannot be told that a replica went offline: the message has
+   * no encoding for such a peer and {@link Session#publish(ReplicationMsg)} drops what it cannot
+   * encode. Nothing can be done about the peer itself - the announcement is not part of the
+   * protocol it speaks - but the writer must not report that drop as a forward, or the shutdown
+   * counts this peer as told by a writer which sent nothing - and, for an announcement no
+   * recipient was recorded for, stops waiting for the peers which do need the message.
+   * <p>
+   * The drop is reported to the operator as well, once for the copy the writer took, and that
+   * record is the only account of a peer the topology leaves behind: what it names is asserted
+   * here next to what the writer told the shutdown. The peer is waited on until it is served
+   * from its queue - see {@link #waitForFollowing(ReplicationServerDomain, int)} - so that the
+   * writer is handed the announcement once and the count is not a race.
+   * <p>
+   * The change published after the announcement is the synchronization point: the writer of a
+   * peer takes from one queue in order, so a peer which has received that change has already
+   * dealt with the announcement which precedes it - forwarded it, or dropped it, logged the drop
+   * and told the shutdown to stop waiting for this peer.
+   */
+  @Test(dataProvider = "peerProtocolVersions")
+  public void onlyAPeerWhichCanDecodeTheMessageIsToldTheReplicaWentOffline(
+      short peerVersion, boolean expectedToBeTold) throws Exception
+  {
+    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+    final RecordingShutdownSync shutdownSync = new RecordingShutdownSync();
+    ReplicationServer replicationServer = null;
+    ReplicationBroker broker = null;
+    FakePeerReplicationServer peer = null;
+    try
+    {
+      final int replicationPort = TestCaseUtils.findFreePort();
+      replicationServer = newReplicationServer(shutdownSync,
+          "shutdownSyncPeerVersion" + peerVersion + "Db", PEER_VERSION_RS_ID + peerVersion,
+          replicationPort);
+      broker = openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID);
+      peer = FakePeerReplicationServer.connected(
+          replicationPort, REMOTE_RS_ID, baseDN, EMPTY_DN_GENID, PEER_WINDOW, peerVersion);
+
+      final ReplicationServerDomain domain =
+          replicationServer.getReplicationServerDomain(baseDN, true);
+      waitForConnectedReplicationServer(domain, REMOTE_RS_ID);
+      waitForFollowing(domain, REMOTE_RS_ID);
+      final Future<List<ReplicationMsg>> received = peer.receiveUntil(DeleteMsg.class);
+
+      final CSNGenerator csns = new CSNGenerator(LOCAL_DS_ID, 0);
+      final CSN offlineCSN = csns.newCSN();
+      final ReplicationBroker publishingBroker = broker;
+      final AtomicReference<List<ReplicationMsg>> drained = new AtomicReference<>();
+      final List<String> records = errorLogRecordsOf(() -> {
+        shutdownSync.replicaOfflineMsgSent(baseDN, offlineCSN);
+        publishingBroker.publish(new ReplicaOfflineMsg(offlineCSN));
+        publishingBroker.publish(new DeleteMsg(DN.valueOf("uid=marker," + TEST_ROOT_DN_STRING),
+            csns.newCSN(), "22222222-2222-2222-2222-222222222222"));
+        drained.set(received.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS));
+        return null;
+      });
+
+      final List<ReplicationMsg> msgs = drained.get();
+      assertThat(lastOf(msgs))
+          .as("the peer never received the change published after the announcement, its read "
+              + "ended with: %s", peer.failure())
+          .isInstanceOf(DeleteMsg.class);
+      assertThat(containsReplicaOfflineMsg(msgs, offlineCSN))
+          .as("the peer speaking protocol version %s was %stold that the replica went offline",
+              peerVersion, expectedToBeTold ? "not " : "")
+          .isEqualTo(expectedToBeTold);
+      assertThat(shutdownSync.dispatchedTo())
+          .as("the message was never queued for the peer, so its writer had nothing to report")
+          .containsExactly(REMOTE_RS_ID);
+      if (expectedToBeTold)
+      {
+        assertThat(shutdownSync.forwardedBy())
+            .as("the writer of the peer speaking protocol version %s never reported the forward "
+                + "of a message which did reach it", peerVersion)
+            .containsExactly(REMOTE_RS_ID);
+      }
+      else
+      {
+        assertThat(shutdownSync.forwardedBy())
+            .as("the writer of the peer speaking protocol version %s reported a forward of a "
+                + "message which never reached it", peerVersion)
+            .isEmpty();
+        assertThat(shutdownSync.gaveUpOn())
+            .as("the writer of the peer speaking protocol version %s dropped the message without "
+                + "telling the shutdown to stop waiting for the peer it was queued for",
+                peerVersion)
+            .containsExactly(REMOTE_RS_ID);
+      }
+
+      /*
+       * The records are picked by the id of the message and the peer they name rather than by
+       * their whole text: the address the writer logs is the one the replication server sees the
+       * peer on, which this side does not hold. What the record has to carry is asserted below,
+       * so that an argument dropped or permuted does not pass here.
+       */
+      final List<String> drops = new ArrayList<>();
+      for (String record : records)
+      {
+        if (record.contains("msgID=" + WARN_IGNORING_UPDATE_UNSUPPORTED_BY_PEER.ordinal())
+            && record.contains("to server " + REMOTE_RS_ID))
+        {
+          drops.add(record);
+        }
+      }
+      if (expectedToBeTold)
+      {
+        assertThat(drops)
+            .as("a message the peer speaking protocol version %s did receive was reported as "
+                + "dropped", peerVersion)
+            .isEmpty();
+      }
+      else
+      {
+        assertThat(drops)
+            .as("the drop for the peer speaking protocol version %s should be logged once",
+                peerVersion)
+            .hasSize(1);
+        /*
+         * The severity is part of what this reports: an operator reads the warnings, and a drop
+         * demoted to a trace or an info leaves the announcement lost with nothing said about it.
+         */
+        assertThat(drops.get(0))
+            .as("the drop should be reported as a warning naming this replication server, the "
+                + "domain, the change and the version it has no encoding for")
+            .contains("severity=WARNING")
+            .contains("RS(" + (PEER_VERSION_RS_ID + peerVersion) + ")")
+            .contains(offlineCSN.toString())
+            .contains("domain \"" + baseDN + "\"")
+            .contains("protocol version " + peerVersion);
+      }
+    }
+    finally
+    {
+      closeQuietly(peer);
+      stop(broker);
+      removeQuietly(replicationServer);
+    }
+  }
+
   /**
    * A peer replication server whose handshake is in flight when the shutdown starts must live
    * long enough to be told: it is one of the servers the ReplicaOfflineMsg is forwarded to, and
@@ -900,6 +1062,99 @@
   }
 
   /**
+   * A peer replication server which cannot decode the ReplicaOfflineMsg is nobody to forward it
+   * to either: the message has no encoding for the protocol version that peer negotiated, so no
+   * writer can ever report a forward for it and waiting would only delay the shutdown by the
+   * whole grace period.
+   */
+  @Test
+  public void shutdownIsNotDelayedWhenNoConnectedPeerCanDecodeTheMessage() throws Exception
+  {
+    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
+    ReplicationServer replicationServer = null;
+    try (ServerSocket listen = TestCaseUtils.bindFreePort())
+    {
+      listen.setSoTimeout(SOCKET_TIMEOUT_MS);
+      replicationServer = newReplicationServer(shutdownSync, "shutdownSyncOldPeerDb", 8236);
+      final Session[] sessionPair = connectSessionPair(listen, getReplSessionSecurity());
+      try (Session remoteEnd = sessionPair[0];
+          Session session = sessionPair[1])
+      {
+        session.setProtocolVersion(ProtocolVersion.REPLICATION_PROTOCOL_V7);
+        registerConnectedReplicationServer(replicationServer, baseDN, session);
+
+        final long startTime = System.nanoTime();
+        shutdownSync.replicaOfflineMsgSent(baseDN, newOfflineCSN());
+        replicationServer.shutdown();
+        final long elapsed = elapsedMillis(startTime);
+
+        assertThat(elapsed)
+            .as("the shutdown waited for a peer which cannot decode the message")
+            .isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
+      }
+    }
+    finally
+    {
+      removeQuietly(replicationServer);
+    }
+  }
+
+  /**
+   * A peer which cannot decode the message says nothing about the peers which can: a topology
+   * where one of the connected peers still speaks a protocol version without the
+   * ReplicaOfflineMsg is waited for, for the sake of the peers which do get it.
+   * <p>
+   * What decides is whether <em>any</em> connected peer can be told, not which handler is
+   * iterated first - {@code connectedRSs} is a ConcurrentHashMap, so the order the two peers
+   * below are registered in is not the order they are read in. This case asserts the wait
+   * rather than the peer which caused it for that reason.
+   */
+  @Test
+  public void shutdownWaitsWhenOnlySomeOfTheConnectedPeersCanDecodeTheMessage() throws Exception
+  {
+    final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
+    ReplicationServer replicationServer = null;
+    try (ServerSocket oldPeerListen = TestCaseUtils.bindFreePort();
+        ServerSocket currentPeerListen = TestCaseUtils.bindFreePort())
+    {
+      oldPeerListen.setSoTimeout(SOCKET_TIMEOUT_MS);
+      currentPeerListen.setSoTimeout(SOCKET_TIMEOUT_MS);
+      replicationServer = newReplicationServer(shutdownSync, "shutdownSyncMixedPeersDb", 8237);
+      final Session[] oldPeerPair = connectSessionPair(oldPeerListen, getReplSessionSecurity());
+      final Session[] currentPeerPair =
+          connectSessionPair(currentPeerListen, getReplSessionSecurity());
+      try (Session oldPeerRemoteEnd = oldPeerPair[0];
+          Session oldPeerSession = oldPeerPair[1];
+          Session currentPeerRemoteEnd = currentPeerPair[0];
+          Session currentPeerSession = currentPeerPair[1])
+      {
+        oldPeerSession.setProtocolVersion(ProtocolVersion.REPLICATION_PROTOCOL_V7);
+        currentPeerSession.setProtocolVersion(ProtocolVersion.getCurrentVersion());
+        registerConnectedReplicationServer(
+            replicationServer, baseDN, oldPeerSession, OLD_PEER_RS_ID);
+        registerConnectedReplicationServer(
+            replicationServer, baseDN, currentPeerSession, REMOTE_RS_ID);
+
+        final long startTime = System.nanoTime();
+        shutdownSync.replicaOfflineMsgSent(baseDN, newOfflineCSN());
+        replicationServer.shutdown();
+        final long elapsed = elapsedMillis(startTime);
+
+        assertThat(elapsed)
+            .as("the shutdown stopped waiting although a connected peer can still be told that "
+                + "the replica went offline")
+            .isGreaterThanOrEqualTo(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD);
+      }
+    }
+    finally
+    {
+      removeQuietly(replicationServer);
+    }
+  }
+
+  /**
    * The writer serving a directory server must not hold back the shutdown either: it used to
    * loop on the pending message until the grace period expired, although its handler had already
    * been shut down - which deactivates its consumer and leaves the loop nothing to take.
@@ -1129,6 +1384,37 @@
   }
 
   /**
+   * Waits for the handler of the peer to be served from its in-memory queue rather than from the
+   * changelog.
+   * <p>
+   * A freshly connected peer is behind by definition, and a handler which is catching up reads
+   * its updates from the changelog, where {@code ReplicaCursor} synthesizes a
+   * {@link ReplicaOfflineMsg} from the offline CSN of the replica - on top of the copy
+   * {@code ReplicationServerDomain.put()} queues for it, and again on every refill of its late
+   * queue, since an offline CSN never enters the state of a handler. How many copies of one
+   * announcement its writer takes is then a race, which is a race on how many times the writer
+   * reports the drop.
+   */
+  private void waitForFollowing(final ReplicationServerDomain domain, final int serverId)
+      throws Exception
+  {
+    newConnectionTimer().repeatUntilSuccess(new TestTimer.CallableVoid()
+    {
+      @Override
+      public void call() throws Exception
+      {
+        final ReplicationServerHandler rsHandler = domain.getConnectedRSs().get(serverId);
+        assertThat(rsHandler)
+            .as("the peer replication server %s is no longer connected", serverId).isNotNull();
+        assertThat(rsHandler.isFollowing())
+            .as("the peer replication server %s is still catching up from the changelog",
+                serverId)
+            .isTrue();
+      }
+    });
+  }
+
+  /**
    * Waits for the collocated directory server to be connected and for its handshake to be over -
    * {@link #waitForConnectedReplicationServer(ReplicationServerDomain, int)} says what the
    * registration alone leaves open.
@@ -1262,6 +1548,25 @@
     return TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startTime);
   }
 
+  /** The last message a peer read, null if it read none at all. */
+  private static ReplicationMsg lastOf(List<ReplicationMsg> msgs)
+  {
+    return msgs.isEmpty() ? null : msgs.get(msgs.size() - 1);
+  }
+
+  private static boolean containsReplicaOfflineMsg(List<ReplicationMsg> msgs, CSN offlineCSN)
+  {
+    for (ReplicationMsg msg : msgs)
+    {
+      if (msg instanceof ReplicaOfflineMsg
+          && offlineCSN.equals(((ReplicaOfflineMsg) msg).getCSN()))
+      {
+        return true;
+      }
+    }
+    return false;
+  }
+
   /** The CSN of a message the collocated replica announces, as PendingChanges generates it. */
   private static CSN newOfflineCSN()
   {
@@ -1526,8 +1831,22 @@
     static FakePeerReplicationServer connected(int replicationPort, int serverId, DN baseDN,
         long generationId, int windowSize) throws Exception
     {
+      return connected(replicationPort, serverId, baseDN, generationId, windowSize,
+          ProtocolVersion.getCurrentVersion());
+    }
+
+    /**
+     * A connected peer speaking the given version of the replication protocol.
+     *
+     * @param protocolVersion
+     *          the replication protocol version this peer announces in its start message, so
+     *          that a peer which predates a message type can be reproduced
+     */
+    static FakePeerReplicationServer connected(int replicationPort, int serverId, DN baseDN,
+        long generationId, int windowSize, short protocolVersion) throws Exception
+    {
       final FakePeerReplicationServer peer = new FakePeerReplicationServer(
-          replicationPort, serverId, baseDN, generationId, windowSize);
+          replicationPort, serverId, baseDN, generationId, windowSize, protocolVersion);
       boolean handshaken = false;
       try
       {
@@ -1549,12 +1868,12 @@
     static FakePeerReplicationServer handshaking(
         int replicationPort, int serverId, DN baseDN, long generationId) throws Exception
     {
-      return new FakePeerReplicationServer(
-          replicationPort, serverId, baseDN, generationId, PEER_WINDOW);
+      return new FakePeerReplicationServer(replicationPort, serverId, baseDN, generationId,
+          PEER_WINDOW, ProtocolVersion.getCurrentVersion());
     }
 
     private FakePeerReplicationServer(int replicationPort, int serverId, DN baseDN,
-        long generationId, int windowSize) throws Exception
+        long generationId, int windowSize, short protocolVersion) throws Exception
     {
       this.serverId = serverId;
       this.generationId = generationId;
@@ -1567,6 +1886,8 @@
         socket.setTcpNoDelay(true);
         socket.connect(new InetSocketAddress("127.0.0.1", replicationPort), SOCKET_TIMEOUT_MS);
         newSession = getReplSessionSecurity().createClientSession(socket, SOCKET_TIMEOUT_MS);
+        // the version this peer speaks: the replication server negotiates the oldest of the two
+        newSession.setProtocolVersion(protocolVersion);
 
         newServerURL = "127.0.0.1:" + socket.getLocalPort();
         newSession.publish(new ReplServerStartMsg(serverId, newServerURL, baseDN, windowSize,
@@ -1611,6 +1932,39 @@
     }
 
     /**
+     * Returns every message this peer receives up to and including the first one of the provided
+     * type, or whatever it managed to read if its session ends before that one arrives.
+     */
+    Future<List<ReplicationMsg>> receiveUntil(final Class<? extends ReplicationMsg> markerType)
+    {
+      return reader.submit(new Callable<List<ReplicationMsg>>()
+      {
+        @Override
+        public List<ReplicationMsg> call()
+        {
+          final List<ReplicationMsg> received = new ArrayList<>();
+          try
+          {
+            while (true)
+            {
+              final ReplicationMsg msg = session.receive();
+              received.add(msg);
+              if (markerType.isInstance(msg))
+              {
+                return received;
+              }
+            }
+          }
+          catch (Exception e)
+          {
+            failed(e);
+            return received;
+          }
+        }
+      });
+    }
+
+    /**
      * Returns the first message of the given type this peer receives, or null if its session
      * ends first.
      */

--
Gitblit v1.10.0