From 600df926522919b63bfca2816ef9588b6f1c6e34 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Tue, 15 Sep 2026 17:09:57 +0000
Subject: [PATCH] [#950] Announce a ReplicaOfflineMsg before it is published, not after it may have been forwarded (#978)

---
 opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/PendingChangesTest.java |  159 +++++++++++++++++++++++++++++++++++++++++++++++++++-
 1 files changed, 154 insertions(+), 5 deletions(-)

diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/PendingChangesTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/PendingChangesTest.java
index 23c026f..1842e63 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/PendingChangesTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/PendingChangesTest.java
@@ -16,12 +16,14 @@
 package org.opends.server.replication.plugin;
 
 import static org.mockito.Matchers.any;
+import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 import static org.testng.Assert.*;
 
+import org.forgerock.opendj.ldap.DN;
 import org.mockito.ArgumentCaptor;
 import org.opends.server.DirectoryServerTestCase;
 import org.opends.server.replication.common.CSN;
@@ -29,23 +31,36 @@
 import org.opends.server.replication.protocol.LDAPUpdateMsg;
 import org.opends.server.replication.protocol.ReplicaOfflineMsg;
 import org.opends.server.replication.protocol.UpdateMsg;
+import org.opends.server.replication.service.DSRSShutdownSync;
 import org.opends.server.replication.service.ReplicationDomain;
 import org.opends.server.types.operation.PluginOperation;
+import org.testng.annotations.BeforeClass;
 import org.testng.annotations.Test;
 
 /**
  * Tests the bookkeeping a replica does on its own changes: they are published in the order of
- * their CSNs, and the announcement that the replica goes offline is only reported as sent when
- * it really was published.
+ * their CSNs, and the announcement that the replica goes offline is made before the message is
+ * published, since the shutdown of a collocated replication server waits for that message to be
+ * forwarded - and stands only for a message the broker reports as written.
  * <p>
- * These tests need no server: the changes are built by a CSNGenerator, which reads the time
- * service, and the time service is up as soon as its class is loaded.
+ * These tests need no server: the changes are built by a CSNGenerator, which needs nothing but
+ * a server id.
  */
 @SuppressWarnings("javadoc")
 @Test(groups = { "precommit", "replication" }, sequential = true)
 public class PendingChangesTest extends DirectoryServerTestCase
 {
   private static final int SERVER_ID = 42;
+  /** A peer replication server the collocated one relays the message to. */
+  private static final int RS_ID = 11;
+
+  private static DN baseDN;
+
+  @BeforeClass
+  public static void classSetup() throws Exception
+  {
+    baseDN = DN.valueOf("dc=example,dc=com");
+  }
 
   @Test
   public void replicaOfflineMsgTheBrokerPublishedIsReportedAsSent() throws Exception
@@ -62,6 +77,47 @@
   }
 
   /**
+   * The collocated replication server forwards the message as soon as it is on the wire, so an
+   * announcement made after the publish is one the forward found nothing to clear: nothing will
+   * ever remove it, and the shutdown waits out its whole grace period for a message which the
+   * topology already has.
+   * <p>
+   * The forward is reported from inside publish(), which is where the message reaches the
+   * session, so the race is reproduced rather than waited for.
+   */
+  @Test
+  public void theReplicaOfflineMsgIsAnnouncedBeforeItIsPublished() throws Exception
+  {
+    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
+    final ReplicationDomain domain = mock(ReplicationDomain.class);
+    forwardWhilePublishing(domain, shutdownSync);
+    final PendingChanges pendingChanges = newPendingChanges(domain, shutdownSync);
+
+    pendingChanges.putReplicaOfflineMsg();
+
+    assertTrue(shutdownSync.canShutdown(baseDN),
+        "the message was forwarded, so nothing must hold the shutdown back any longer");
+  }
+
+  /**
+   * The announcement of a message the broker took stands until a peer forwards it: the
+   * withdrawal is for the message the broker refused, and an announcement taken back after a
+   * publish which succeeded would leave the shutdown nothing to wait for.
+   */
+  @Test
+  public void theAnnouncementOfAPublishedMessageStandsUntilItIsForwarded() throws Exception
+  {
+    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
+    final PendingChanges pendingChanges =
+        newPendingChanges(domainWhichPublishes(true), shutdownSync);
+
+    pendingChanges.putReplicaOfflineMsg();
+
+    assertFalse(shutdownSync.canShutdown(baseDN),
+        "the message went out and nobody has forwarded it yet, so the shutdown must wait for it");
+  }
+
+  /**
    * The broker writes nothing when it has no usable session, when the changes which come before
    * this one still have to be republished by the recovery, or when it is stopped in between - and
    * what was not written must not be reported as sent: the shutdown of a collocated replication
@@ -80,6 +136,27 @@
   }
 
   /**
+   * The announcement is made before the message is published, and the broker may refuse it
+   * once it is: an announcement which stayed would be one nobody will ever forward, and the
+   * shutdown would wait out its whole grace period for a message which never left. It is
+   * therefore withdrawn - and it is a withdrawal, not an announcement which was never made: the
+   * shutdown is held back while the broker holds the message.
+   */
+  @Test
+  public void theAnnouncementOfAReplicaOfflineMsgTheBrokerRefusedIsWithdrawn() throws Exception
+  {
+    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
+    final ReplicationDomain domain = mock(ReplicationDomain.class);
+    refuseWhilePublishing(domain, shutdownSync);
+    final PendingChanges pendingChanges = newPendingChanges(domain, shutdownSync);
+
+    assertNull(pendingChanges.putReplicaOfflineMsg(), "the broker refused the message");
+
+    assertTrue(shutdownSync.canShutdown(baseDN),
+        "the message never reached the wire, so nothing must hold the shutdown back");
+  }
+
+  /**
    * The message carries the newest CSN of the replica, so a change which is still in flight
    * holds it back, and the broker is never even asked to publish it.
    */
@@ -117,6 +194,32 @@
   }
 
   /**
+   * The announcement follows the publication rather than the queueing, so a message which a
+   * change in flight holds back is not announced: neither while it waits, nor when the change
+   * which held it back completes and the message is given up on. Announcing it either time
+   * would leave the shutdown waiting out its whole grace period for a forward which cannot
+   * happen.
+   */
+  @Test
+  public void theReplicaOfflineMsgHeldBackByAChangeInFlightIsNeverAnnounced() throws Exception
+  {
+    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
+    final ReplicationDomain domain = domainWhichPublishes(true);
+    final PendingChanges pendingChanges = newPendingChanges(domain, shutdownSync);
+    final CSN inFlight = pendingChanges.putLocalOperation(newLocalOperation());
+
+    pendingChanges.putReplicaOfflineMsg();
+
+    assertTrue(shutdownSync.canShutdown(baseDN),
+        "the message is still queued behind a change in flight, and nothing was announced");
+
+    pendingChanges.commitAndPushCommittedChanges(inFlight, mock(LDAPUpdateMsg.class));
+
+    assertTrue(shutdownSync.canShutdown(baseDN),
+        "the message was given up on with the change which held it back, and never announced");
+  }
+
+  /**
    * A change the broker refused leaves the pending changes all the same: the replica has done
    * it, its ServerState says so, and it is by finding that state ahead of the one its
    * replication server reports that the next session republishes the change from the historical
@@ -137,9 +240,55 @@
     assertEquals(pendingChanges.size(), 0, "and is not queued for a second attempt");
   }
 
+  /**
+   * Reports the forward of the message from within the publish which puts it on the wire, after
+   * checking from there that the announcement is already in place, and publishes it. The forward
+   * is the one of a peer the message was never recorded as queued for, which is what a forward
+   * racing the announcement looks like.
+   */
+  private void forwardWhilePublishing(
+      final ReplicationDomain domain, final DSRSShutdownSync shutdownSync)
+  {
+    doAnswer(invocation -> {
+      final UpdateMsg msg = (UpdateMsg) invocation.getArguments()[0];
+      if (msg instanceof ReplicaOfflineMsg)
+      {
+        assertFalse(shutdownSync.canShutdown(baseDN),
+            "the message must be announced before it is published");
+        shutdownSync.replicaOfflineMsgForwarded(baseDN, msg.getCSN(), RS_ID);
+      }
+      return true;
+    }).when(domain).publish(any(UpdateMsg.class));
+  }
+
+  /**
+   * Refuses to publish the message, the way a broker with no usable session does, after checking
+   * from within the publish that the announcement is already in place.
+   */
+  private void refuseWhilePublishing(
+      final ReplicationDomain domain, final DSRSShutdownSync shutdownSync)
+  {
+    doAnswer(invocation -> {
+      assertFalse(shutdownSync.canShutdown(baseDN),
+          "the message must be announced before it is published");
+      return false;
+    }).when(domain).publish(any(UpdateMsg.class));
+  }
+
   private PendingChanges newPendingChanges(ReplicationDomain domain)
   {
-    return new PendingChanges(new CSNGenerator(SERVER_ID, 0), domain);
+    return newPendingChanges(domain, new DSRSShutdownSync());
+  }
+
+  /**
+   * The pending changes of a replica whose domain announces itself through the shutdown sync,
+   * with the announcer LDAPReplicationDomain hands its own pending changes.
+   */
+  private PendingChanges newPendingChanges(
+      final ReplicationDomain domain, final DSRSShutdownSync shutdownSync)
+  {
+    return new PendingChanges(new CSNGenerator(SERVER_ID, 0), domain,
+        new ShutdownSyncAnnouncer(shutdownSync, baseDN));
   }
 
   /** A domain whose broker accepts, or refuses, whatever it is given to publish. */

--
Gitblit v1.10.0