From a2a74542282fa2eba683661058786625b50c6dc7 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Thu, 24 Sep 2026 06:37:59 +0000
Subject: [PATCH] [#1048] Hold the session restart a released change asks for while a total update runs (#1049)
---
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java | 175 +++++++++++++++++++++++++++++++++++++++++++++-------------
1 files changed, 136 insertions(+), 39 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 1238dbd..678fc25 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
@@ -2786,7 +2786,7 @@
* another road asked for with the backoff, or one given back with it, keeps its
* wait whichever thread runs it.
*/
- final boolean parkedGivenBack = giveBackParkedChanges(
+ final List<CSN> parkedGivenBack = giveBackParkedChanges(
replayThreadShutdown.get() || t instanceof OutOfMemoryError
? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF);
if (owned != null)
@@ -2822,8 +2822,17 @@
*/
recoverFromReplayFailure(owned, replayThreadShutdown, t instanceof OutOfMemoryError);
}
+ /*
+ * The change is handed back and counted as the road it took counts it: the last
+ * resort below speaks for a give-back which did not run, and a throw out of what
+ * follows - the restart the parked changes are run with, or the line which says it
+ * is held - is not one. Left set, that throw would have the last resort report the
+ * change as "released without its failure being counted", which is the one line an
+ * operator acts on, and it would be false.
+ */
+ owned = null;
}
- if (parkedGivenBack && !replayThreadShutdown.get() && !sessionHasAnOwner())
+ if (!parkedGivenBack.isEmpty() && !replayThreadShutdown.get() && !sessionHasAnOwner())
{
/*
* The road the change this thread was replaying took may have run the restart the
@@ -2838,11 +2847,31 @@
* 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.
+ * has an owner is left alone the way the give-back left it - it asked for nothing
+ * there, which is why it handed back no change to run a restart for - and a
+ * request another thread left standing is not this one's to spend on a restart
+ * which is refused where it runs.
*/
- runRequestedSessionRestarts();
+ if (!runRequestedSessionRestarts())
+ {
+ /*
+ * The changes above were reported as given back to a replication server which
+ * "still owns it and sends it again", and it does not send them yet: a total
+ * update is being processed over the session, the restart which brings them back
+ * waits for it, and the state checkpointer runs it once it is over - the same
+ * hold, and the same line, a change whose replay failed is reported with. One
+ * line per change, the way the give-back reports them: these are the changes
+ * this one is about, and nothing else says they wait.
+ *
+ * Built on the road out of a JVM which has run out of memory too, where the line
+ * of a failed replay is not: the give-back has already built one line per change
+ * on that road, so what this asks the JVM for is not memory it was spared.
+ */
+ for (CSN csn : parkedGivenBack)
+ {
+ logger.info(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE, csn, getBaseDN());
+ }
+ }
}
}
catch (Throwable recoveryFailure)
@@ -2856,7 +2885,8 @@
* 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: one which threw asks for itself again on
- * its way out, and one a domain whose session has an owner would refuse is not
+ * its way out, one a total update is being processed over the session is not run
+ * while it lasts, 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
@@ -2888,12 +2918,14 @@
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
+ * Outside the guard above: three roads reach here with a request standing and no
+ * change of this thread's to hand back, and all 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
+ * - leaves the request standing; a restart the parked road ran and which threw
+ * has asked for one again on its way out; and a throw out of the line which says
+ * a total update holds that restart leaves the request standing too, where it is
+ * refused until the total update is over. 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
@@ -3540,8 +3572,10 @@
if (replayFailed && recoverFromReplayFailure(msg.getCSN(), replayThreadShutdown))
{
// The ack has been published and the change is given back: the replication server
- // delivers it again, now or - while a total update owns the session - after the
- // import restarts it. There is nothing left to replay here.
+ // delivers it again, now or - while a total update is being processed over the
+ // session - once it is over, when the restart which was held runs or, on the import
+ // direction, when the session is started from the reloaded state. There is nothing
+ // left to replay here.
return;
}
@@ -3822,8 +3856,11 @@
*
* @param csn the CSN of the change which could not be replayed
* @param failure how long, and over how many deliveries, its replay has been failing
+ * @return whether the warning was written, so that a line which qualifies it - the one
+ * which says the restart it announced is held - is written with it and folded
+ * with it
*/
- private void logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure failure)
+ private boolean logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure failure)
{
final long now = monotonicNowInMs();
final long lastLogged = lastReplayRetryWarningTime.get();
@@ -3832,6 +3869,7 @@
{
logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts(),
failure.getFailingForMs(), foldedReplayRetryWarnings.getAndSet(0));
+ return true;
}
else
{
@@ -3845,6 +3883,7 @@
foldedReplayRetryWarnings.incrementAndGet();
logger.trace("Could not replay change %s in domain %s: delivery %d, failing for %d ms",
csn, getBaseDN(), failure.getAttempts(), failure.getFailingForMs());
+ return false;
}
}
@@ -4001,6 +4040,7 @@
return true;
}
+ boolean warned = false;
if (!outOfMemory)
{
/*
@@ -4012,7 +4052,7 @@
* unlogged: the error ends the replay thread, and the uncaught exception handler of
* DirectoryThread writes the line and raises the alert for it.
*/
- logReplayRetryWarning(csn, failure);
+ warned = logReplayRetryWarning(csn, failure);
}
/*
* This change is not owned by anyone anymore, so the session has to be restarted for
@@ -4031,7 +4071,32 @@
*/
sessionRestarts.request(replayThreadShutdown.get() || outOfMemory
? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF);
- runRequestedSessionRestarts();
+ if (!runRequestedSessionRestarts() && warned)
+ {
+ /*
+ * The warning above said the session is being restarted for the change, and it is not
+ * yet: a total update is being processed over that session - almost always an export
+ * from this replica, since a total update into it owns the session and is refused
+ * above, except by an import which claims its context between that read and this one
+ * (issue #1041) - and the restart waits for it, for as long as the total update takes.
+ * Said on its own, so that a change which is not delivered again for minutes is not a
+ * change nobody asked for.
+ *
+ * Written where that warning was written and nowhere else: it qualifies that line, so
+ * a backend which fails every delivery of a long export would otherwise be one of
+ * these per delivery while the warnings they qualify are folded into a count - the
+ * very repetition the throttle is there to fold. It is not built at all on the road
+ * out of a JVM which has run out of memory, for the reason the warning is not built
+ * there: warned is false on it.
+ *
+ * On the import road of the window above the request is cleared by importBackend()
+ * rather than run, so the restart this line announces never runs there. The change is
+ * delivered again all the same, which is what an operator reads this for: the import
+ * loads the ServerState of the exporter, and the session started at its end asks for
+ * everything that state does not cover.
+ */
+ logger.info(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE, csn, getBaseDN());
+ }
return true;
}
@@ -4051,15 +4116,48 @@
*/
void giveBackChangesParkedByStoppingThread()
{
+ // What it handed back is not read here: this road runs no restart for them, so it has
+ // nothing to say about one which waits.
giveBackParkedChanges(SessionRestart.NOW);
}
/**
* Restarts the session as long as changes which could not be replayed are waiting to be
- * delivered again.
+ * delivered again - unless a total update is being processed, in which case the requests
+ * are left standing for the state checkpointer to run once it is over.
+ * <p>
+ * A restart stops the session the total update runs over, in either direction. An import
+ * into this replica reads its entries from that session and would end on the ones which
+ * had arrived - and {@code disabled} does not say an import is running, since
+ * {@code preBackendImport()} keeps the backend events this domain is the cause of from
+ * disabling it. An export from this replica publishes its entries over it, and
+ * {@code exportLDIFEntry()} gives the export up as
+ * {@code ERR_INIT_RS_DISCONNECTION_DURING_EXPORT} once the broker has been stopped under
+ * it, which leaves the replica it was initializing to be initialized again: minutes on a
+ * large backend, spent for a change which would have waited. So the change waits: the
+ * request stays standing, the state checkpointer comes for it once a second and runs it
+ * as soon as the total update is over ({@link #runPendingSessionRestart()}), and the
+ * replication server delivers the change again then. The ServerState waits with it, and
+ * the replay of this domain keeps running in the meantime.
+ * <p>
+ * A total update which begins between this read and any of the stops this call makes is
+ * not seen here, and is cut by it: the read and the claim of the import/export context
+ * share no lock, which is issue #1041 on the import side. The window is the call rather
+ * than a few statements of it: {@code restartSession()} stops the session before it waits
+ * its backoff out, so a request taken later in the loop stops the session a backoff wait
+ * after the read - up to {@link #MAX_REPLAY_RETRY_DELAY_IN_MS}, and the restart of a
+ * request another thread has just been told is held is one of those. Before this the whole
+ * of the total update was that window.
+ *
+ * @return {@code false} when a total update is being processed and the requests were left
+ * standing for the state checkpointer, {@code true} otherwise
*/
- private void runRequestedSessionRestarts()
+ private boolean runRequestedSessionRestarts()
{
+ if (ieRunning())
+ {
+ return false;
+ }
/*
* The outer loop is what makes a request which was made while this thread was giving
* up the recovery its own: the thread which made it found the recovery taken and left
@@ -4099,6 +4197,7 @@
replayFailureRecovery.set(false);
}
}
+ return true;
}
/**
@@ -4114,22 +4213,15 @@
* topology, with the changes it did not replay owned by the replication server and its
* ServerState stopped behind them.
* <p>
- * Not run while a total update is being processed, in either direction: a restart stops
- * the session the total update runs over. An import into this replica reads its entries
- * from that session and would end on the ones which had arrived - and {@code disabled}
- * does not say an import is running, since {@code preBackendImport()} keeps the backend
- * events this domain is the cause of from disabling it. An export from this replica
- * publishes its entries over it, and {@code exportLDIFEntry()} gives the export up as
- * {@code ERR_INIT_RS_DISCONNECTION_DURING_EXPORT} once the broker has been stopped
- * under it, which leaves the replica it was initializing to be initialized again. This
- * thread is the one which can afford to wait: the request stays standing, and it comes
- * back here once a second, so the restart is run as soon as the total update is over.
- * The change the restart was asked for waits for as long as the total update takes, and
- * the ServerState with it; the replay of this domain keeps running in the meantime.
+ * While a total update is being processed, in either direction, the restart is not run -
+ * no restart asked for by a released change is, see {@link #runRequestedSessionRestarts()}
+ * - and this thread is the one which can afford to wait for it: the request stays
+ * standing, and it comes back here once a second, so the restart is run as soon as the
+ * total update is over.
*/
private void runPendingSessionRestart()
{
- if (shutdown.get() || disabled || ieRunning() || !sessionRestarts.isPending())
+ if (shutdown.get() || disabled || !sessionRestarts.isPending())
{
return;
}
@@ -4210,22 +4302,27 @@
* @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
+ * @return the changes it handed back and asked the restart for, oldest first: 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 - and reports what that restart did for them. Empty when
+ * this thread had parked none, and empty on a domain whose session has an owner,
+ * where nothing is asked for. The list is the one the release allocated: nothing is
+ * allocated for the answer on the road out of a JVM which has run out of memory
*/
- private boolean giveBackParkedChanges(SessionRestart restart)
+ private List<CSN> giveBackParkedChanges(SessionRestart restart)
{
final List<CSN> parked = remotePendingChanges.releaseParkedChangesOwnedByCurrentThread();
if (parked.isEmpty())
{
- return false;
+ return parked;
}
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;
+ // changes, and neither leaves a session for this thread to restart. Nothing was asked
+ // for here, so the caller has nothing of this road's to run or to report.
+ return Collections.emptyList();
}
/*
* Asked for before the changes are reported: a throw out of the report - the JVM which
@@ -4238,7 +4335,7 @@
incProcessedUpdates();
logger.info(NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK, csn, getBaseDN());
}
- return true;
+ return parked;
}
/**
--
Gitblit v1.10.0