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/LDAPReplicationDomain.java |  217 +++++++++++++++++++++++++++++++++++++++++++++++------
 1 files changed, 191 insertions(+), 26 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
index 5e75be0..f68df24 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
@@ -2760,20 +2760,35 @@
        * failing and is eventually given up on: handing it back bare would have this domain
        * ask for it, and restart its session for it, for as long as the server is up.
        *
-       * The changes this thread parked as waiting for another change are left alone: they
-       * are handed to whichever thread clears the change they are waiting for, and that
-       * thread takes them over.
+       * The changes this thread parked as waiting for another change are given back too,
+       * and before the road above runs: that road restarts the session, and a change which
+       * is still owned when the replication server sends it again over it is turned down as
+       * a duplicate - the one delivery which could have taken it over (issue #954).
        *
-       * Which change this thread owns is read before anything is done with it, and that
-       * read takes no lock and allocates nothing: everything below is gated on the answer,
-       * so a lookup which threw in its turn - on the road out of a JVM which has just
-       * refused an allocation - would leave the change listed, uncommitted and owned by a
-       * thread which is about to end, which is the state this whole issue is about.
+       * Which change this thread owns is read first of all, and that read takes no lock and
+       * allocates nothing: everything below is gated on the answer, so a lookup which threw
+       * in its turn - on the road out of a JVM which has just refused an allocation - would
+       * leave the change listed, uncommitted and owned by a thread which is about to end,
+       * which is the state this whole issue is about. It is read before the parked changes
+       * are given back rather than after, because that give-back allocates and can throw on
+       * the same road, and the last resort below can only hand back a change it was told
+       * about.
        */
       CSN owned = null;
       try
       {
         owned = remotePendingChanges.getChangeOwnedByCurrentThread();
+        /*
+         * Asked for without the backoff on the two roads recoverFromReplayFailure() asks
+         * for it without on: a thread which is stopping, and one which an OutOfMemoryError
+         * is ending - the backend is not what is going away. What is asked for is not what
+         * is run: a request is never answered by less than it asked for, so a restart
+         * another road asked for with the backoff, or one given back with it, keeps its
+         * wait whichever thread runs it.
+         */
+        final boolean parkedGivenBack = giveBackParkedChanges(
+            replayThreadShutdown.get() || t instanceof OutOfMemoryError
+                ? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF);
         if (owned != null)
         {
           if (replayThreadShutdown.get() || shutdown.get() || disabled)
@@ -2808,6 +2823,27 @@
             recoverFromReplayFailure(owned, replayThreadShutdown, t instanceof OutOfMemoryError);
           }
         }
+        if (parkedGivenBack && !replayThreadShutdown.get() && !sessionHasAnOwner())
+        {
+          /*
+           * The road the change this thread was replaying took may have run the restart the
+           * give-back asked for - they ask for the same one - and it may have had none to
+           * run: this thread owned no change, or the change it owned was given up on. What
+           * is still requested is run here rather than left standing: the state
+           * checkpointer would run it within its tick, so this is the latency of the
+           * delivery the changes which were handed back wait for, and nothing more - a
+           * request this thread leaves is not lost.
+           *
+           * A thread which is stopping leaves it standing, the way abandonReplay() does:
+           * the state checkpointer runs one restart for every change the threads of the
+           * pool hand back on their way out, rather than each of them running one while
+           * the configuration change which is stopping them waits. A domain whose session
+           * has an owner is left alone the way the give-back left it: nothing was asked
+           * for on that road, and a request another thread left standing is not this
+           * one's to spend on a restart which is refused where it runs.
+           */
+          runRequestedSessionRestarts();
+        }
       }
       catch (Throwable recoveryFailure)
       {
@@ -2819,7 +2855,15 @@
          * so that it is delivered again: this is the last resort, the throwable is rethrown
          * whatever happens here, and the thread this runs on may well be ending on it. A
          * restart which can not run here leaves its request standing, and the state
-         * checkpointer of this domain runs it.
+         * checkpointer of this domain runs it: one which threw asks for itself again on
+         * its way out, and one a domain whose session has an owner would refuse is not
+         * run at all rather than spent on the refusal.
+         *
+         * The restart is asked for once the change is released and not before, the way
+         * every road which releases one asks: asked for first, it could be taken and run
+         * by another thread while this one still owned the change, and the delivery the
+         * new session brought would be turned down as the duplicate of a change a replay
+         * thread owns, with nothing left standing to ask for it again.
          */
         if (owned != null)
         {
@@ -2840,20 +2884,42 @@
           {
             suppress(recoveryFailure, reportFailure);
           }
-          try
+        }
+        try
+        {
+          /*
+           * Outside the guard above: two roads reach here with a request standing and no
+           * change of this thread's to hand back, and both are the parked changes' - the
+           * give-back which released them asks for the restart before it reports them, and
+           * a throw out of the report - the JVM which unwound this replay is out of memory
+           * - leaves the request standing; and a restart the parked road ran and which
+           * threw has asked for one again on its way out. The changes it released are
+           * listed, uncommitted and unowned, so the request is what brings them back, and
+           * this thread is the one there to run it (issue #954). A give-back which threw
+           * before it released anything left the parked changes as they were, owned by this
+           * thread and handed out by getNextUpdate() to whichever thread clears what they
+           * wait for: nothing here can do better for those.
+           *
+           * Not run on a domain whose session has an owner, the way no road of a failed
+           * replay runs it there: the restart is refused where it runs and the request
+           * would be spent on the refusal, while a request left standing is run by the
+           * state checkpointer once the owner is gone - or forgotten with the pending
+           * changes it was made for, when the owner forgets them on its way out.
+           */
+          if (!sessionHasAnOwner())
           {
             runRequestedSessionRestarts();
           }
-          catch (Throwable restartFailure)
-          {
-            /*
-             * Nothing is left to try here: the change is listed, uncommitted and unowned,
-             * and the restart which threw has asked for one again, so the session restart
-             * the state checkpointer runs delivers it again. This goes with the throwable
-             * which is rethrown below rather than being reported on its own.
-             */
-            suppress(recoveryFailure, restartFailure);
-          }
+        }
+        catch (Throwable restartFailure)
+        {
+          /*
+           * Nothing is left to try here: the changes are listed, uncommitted and unowned,
+           * and the restart which threw has asked for one again, so the session restart
+           * the state checkpointer runs delivers them again. This goes with the throwable
+           * which is rethrown below rather than being reported on its own.
+           */
+          suppress(recoveryFailure, restartFailure);
         }
         // The error which unwound the replay is the one reported, whatever the give-back
         // ran into on top of it.
@@ -3387,12 +3453,14 @@
              * an OutOfMemoryError of its own, still owns its change: it is given back
              * counted, and the thread ends on this error rather than on the one it stepped
              * over. A replay which committed owns nothing anymore - commit() cleared the
-             * owner, and the index the give-back reads, in the same step - so the give-back
-             * is a no-op, and rightly so: a change which is in the data is not one to ask
-             * for again. What that road steps over is getNextUpdate() below, so the changes
-             * parked behind the committed change wait for the next replay of this domain to
-             * hand them out. That is the trade #923 asks for: a thread which met this error
-             * is not to carry on, not even for them.
+             * owner, and the index the give-back reads, in the same step - so the change it
+             * was replaying is not given back, and rightly so: a change which is in the data
+             * is not one to ask for again. What that road steps over is getNextUpdate()
+             * below, which hands out the changes parked behind the committed change: the
+             * ones this thread parked are given back on the way out of replay() and the
+             * session is restarted for them (issue #954), the ones other threads parked wait
+             * for the next replay of this domain to hand them out. That is the trade #923
+             * asks for: a thread which met this error is not to carry on, not even for them.
              */
             throw e;
           }
@@ -3968,6 +4036,25 @@
   }
 
   /**
+   * Gives back the changes a replay thread which is stopping parked in this domain.
+   * <p>
+   * Called by that thread on its way out (issue #986). The restart which brings them back
+   * is asked for and left standing, the way the thread leaves the request it makes for the
+   * change it abandons: the threads of the pool are stopped one after the other and joined,
+   * and each running a restart on its way out would have the configuration change which is
+   * stopping them wait for one restart per thread. The state checkpointer of this domain
+   * runs one restart for the lot within its tick, and holds it while a total update runs
+   * over the session, in either direction - a change delivered again before the pool which
+   * replaces this thread is up waits in the replay queue for it. Asked for without the
+   * backoff: what went away is a replay thread, not the backend, and these changes were
+   * never applied here.
+   */
+  void giveBackChangesParkedByStoppingThread()
+  {
+    giveBackParkedChanges(SessionRestart.NOW);
+  }
+
+  /**
    * Restarts the session as long as changes which could not be replayed are waiting to be
    * delivered again.
    */
