From b8c3c195995500ba0351a3736ad43147a6036dd0 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Fri, 25 Sep 2026 06:26:29 +0000
Subject: [PATCH] [#1055] Report a ReplicaOfflineMsg forwarded once it is written to the peer, not once it is queued (#1057)

---
 opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java |   54 ++++++++++++++++++++++++++++++++++++++++--------------
 1 files changed, 40 insertions(+), 14 deletions(-)

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 d1edac1..9f68e4b 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
@@ -17,9 +17,11 @@
  */
 package org.opends.server.replication.server;
 
+import java.io.IOException;
 import java.net.SocketException;
 
 import org.forgerock.i18n.LocalizableMessage;
+import org.forgerock.opendj.ldap.DN;
 import org.opends.server.api.DirectoryThread;
 import org.forgerock.i18n.slf4j.LocalizedLogger;
 import org.opends.server.replication.common.ServerStatus;
@@ -120,24 +122,14 @@
                 replicationServerDomain.getBaseDN(), handler.getServerId());
           }
         }
+        else if (updateMsg instanceof ReplicaOfflineMsg && !handler.isDataServer())
+        {
+          forwardReplicaOfflineMsg((ReplicaOfflineMsg) updateMsg);
+        }
         else
         {
           // Publish the update to the remote server using a protocol version it supports
           session.publish(updateMsg);
-          /*
-           * Only the forward to a peer RS ends the wait of the shutdown: what the grace period
-           * buys is the rest of the topology learning that the replica went offline. A directory
-           * server is never handed this message - ReplicationServerDomain.put() does not queue
-           * it for one, and DataServerHandler.updateServerState() drops the one the changelog
-           * cursor of a directory server which is catching up synthesizes from the offline CSN
-           * of the replica (issue #1029) - so the guard says whose forward counts rather than
-           * telling two deliveries apart.
-           */
-          if (updateMsg instanceof ReplicaOfflineMsg && !handler.isDataServer())
-          {
-            dsrsShutdownSync.replicaOfflineMsgForwarded(
-                replicationServerDomain.getBaseDN(), updateMsg.getCSN(), handler.getServerId());
-          }
         }
       }
     }
@@ -170,6 +162,40 @@
     }
   }
 
+  /**
+   * Publishes a ReplicaOfflineMsg to the peer replication server, and reports the forward to the
+   * shutdown which may be waiting for it.
+   * <p>
+   * Only the forward to a peer RS ends the wait of the shutdown: what the grace period buys is
+   * the rest of the topology learning that the replica went offline. A directory server is never
+   * handed this message - ReplicationServerDomain.put() does not queue it for one, and
+   * DataServerHandler.updateServerState() drops the one the changelog cursor of a directory
+   * server which is catching up synthesizes from the offline CSN of the replica (issue #1029) -
+   * so the guard of the caller says whose forward counts rather than telling two deliveries apart.
+   * <p>
+   * The forward is reported once the message has been written to the peer, not once it is queued
+   * for the thread of the session: the shutdown closes the session as soon as its wait ends, and
+   * Session.close() sends what is still queued only once the write it joins has returned, and
+   * only within a budget of its own, so a message reported forwarded while it was queued behind
+   * one the peer had not read yet would end the wait for a peer which had not been told, and
+   * leave its delivery to that budget rather than to the grace period. A message the session
+   * refuses - one published while the session is being closed - will never be written, and the
+   * shutdown must not wait for it. One the protocol version of the peer cannot carry is refused
+   * by the session as well, but does not get this far: isUpdateMsgFiltered() drops it and says so
+   * first (issue #1014).
+   */
+  private void forwardReplicaOfflineMsg(final ReplicaOfflineMsg msg) throws IOException
+  {
+    final DN baseDN = replicationServerDomain.getBaseDN();
+    final int serverId = handler.getServerId();
+    final boolean accepted = session.publish(msg,
+        () -> dsrsShutdownSync.replicaOfflineMsgForwarded(baseDN, msg.getCSN(), serverId));
+    if (!accepted)
+    {
+      dsrsShutdownSync.replicaOfflineMsgNotForwarded(baseDN, serverId);
+    }
+  }
+
   private boolean isUpdateMsgFiltered(UpdateMsg updateMsg)
   {
     if (!updateMsg.isEncodableFor(handler.getProtocolVersion()))

--
Gitblit v1.10.0