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/RemotePendingChanges.java |  112 ++++++++++++++++++++++++++++++++++++++++++++++++++++---
 1 files changed, 105 insertions(+), 7 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/RemotePendingChanges.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/RemotePendingChanges.java
index ec3ef7b..62629a0 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/RemotePendingChanges.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/RemotePendingChanges.java
@@ -17,7 +17,11 @@
  */
 package org.opends.server.replication.plugin;
 
+import static java.util.Collections.*;
+
+import java.util.ArrayList;
 import java.util.Iterator;
+import java.util.List;
 import java.util.NoSuchElementException;
 import java.util.SortedMap;
 import java.util.SortedSet;
@@ -90,8 +94,10 @@
    * this issue is about (issue #922).
    * <p>
    * A thread is entered here when it takes a change over and removed when it gives it back,
-   * applies it, or parks it as waiting for another change - the parked ones are handed to
-   * whichever thread clears what they wait for, so they are not this one's to give back.
+   * applies it, or parks it as waiting for another change - a parked change is not the one
+   * this thread is replaying, and giving it back is
+   * {@link #releaseParkedChangesOwnedByCurrentThread()}, which reads the changes which are
+   * waiting rather than this index (issue #954).
    * <p>
    * The entry of a thread is written by that thread and by nobody else, and that - not the
    * lock - is what keeps the writes apart: the park in {@link #addDependency(PendingChange)}
@@ -545,9 +551,10 @@
    * Returns the CSN of the change the calling thread is replaying, when it still owns one.
    * <p>
    * A thread owns the change it is replaying and the ones it parked as waiting for another
-   * change. The parked ones are left out: they are handed to whichever thread clears the
-   * change they are waiting for, and that thread takes them over, so giving one back here
-   * would have the same change handed to two threads (issue #922).
+   * change. The parked ones are left out: they are not the change this thread is replaying,
+   * and giving one back is more than dropping its owner - it has to be unparked in the same
+   * step, or it would be handed out by two roads at once, which is what
+   * {@link #releaseParkedChangesOwnedByCurrentThread()} does (issues #922 and #954).
    * <p>
    * It is a plain read of {@link #changeBeingReplayed}: no lock is taken and nothing is
    * allocated. This is what the give-back on the way out of an unwound replay asks first,
@@ -570,6 +577,95 @@
   }
 
   /**
+   * Gives back the changes the calling thread parked as waiting for another change, and
+   * takes them out of the changes which are waiting in the same step.
+   * <p>
+   * A parked change stays owned by the thread which parked it while that thread goes on
+   * to the changes which follow: {@link #getNextUpdate()} is what hands it out again, to
+   * whichever replay thread clears the change it was waiting for, and that thread takes it
+   * over. A replay which is unwound leaves the thread which parked it without that road -
+   * it takes the next delivery off the replay queue instead - so the change would be left
+   * owned by a thread which is never coming back to it, and every redelivery of a change a
+   * replay thread owns is refused as a duplicate (issue #954).
+   * <p>
+   * Unparking a change and giving it back is one step, under both locks, so that only one
+   * road can hand it out: a change which was released while it is still listed as waiting
+   * would be handed to the thread {@link #getNextUpdate()} gives it to and to the thread
+   * which takes over the delivery which follows - the double replay the ownership is there
+   * to prevent (OPENDJ-1115).
+   * <p>
+   * The changes stay listed and uncommitted, and stay among the changes the newer ones are
+   * checked against, the way a change whose replay failed does: they are not in the data,
+   * so they hold this domain's ServerState back and the changes which follow them keep
+   * waiting for them.
+   * <p>
+   * The changes another thread parked are left alone: a change is given back by the thread
+   * which owns it and by nobody else (issue #922). That thread may be inside the dependency
+   * checks which parked it - they park a change once per dependency it has - so a change
+   * released under it would be listed as waiting again a moment later, and handed out while
+   * the delivery which took it over is being replayed.
+   *
+   * @return the CSNs of the changes it gave back, oldest first; empty when this thread owns
+   *         no parked change - the changes a thread parked stay its own, whichever replay
+   *         parked them, until {@link #getNextUpdate()} hands them to the thread which
+   *         cleared what they wait for or they are given back here
+   */
+  List<CSN> releaseParkedChangesOwnedByCurrentThread()
+  {
+    final Thread current = Thread.currentThread();
+    /*
+     * The second lock is taken inside the try of the first: taking a lock which is held
+     * by another thread allocates the node this one waits on, and this runs on the road
+     * out of a JVM which has just refused an allocation. A throw out of the second lock
+     * would otherwise unwind past the first with that one held by a thread which is
+     * ending, and every road which lists, commits or gives back a change would wait for
+     * it for good.
+     */
+    pendingChangesWriteLock.lock();
+    try
+    {
+      dependentChangesLock.lock();
+      try
+      {
+        if (dependentChanges.isEmpty())
+        {
+          // Nothing is waiting, which is the state every replay but a handful leaves behind.
+          return emptyList();
+        }
+        /*
+         * Sized for every change which is waiting, so that the one allocation of this
+         * method past the locks is made before anything is taken out of the set. The rule
+         * getNextUpdate() states for itself holds here: an allocation which fails once a
+         * change has been unparked and released - and this runs on the road out of a JVM
+         * which has just refused one - would have moved that change out of the hands which
+         * hand it out again, with the caller never told that it did.
+         */
+        final List<CSN> released = new ArrayList<>(dependentChanges.size());
+        final Iterator<PendingChange> it = dependentChanges.iterator();
+        while (it.hasNext())
+        {
+          final PendingChange change = it.next();
+          if (change.isOwnedBy(current))
+          {
+            it.remove();
+            change.setOwner(null);
+            released.add(change.getCSN());
+          }
+        }
+        return released;
+      }
+      finally
+      {
+        dependentChangesLock.unlock();
+      }
+    }
+    finally
+    {
+      pendingChangesWriteLock.unlock();
+    }
+  }
+
+  /**
    * Get the first update in the list that have some dependencies cleared.
    * <p>
    * The change is handed to the calling thread, which owns it from then on: it is
@@ -672,8 +768,10 @@
        * parked one is handed to the thread which clears what it waits for, and one which is
        * not listed here anymore is gone with the pending changes of a domain which was
        * disabled. The owner stays as it is - it is what has getNextUpdate() hand the change
-       * over rather than leave it to nobody - and the give-back on the way out of an
-       * unwound replay leaves it alone (issue #922).
+       * over rather than leave it to nobody - and the give-back of the change a replay was
+       * unwound on leaves it alone (issue #922). What hands a parked change back is
+       * releaseParkedChangesOwnedByCurrentThread(), which unparks it in the same step so
+       * that the two roads can not hand it out at once (issue #954).
        */
       changeBeingReplayed.remove(Thread.currentThread(), dependentChange.getCSN());
     }

--
Gitblit v1.10.0