From bf0f4f294f477e393b59ef7f160795e00aa02cc9 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Thu, 24 Sep 2026 07:57:37 +0000
Subject: [PATCH] Send what a replication session's publisher left queued when it is closed (#1035)

---
 opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java |  263 ++++++++++++++++++++++++++++++++++++++++++++++++++-
 1 files changed, 254 insertions(+), 9 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java
index 5965219..bd867c0 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java
@@ -13,6 +13,7 @@
  *
  * Copyright 2006-2009 Sun Microsystems, Inc.
  * Portions Copyright 2011-2016 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
  */
 package org.opends.server.replication.protocol;
 
@@ -30,12 +31,14 @@
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.locks.Lock;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.zip.DataFormatException;
 
 import javax.net.ssl.SSLSocket;
 
+import org.forgerock.i18n.LocalizableMessage;
 import org.forgerock.i18n.slf4j.LocalizedLogger;
 import org.opends.server.api.DirectoryThread;
 import org.opends.server.types.HostPort;
@@ -48,6 +51,19 @@
 {
   private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
 
+  /**
+   * How long a close spends sending what the publisher thread left queued in {@code sendQueue},
+   * in milliseconds.
+   * <p>
+   * The same 5 s as {@code DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD}, which is how long a
+   * shutdown is already willing to wait for one of these messages - the announcement that a
+   * replica went offline - to be forwarded. A close has no reason to wait longer for it than the
+   * shutdown which is waiting on the close. It is a budget per close, though, not per shutdown:
+   * a shutdown closes its sessions one after another, so it can pay this once for each peer which
+   * is alive and not reading.
+   */
+  private static final long DRAIN_BUDGET_MS = 5000;
+
   private final Socket plainSocket;
   private final SSLSocket secureSocket;
   private final InputStream plainInput;
@@ -101,6 +117,14 @@
   private final CountDownLatch latch = new CountDownLatch(1);
 
   /**
+   * How many buffers the publisher thread took off {@code sendQueue} and failed to write. After a
+   * failed write that thread goes on taking the queue until the session is closed, and every
+   * write after the first fails as well, so these are messages the peer was not told about just
+   * as much as what is still queued - see {@link #close()}.
+   */
+  private final AtomicInteger publisherFailedWrites = new AtomicInteger();
+
+  /**
    * Creates a new Session.
    *
    * @param socket
@@ -140,6 +164,10 @@
   /**
    * This method is called when the session with the remote must be closed.
    * This object won't be used anymore after this method is called.
+   * <p>
+   * A message which was published on this session but which its publisher thread had not sent yet
+   * is sent here rather than dropped, within the budget of {@link #DRAIN_BUDGET_MS}. See {@link
+   * #sendWhatThePublisherLeftQueued()}.
    */
   @Override
   public void close()
@@ -165,6 +193,18 @@
       Thread.currentThread().interrupt();
     }
 
+    /*
+     * Re-read the error rather than answer with the snapshot taken before the join: a publisher
+     * whose send() failed while this thread was joining it recorded the error there, and what
+     * follows - the drain and the StopMsg - is what must not be written to a socket which has
+     * already failed. Reading it before the join left both writing to one, the drain naming its
+     * own failure rather than the one the publisher had recorded.
+     */
+    synchronized (stateLock)
+    {
+      localSessionError = sessionError;
+    }
+
     // Perform close outside of critical section.
     if (logger.isTraceEnabled())
     {
@@ -186,18 +226,72 @@
       }
     }
 
