From 80481f756d71bd58b4bda627758dcd774e9d5dcd Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Fri, 18 Sep 2026 14:45:54 +0000
Subject: [PATCH] [#986] Give back the changes a replay thread the pool stopped had parked (#988)

---
 opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/ReplayThread.java |   87 +++++++++++++++++++++++++++++++++++++++++--
 1 files changed, 83 insertions(+), 4 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/ReplayThread.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/ReplayThread.java
index 0dddf3a..27cffe6 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/ReplayThread.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/ReplayThread.java
@@ -24,6 +24,7 @@
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.locks.ReentrantLock;
+import java.util.function.Consumer;
 
 import org.opends.server.api.DirectoryThread;
 import org.forgerock.i18n.slf4j.LocalizedLogger;
@@ -39,6 +40,15 @@
 public class ReplayThread extends DirectoryThread
 {
   private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
+  /**
+   * The give-back a thread runs on every domain of this server on its way out, held here
+   * rather than written where it is run: a method reference is linked, and its instance
+   * made, where it is first run, and this one is first run on the way out of a thread -
+   * which an OutOfMemoryError may be ending, on the road this give-back is there for. Made
+   * when this class is loaded instead, on a thread which can allocate (issue #986).
+   */
+  private static final Consumer<LDAPReplicationDomain> GIVE_BACK_PARKED_CHANGES =
+      LDAPReplicationDomain::giveBackChangesParkedByStoppingThread;
 
   private final BlockingQueue<UpdateToReplay> updateToReplayQueue;
   private final ReentrantLock switchQueueLock;
@@ -77,6 +87,53 @@
       logger.trace("Replication Replay thread starting.");
     }
 
+    try
+    {
+      replayUntilStopped();
+    }
+    finally
+    {
+      /*
+       * The changes this thread parked as waiting for another change are handed out again
+       * by getNextUpdate() alone, which every replay loop of a domain runs once it is done
+       * with a change: a parked change is replayed by whichever thread clears the change it
+       * was waiting for. A thread which is stopping is not on that road anymore, so what it
+       * parked would be left owned by a thread which does not exist, while every redelivery
+       * of a change a replay thread owns is refused as a duplicate: on a domain which then
+       * goes quiet that change is where the ServerState of this replica, and every change
+       * behind it from every master, stops (issue #986).
+       *
+       * Given back by the thread which owns them, so that the rule every road which reads
+       * ownership follows holds on this one as well: a change is given back by the thread it
+       * was handed to and by nobody else (issue #922). It is also the one place which sees
+       * them all - the pool is shared by every domain of this server, while a replay knows
+       * only the domain it was replaying for.
+       *
+       * The session which brings them back is asked for and left standing, in every domain
+       * which got something back, and the state checkpointer of each of them runs it within
+       * its tick: a thread on its way out is not held for a session - the threads of the
+       * pool are stopped one after the other and joined, and each running a restart of its
+       * own would have the configuration change which is stopping them wait for one restart
+       * per thread - and a change delivered again before the pool which replaces this one
+       * is up waits in the replay queue for it. A thread which an OutOfMemoryError is ending
+       * gives back here what it parked in the domains it was not replaying for, on the same
+       * terms; the change it was replaying, and what it had parked in that same domain, were
+       * given back and asked for again on its way out of replay().
+       */
+      giveBackParkedChanges();
+    }
+    if (logger.isTraceEnabled())
+    {
+      logger.trace("Replication Replay thread stopping.");
+    }
+  }
+
+  /**
+   * Takes the deliveries of the domains of this server off the shared replay queue and
+   * replays them, until this thread is stopped.
+   */
+  private void replayUntilStopped()
+  {
     while (!shutdown.get())
     {
       try
@@ -145,9 +202,31 @@
         logger.error(ERR_EXCEPTION_REPLAYING_REPLICATION_MESSAGE, stackTraceToSingleLineString(t));
       }
     }
-    if (logger.isTraceEnabled())
-    {
-      logger.trace("Replication Replay thread stopping.");
-    }
+  }
+
+  /**
+   * Gives back the changes this thread parked as waiting for another change, in every
+   * domain of this server.
+   * <p>
+   * A change which is given back stays listed and uncommitted, the way a change whose replay
+   * failed does: it is not in the data, so it holds the ServerState of its domain back and
+   * the changes which follow it keep waiting for it, until the delivery which takes it over
+   * replays it.
+   * <p>
+   * Every domain gets its turn whatever one of them threw: what can throw here is an
+   * allocation, on the way out of a thread an OutOfMemoryError may be ending - the iterator
+   * over the domains, before any of them is reached, then for each of them the list of what
+   * it released, made before anything is released, or the report of a change once it is,
+   * and between two domains the list a second failure is recorded in under the first, which
+   * the loop guards on its own - and the domains which follow would otherwise be left with
+   * changes owned by a thread which does not exist anymore, the state this give-back is
+   * for. A domain which threw past the release has asked for its restart already: the
+   * request is made before the report. The first failure is thrown once the loop is over,
+   * so that the uncaught exception handler of {@link DirectoryThread} writes the line and
+   * raises the alert.
+   */
+  private void giveBackParkedChanges()
+  {
+    MultimasterReplication.forEachDomain(GIVE_BACK_PARKED_CHANGES);
   }
 }

--
Gitblit v1.10.0