From 776339a8c63c0bd8c83cfe8a619f6146794358e8 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Sat, 12 Sep 2026 11:53:58 +0000
Subject: [PATCH] [#922] Give a change back when the replay which owns it is unwound (#958)
---
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/RemotePendingChanges.java | 186 ++++++++++++++++++++++++++++++++++++++++++----
1 files changed, 169 insertions(+), 17 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 1d54c51..ec3ef7b 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
@@ -23,6 +23,8 @@
import java.util.SortedSet;
import java.util.TreeMap;
import java.util.TreeSet;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ConcurrentSkipListSet;
import java.util.concurrent.locks.ReentrantLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
@@ -73,6 +75,32 @@
*/
private final ConcurrentSkipListSet<PendingChange> activeAndDependentChanges = new ConcurrentSkipListSet<>();
+ /**
+ * The change each replay thread is replaying right now, read by the give-back on the way
+ * out of a replay which was unwound.
+ * <p>
+ * It is an index of what {@link PendingChange#isOwnedBy(Thread)} already says rather than
+ * a second copy of it: it is written in the same locked step wherever a change is taken
+ * over or given back, and every road which acts on what it answers checks the ownership
+ * of the change again. What it buys is the read. That read is made on a road an
+ * {@code OutOfMemoryError} leads to, and looking the change up by walking the pending
+ * changes takes two locks and allocates an iterator - an allocation a JVM which has just
+ * refused one may well refuse again, and the change would then be left listed,
+ * uncommitted and owned by a thread which is not replaying it anymore, which is the wedge
+ * 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.
+ * <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)}
+ * clears it under the read lock, where every other writer holds the write lock, and it
+ * races nothing for it. {@link #clear()} is the one writer of every entry, and it holds
+ * both locks.
+ */
+ private final ConcurrentMap<Thread, CSN> changeBeingReplayed = new ConcurrentHashMap<>();
+
private final ReentrantReadWriteLock pendingChangesLock = new ReentrantReadWriteLock(true);
private final ReentrantReadWriteLock.ReadLock pendingChangesReadLock = pendingChangesLock.readLock();
private final ReentrantReadWriteLock.WriteLock pendingChangesWriteLock = pendingChangesLock.writeLock();
@@ -209,9 +237,25 @@
/**
* Mark an update message as committed.
+ * <p>
+ * A change another replay thread owns is not this thread's to record: it is reported as
+ * a change which is not here, the way one which is not listed anymore is. A thread
+ * decides that a change is in the data, or that it is to be given up on, and records
+ * that decision a turn of this lock later - long enough for the change to have been
+ * handed back, delivered again and taken over in between. Recording it then would
+ * advance the ServerState over a change which is not in the data yet (issue #889) and
+ * have the thread which is applying it right now fail to commit.
+ * <p>
+ * The give-back on the way out of an unwound replay is what makes this reachable: it
+ * runs wherever the replay was left, so the checks which keep it from taking a change
+ * away from the thread which owns it now belong on every road which reads ownership,
+ * not only on {@link #replayFailed(CSN)} (issue #922).
*
* @param csn
* The CSN of the update message that must be set as committed.
+ * @throws NoSuchElementException
+ * if there is no change with that CSN for this thread to record: it is not
+ * listed as pending anymore, or another replay thread owns it
*/
public void commit(CSN csn)
{
@@ -219,12 +263,14 @@
try
{
PendingChange curChange = pendingChanges.get(csn);
- if (curChange == null)
+ if (curChange == null
+ || (curChange.isOwned() && !curChange.isOwnedBy(Thread.currentThread())))
{
throw new NoSuchElementException();
}
curChange.setCommitted(true);
- curChange.setOwned(false);
+ curChange.setOwner(null);
+ changeBeingReplayed.remove(Thread.currentThread(), csn);
activeAndDependentChanges.remove(curChange);
final Iterator<PendingChange> it = pendingChanges.values().iterator();
@@ -271,6 +317,12 @@
* The changes another replay thread is applying right now are left alone: they are
* about to commit, and forgetting them would have their {@code commit()} fail, the
* ServerState stay behind them and the replication server replay them a second time.
+ * <p>
+ * Only the thread which owns the change gives it back. A release which arrives from
+ * another one is a release of a change which has been taken over since - the failure of
+ * a delivery is reported after the change it carried was handed to the delivery which
+ * follows it - and taking the change away from the thread which is replaying it right
+ * now is the double replay the ownership is there to prevent (issue #922).
*
* @param csn the CSN of the change whose replay failed
*/
@@ -280,9 +332,10 @@
try
{
final PendingChange change = pendingChanges.get(csn);
- if (change != null && !change.isCommitted())
+ if (change != null && !change.isCommitted() && change.isOwnedBy(Thread.currentThread()))
{
- change.setOwned(false);
+ change.setOwner(null);
+ changeBeingReplayed.remove(Thread.currentThread(), csn);
}
}
finally
@@ -339,9 +392,13 @@
* the CSN of the change whose replay failed
* @param nowMs
* when it failed, on a clock which only moves forward
- * @return the failures of the change, or {@code null} when it is not listed as an
- * uncommitted change anymore, which happens when the domain was disabled while
- * it was being replayed: there is no change left here to give up on
+ * @return the failures of the change, or {@code null} when there is no change left here
+ * for this thread to give up on: it is not listed as an uncommitted change
+ * anymore, which happens when the domain was disabled while it was being
+ * replayed, or another replay thread owns it, which is that change having been
+ * delivered again and taken over while this thread was on its way to reporting
+ * on it. Spending the give-up budget of a change another thread is applying
+ * would have this replica skip a change which is being written (issue #922)
*/
public ReplayFailure recordReplayFailure(CSN csn, long nowMs)
{
@@ -349,7 +406,8 @@
try
{
final PendingChange change = pendingChanges.get(csn);
- if (change == null || change.isCommitted())
+ if (change == null || change.isCommitted()
+ || (change.isOwned() && !change.isOwnedBy(Thread.currentThread())))
{
return null;
}
@@ -434,6 +492,7 @@
pendingChanges.clear();
dependentChanges.clear();
activeAndDependentChanges.clear();
+ changeBeingReplayed.clear();
failingChanges = 0;
}
finally
@@ -465,8 +524,16 @@
{
return false;
}
- change.setOwned(true);
+ /*
+ * Listed as being replayed before it is owned, and not the other way round: the
+ * caller enters the replay - where the give-back on the way out lives - once this
+ * returns, so nothing which allocates must run between the owner being stamped and
+ * that. A change listed here without an owner is the state a failed replay leaves
+ * behind, and the dependency checks which read this set do not read the owner.
+ */
activeAndDependentChanges.add(change);
+ changeBeingReplayed.put(Thread.currentThread(), change.getCSN());
+ change.setOwner(Thread.currentThread());
return true;
}
finally
@@ -475,7 +542,40 @@
}
}
/**
+ * 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).
+ * <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,
+ * and it runs on a road an {@code OutOfMemoryError} leads to - a lookup which allocated
+ * could be refused in its turn, and a give-back which does not know which change to give
+ * back leaves it listed, uncommitted and owned by a thread which is not replaying it
+ * anymore, which is the wedge this issue is about.
+ * <p>
+ * An answer which is out of date is safe: every road which acts on it - {@code commit()},
+ * {@link #recordReplayFailure(CSN, long)} and {@link #replayFailed(CSN)} - checks the
+ * ownership of the change again under the write lock, and is a no-op for a change this
+ * thread does not own anymore.
+ *
+ * @return the CSN of the change this thread is replaying, or {@code null} when it does
+ * not own one anymore - the road its replay took gave it back, or it was applied
+ */
+ CSN getChangeOwnedByCurrentThread()
+ {
+ return changeBeingReplayed.get(Thread.currentThread());
+ }
+
+ /**
* 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
+ * replayed by whichever replay thread cleared the change it was waiting for rather than
+ * by the one which parked it, and a change is given back by the thread which owns it
+ * and by nobody else (issue #922).
*
* @return The LDAPUpdateMsg to be handled.
*/
@@ -485,22 +585,65 @@
dependentChangesLock.lock();
try
{
- if (!dependentChanges.isEmpty() && !pendingChanges.isEmpty())
+ if (!hasChangeToHandOut())
{
- PendingChange firstDependentChange = dependentChanges.first();
- if (pendingChanges.firstKey().isNewerThanOrEqualTo(firstDependentChange.getCSN()))
- {
- dependentChanges.remove(firstDependentChange);
- return firstDependentChange.getLDAPUpdateMsg();
- }
+ /*
+ * Nothing is waiting, or what waits is still held back by the changes before it.
+ * This is called at the end of every replay, by every replay thread, so the answer
+ * is looked for under the read lock: taking the write lock here would have a
+ * backlog of waiting changes serialize the replay of the changes which have none.
+ */
+ return null;
}
- return null;
}
finally
{
dependentChangesLock.unlock();
pendingChangesReadLock.unlock();
}
+
+ /*
+ * There is one to hand out, and handing it out writes its owner, which is written
+ * under the write lock as the rest of the state of a change is. It is looked for again
+ * under that lock: another replay thread may have been handed it in between.
+ */
+ pendingChangesWriteLock.lock();
+ dependentChangesLock.lock();
+ try
+ {
+ if (hasChangeToHandOut())
+ {
+ final PendingChange firstDependentChange = dependentChanges.first();
+ /*
+ * Entered as the change this thread is replaying before it is taken out of the
+ * ones which are waiting, for the same reason markInProgress() enters it before it
+ * stamps the owner: an allocation which fails here must leave the change where it
+ * was rather than take it out of the hands which would hand it out again.
+ */
+ changeBeingReplayed.put(Thread.currentThread(), firstDependentChange.getCSN());
+ dependentChanges.remove(firstDependentChange);
+ firstDependentChange.setOwner(Thread.currentThread());
+ return firstDependentChange.getLDAPUpdateMsg();
+ }
+ return null;
+ }
+ finally
+ {
+ dependentChangesLock.unlock();
+ pendingChangesWriteLock.unlock();
+ }
+ }
+
+ /**
+ * Returns whether the first change waiting for another one can be replayed now, that is
+ * whether every change before it has left the pending changes.
+ */
+ @GuardedBy("pendingChangesLock, dependentChangesLock")
+ private boolean hasChangeToHandOut()
+ {
+ return !dependentChanges.isEmpty()
+ && !pendingChanges.isEmpty()
+ && pendingChanges.firstKey().isNewerThanOrEqualTo(dependentChanges.first().getCSN());
}
/**
@@ -524,6 +667,15 @@
{
dependentChanges.add(dependentChange);
}
+ /*
+ * Whichever of the two it was, this thread is not replaying that change anymore: a
+ * 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).
+ */
+ changeBeingReplayed.remove(Thread.currentThread(), dependentChange.getCSN());
}
finally
{
--
Gitblit v1.10.0