-    // V4 protocol introduces a StopMsg to properly end communications.
-    if (localSessionError == null
-        && protocolVersion >= ProtocolVersion.REPLICATION_PROTOCOL_V4)
+    if (localSessionError != null)
     {
-      try
+      /*
+       * Nothing more is written to a session which already failed, neither what the publisher left
+       * queued nor the StopMsg: writing more to it cannot work. What it gives up on is still
+       * reported: what is left in the queue, and what the publisher took off it and failed to
+       * write - the write which failed first, and every one it went on failing until this close
+       * stopped it, which the join above makes complete. It is done without taking publishLock,
+       * so that a thread blocked in a write of this socket is released by the close of the
+       * sockets below rather than waited for.
+       */
+      isRunning.set(false);
+      final int notSent = publisherFailedWrites.get() + takeWhatIsLeftQueued();
+      if (notSent > 0)
       {
-        publish(new StopMsg());
+        reportQueueNotSent(notSent, "the session had already failed: " + localSessionError);
       }
-      catch (final IOException ignored)
+      StaticUtils.close(plainSocket, secureSocket);
+      return;
+    }
+
+    /*
+     * The publisher thread has stopped and what it had not sent is still in the queue. Send it,
+     * rather than let the close drop it: nothing publishes these again, and the StopMsg which
+     * follows leaves the peer reading an orderly close with no sign that anything was missing.
+     *
+     * This thread is not the only one which can write the socket here - a ServerWriter or a
+     * HeartbeatThread outlives this close. It is this thread, not run(), which takes the session
+     * off the queueing branch of publish(), and it does so under publishLock and keeps the lock
+     * across the queue and the StopMsg: until then a publish() concurrent with the close takes
+     * the queueing branch, and from then on it takes the synchronous branch and waits for the
+     * lock. Either way nothing newer is written between two of the drained messages - see
+     * sendWhatThePublisherLeftQueued() below.
+     */
+    publishLock.lock();
+    try
+    {
+      isRunning.set(false);
+      sendWhatThePublisherLeftQueued();
+
+      /*
+       * Re-read again: a write of the drain which failed recorded its error, and the StopMsg must
+       * not be written to a socket which has already failed either.
+       */
+      synchronized (stateLock)
       {
-        // Ignore errors on close.
+        localSessionError = sessionError;
       }
+
+      // V4 protocol introduces a StopMsg to properly end communications.
+      if (localSessionError == null
+          && protocolVersion >= ProtocolVersion.REPLICATION_PROTOCOL_V4)
+      {
+        try
+        {
+          publish(new StopMsg());
+        }
+        catch (final IOException ignored)
+        {
+          // Ignore errors on close.
+        }
+      }
+    }
+    finally
+    {
+      publishLock.unlock();
     }
 
     StaticUtils.close(plainSocket, secureSocket);