@@ -4076,6 +4163,84 @@
   }
 
   /**
+   * Gives back the changes this replay thread parked as waiting for another change, on the
+   * way out of a replay which was unwound.
+   * <p>
+   * A parked change is handed out again by {@code getNextUpdate()} alone, which every
+   * replay loop of this 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 whose
+   * replay was unwound is not on that road anymore - it takes the next delivery off the
+   * replay queue - so a change it parked would be left owned by a thread which is not
+   * coming back to it, while every redelivery of it is refused as a duplicate. On a domain
+   * which then goes quiet that change is where this replica's ServerState, and every change
+   * behind it from every master, stops (issue #954).
+   * <p>
+   * They are handed back without a failure being counted against them: they were never
+   * applied here, so the give-up budget which decides when this replica skips a change it
+   * can not apply is not this delivery's to spend, the way it is not for a change abandoned
+   * by a replay thread which is stopping.
+   * <p>
+   * The delivery which carried one published no ack - the ack of a parked change is
+   * published by the delivery which replays it - so it is counted as processed here, the
+   * way a delivery which is dropped rather than replayed is: that count is of the
+   * deliveries this replica took off the session, and these are over. The window they hold
+   * is not given back either, and does not need to be: the session they came over is about
+   * to be restarted, and a session which starts is given its receive window anew.
+   * <p>
+   * On a domain whose session has an owner - the domain itself, going away, or a total
+   * update into it, from the moment it is asked for - they are released and nothing more,
+   * the way {@code abandonReplay()} hands a change back on that road (see
+   * {@link #sessionHasAnOwner()}): there is no session of this thread's to restart, and
+   * the restart it would ask for is refused where it runs. The domain forgets its pending
+   * changes on its way down, the import forgets them at its end, and a change released
+   * for a total update which never begins stays listed until the next failed replay of
+   * this domain restarts the session, which has the replication server send it again. A
+   * line which says the replication server sends the change again would not hold on any
+   * of these - a server which is shutting down abandons every change in flight, and none
+   * of them is delivered again before it is started back.
+   * <p>
+   * A replay thread which is stopping calls this through
+   * {@link #giveBackChangesParkedByStoppingThread()}, for every domain of this server: what
+   * it parked would be left owned by a thread which does not exist anymore, and every
+   * redelivery of a change a replay thread owns is refused as a duplicate (issue #986). The
+   * request it makes here is left standing for the state checkpointer, the way that thread
+   * leaves the request it makes for the change it abandons.
+   *
+   * @param restart what the session restart is asked for as: with the backoff a failing
+   *          backend is owed, or without it on a thread which is stopping or which an
+   *          OutOfMemoryError is ending
+   * @return whether any change was handed back: a change which nobody owns is one only a
+   *         new delivery brings back, so the caller runs the restart asked for them - on
+   *         a thread which is not stopping, and on a domain whose session has no owner
+   */
+  private boolean giveBackParkedChanges(SessionRestart restart)
+  {
+    final List<CSN> parked = remotePendingChanges.releaseParkedChangesOwnedByCurrentThread();
+    if (parked.isEmpty())
+    {
+      return false;
+    }
+    if (sessionHasAnOwner())
+    {
+      // The domain owns its session, or a total update does: both forget the pending
+      // changes, and neither leaves a session for this thread to restart.
+      return true;
+    }
+    /*
+     * Asked for before the changes are reported: a throw out of the report - the JVM which
+     * unwound this replay is out of memory - must not lose the restart which is what brings
+     * them back.
+     */
+    sessionRestarts.request(restart);
+    for (CSN csn : parked)
+    {
+      incProcessedUpdates();
+      logger.info(NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK, csn, getBaseDN());
+    }
+    return true;
+  }
+
+  /**
    * Gives a change back to the replication server when this replay thread stops before it
    * could apply it.
    * <p>

--
Gitblit v1.10.0