@@ -206,6 +300,101 @@
 
 
   /**
+   * Sends the buffers the publisher thread had not sent when it stopped, so that a close of the
+   * session does not drop them.
+   * <p>
+   * Called from {@link #close()} once the publisher has been joined, with {@code publishLock}
+   * held. A queued message is already encoded for this peer's protocol version - {@link
+   * #publish(ReplicationMsg)} did that before queueing it - so there is nothing to decide here
+   * beyond how long to keep trying.
+   * <p>
+   * The whole queue goes out under {@code publishLock}, because the publisher is not the only
+   * thread which writes this socket: a {@code ServerWriter} is joined only after the close which
+   * gets here (ServerHandler.shutdown() closes the session before joining it, and ServerReader's
+   * finally closes it before stopping the handler), and a {@code HeartbeatThread} is shut down
+   * after it too. Their {@code publish()} takes the synchronous branch once the session is off
+   * the queueing one, so without the lock a newer message could be written between two of these
+   * older ones - which a peer replication server answers by dropping the older ones at debug
+   * level, its log file refusing a record which would break its key ordering. The close takes
+   * the session off the queueing branch under the same lock it holds here, so such a
+   * {@code publish()} either queues ahead of that and is drained here, or waits for the drain and
+   * lands after it, or returns at the door, having seen the close. One which read the close as
+   * not yet begun can still be descheduled before its buffer is queued and queue it after the
+   * last poll below; it checks for that once it has, and takes the buffer back and reports it -
+   * see {@link #publish(ReplicationMsg)}. The wait for the lock itself is not part of the budget
+   * below, no more than it is for the {@code StopMsg} which follows.
+   * <p>
+   * The budget bounds how many messages a close spends on a peer which is reading slowly. A peer
+   * which is reading pays none of it; a peer which answers with a reset pays one failed write. It
+   * is checked between messages, so a single write which blocks past the budget still runs to
+   * completion - for a peer which has stopped reading, or which vanished without a reset, that is
+   * for as long as TCP keeps the connection alive: bounding it needs a non-blocking socket, which
+   * this session is not. The budget is per close, and a shutdown closes its sessions one after
+   * another - ReplicationServerDomain.stopAllServers() stops each handler in turn on one thread -
+   * so a domain with several peers which are alive and not reading pays it once per such peer.
+   * On the road which does not shut the whole server down the wait is paid under the lock of the
+   * replication domain - ReplicationServerDomain.stopServer() holds it across the handler
+   * shutdown which closes this session - where a handshake meanwhile waiting on that lock times
+   * out and the broker retries. What is given up on is reported rather than dropped in silence,
+   * that being the part of this which cost the most to diagnose.
+   */
+  private void sendWhatThePublisherLeftQueued()
+  {
+    final long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(DRAIN_BUDGET_MS);
+    byte[] buffer;
+    while ((buffer = sendQueue.poll()) != null)
+    {
+      if (System.nanoTime() - deadline >= 0)
+      {
+        reportQueueNotSent(1 + takeWhatIsLeftQueued(), "the peer did not read them within "
+            + DRAIN_BUDGET_MS + " ms");
+        return;
+      }
+      try
+      {
+        send(buffer);
+      }
+      catch (final IOException e)
+      {
+        /*
+         * send() has recorded the error; the rest of the queue cannot go out either. The
+         * exception is named rather than traced: this is reached whenever a peer which has
+         * announced it is leaving closes before the drain reaches it, where what the write
+         * failed with is the whole of what a reader of the log needs - a directory server
+         * re-reads these from the changelog when it reconnects, a replication server does not.
+         */
+        reportQueueNotSent(1 + takeWhatIsLeftQueued(),
+            "the write failed with " + e.getClass().getName() + ": " + e.getMessage());
+        return;
+      }
+    }
+  }
+
+  /**
+   * Empties the queue a close gives up on, and answers how many buffers it held. Taking them
+   * rather than counting them is what keeps a {@code publish()} which queued a buffer late from
+   * reporting it a second time: it reports only a buffer it still finds in the queue.
+   */
+  private int takeWhatIsLeftQueued()
+  {
+    int taken = 0;
+    while (sendQueue.poll() != null)
+    {
+      taken++;
+    }
+    return taken;
+  }
+
+  /** Says which messages a close could not hand to the peer, and why. */
+  private void reportQueueNotSent(final int count, final String reason)
+  {
+    logger.warn(LocalizableMessage.raw(
+        "The replication session %s was closed with %d message(s) which had been published on it "
+        + "but not yet sent, and the peer was not told about them: %s",
+        getName(), count, reason));
+  }
+
+  /**
    * This methods allows to determine if the session close was initiated
    * on this Session.
    *
@@ -322,6 +511,10 @@
           // Avoid blocking forever so that we can check for session closure.
           if (sendQueue.offer(buffer, 100, TimeUnit.MILLISECONDS))
           {
+            if (!isRunning.get())
+            {
+              takeBackWhatWasQueuedTooLate(buffer);
+            }
             return;
           }
         }
@@ -338,6 +531,34 @@
     }
   }
 
+  /**
+   * Takes a buffer back out of the queue if nothing is going to send it, and reports it.
+   * <p>
+   * {@code publish()} reads the close as not yet begun and queues the buffer after that, with no
+   * lock in between: descheduled there, it can queue the buffer after a close has stopped the
+   * publisher thread and drained the queue, and nothing would then send it or say so. Once the
+   * session is off the queueing branch, the lock here is the one the close holds while it takes
+   * the session off that branch and drains the queue, so what is still in the queue after it is
+   * what nothing sends. A buffer queued before the session came off the queueing branch is left
+   * to the drain - it cannot have seen the flag cleared - and a buffer the drain or the close
+   * already took is not found here, so nothing is reported twice.
+   */
+  private void takeBackWhatWasQueuedTooLate(final byte[] buffer)
+  {
+    publishLock.lock();
+    try
+    {
+      if (sendQueue.remove(buffer))
+      {
+        reportQueueNotSent(1, "it was queued after the publisher of the session had stopped");
+      }
+    }
+    finally
+    {
+      publishLock.unlock();
+    }
+  }
+
   /** Sends a replication message already encoded to the socket.
    *
    * @param buffer
@@ -526,7 +747,20 @@
   @Override
   public void run()
   {
-    isRunning.set(true);
+    synchronized (stateLock)
+    {
+      /*
+       * A close which came before the start has already run, and it is not coming back to clear
+       * the flag - nor is the end of this method, which leaves that to the close. Set, the flag
+       * would keep every later publish() on the queueing branch, which returns at once on a
+       * closed session as if the message had been queued; unset, publish() stays on the
+       * synchronous branch, which fails on the closed socket.
+       */
+      if (!closeInitiated)
+      {
+        isRunning.set(true);
+      }
+    }
     latch.countDown();
     if (logger.isTraceEnabled())
     {
@@ -551,10 +785,21 @@
       catch (IOException e)
       {
         setSessionError(e);
+        publisherFailedWrites.incrementAndGet();
         needClosing = true;
       }
     }
-    isRunning.set(false);
+    /*
+     * A close clears the flag itself, under publishLock, once it has joined this thread - see
+     * close(). Clearing it here would open a window between the end of this thread and the drain
+     * of the close, in which a publish() takes the synchronous branch and writes a newer message
+     * ahead of the queue. Only a loop which ended without a close - an interrupt from elsewhere -
+     * clears it here, so that publish() does not go on queueing onto a queue nobody sends.
+     */
+    if (!closeInitiated)
+    {
+      isRunning.set(false);
+    }
     if (needClosing)
     {
       close();

--
Gitblit v1.10.0