From d07bb318976a8e1c65b1c26b431b0909380ce473 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Thu, 17 Sep 2026 11:57:18 +0000
Subject: [PATCH] [#942] Warn once per interval that a change is being retried, not once per delivery (#982)
---
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java | 203 +++++++++++++++++++++++++++++++++++++++++++++-----
1 files changed, 182 insertions(+), 21 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 409083e..5e75be0 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
@@ -340,6 +340,14 @@
*/
private static final long UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MS = 60000;
/**
+ * How long the warning telling that a change is being asked for again is not logged
+ * again, for the reason the alert above is not sent again and for one more: a change
+ * which keeps failing is asked for again every {@link #MAX_REPLAY_RETRY_DELAY_IN_MS}
+ * at the slowest, for as long as its give-up budget lasts - and how long that is has
+ * been the administrator's to set since issue #901.
+ */
+ private static final long REPLAY_RETRY_WARNING_INTERVAL_IN_MS = 60000;
+ /**
* What the ack of a delivery whose replay ran out of memory says did not apply the
* change.
* <p>
@@ -386,11 +394,14 @@
@GuardedBy("sessionRestartBackoff")
private long sessionRestartBackoffWakes;
/**
- * How many times in a row the session was restarted without a change being replayed in
- * between. The backoff is computed from this rather than from the failures of the
- * change which happens to open the recovery: an outage fails every change in flight,
- * and the ones which are sent for the first time would otherwise keep the wait at its
- * shortest for as long as the outage lasts.
+ * How many times in a row the session was restarted since this replica last replayed a
+ * change while nothing was failing, or was last disabled or imported into. A replay of
+ * another change while one still fails does not reset it, and neither does giving up on
+ * the last one failing: see {@link #resetReplayFailureTracking()} and
+ * {@link #skipUnreplayableChange(CSN, LocalizableMessage)}. The backoff is computed from
+ * this rather than from the failures of the change which happens to open the recovery:
+ * an outage fails every change in flight, and the ones which are sent for the first time
+ * would otherwise keep the wait at its shortest for as long as the outage lasts.
*/
private final AtomicInteger consecutiveSessionRestarts = new AtomicInteger();
/**
@@ -481,14 +492,38 @@
*/
private volatile long replayDrainTimeoutInMs = REPLAY_DRAIN_TIMEOUT_IN_MS;
/**
- * Stands for "the alert about a change this replica gave up on was never sent". The
- * time it is compared with only moves forward from an origin which is arbitrary, so
- * zero is not far enough in the past to say it.
+ * How long this domain does not warn again about a change it asks for again. Only the
+ * tests, which can not wait out {@link #REPLAY_RETRY_WARNING_INTERVAL_IN_MS} between two
+ * warnings, set another value.
*/
- private static final long UNREPLAYED_CHANGE_ALERT_NEVER_SENT = Long.MIN_VALUE / 2;
+ private volatile long replayRetryWarningIntervalInMs = REPLAY_RETRY_WARNING_INTERVAL_IN_MS;
+ /**
+ * Stands for "a replay failure was never reported yet", whether by the alert about a
+ * change this replica gave up on or by the warning about a change it asks for again.
+ * The time it is compared with only moves forward from an origin which is arbitrary,
+ * so zero is not far enough in the past to say it.
+ */
+ private static final long REPLAY_FAILURE_NEVER_REPORTED = Long.MIN_VALUE / 2;
/** When the alert about a change this replica gave up on was last sent. */
private final AtomicLong lastUnreplayedChangeAlertTime =
- new AtomicLong(UNREPLAYED_CHANGE_ALERT_NEVER_SENT);
+ new AtomicLong(REPLAY_FAILURE_NEVER_REPORTED);
+ /** When the warning about a change this replica asks for again was last logged. */
+ private final AtomicLong lastReplayRetryWarningTime =
+ new AtomicLong(REPLAY_FAILURE_NEVER_REPORTED);
+ /**
+ * How many failed deliveries were folded into no warning since the last one was logged,
+ * or since this replica last stopped failing: the next warning says how many it stands
+ * for. Forgotten when a change is replayed or given up on and nothing is failing anymore
+ * - see {@link #resetReplayFailureTracking()} and
+ * {@link #skipUnreplayableChange(CSN, LocalizableMessage)}, only the first of which
+ * forgets the session restart backoff with it - and with the pending changes when this
+ * domain is disabled or imported into, and once more when it is enabled back, for the
+ * delivery whose failure was being recorded while the domain went down: it is folded
+ * after {@link #disable()} forgot the count, see {@link #enable()}. A delivery which
+ * fails while the domain is shutting down, disabled or imported into is not counted: it
+ * is not warned about either, see {@link #sessionHasAnOwner()}.
+ */
+ private final AtomicInteger foldedReplayRetryWarnings = new AtomicInteger();
/**
* The result codes conflict resolution knows how to solve. The result code the server
* puts on an internal error is configurable, and every one of these reports a failure -
@@ -2371,7 +2406,7 @@
logger.error(ERR_OPERATION_NOT_FOUND_IN_PENDING, op, curCSN);
return;
}
- resetSessionRestartBackoff();
+ resetReplayFailureTracking();
}
else
{
@@ -3598,12 +3633,13 @@
private void recordChangeResolved(CSN csn)
{
updateError(csn);
- resetSessionRestartBackoff();
+ resetReplayFailureTracking();
}
/**
* Has the change which fails next start the backoff between the session restarts over,
- * if this replica is not failing any change anymore.
+ * and the warning about it count the deliveries it stands for from zero, if this
+ * replica is not failing any change anymore.
* <p>
* A change made it and nothing is failing anymore, so the backend is serving again and
* the session is not being restarted in a row. While something is still failing, a
@@ -3611,12 +3647,19 @@
* applied here fails alone, among changes which replay perfectly well, and letting
* those reset the wait would have this domain tear its session down every second for as
* long as that one change takes to be given up on.
+ * <p>
+ * The deliveries which were folded into no warning go with the backoff rather than into
+ * the next warning: a line logged when this domain fails again - a day later, over
+ * another change - would have them read as deliveries of that failure. How long the
+ * warning is not logged again is deliberately left alone, so that a backend which fails
+ * and recovers in turn is not one warning per failure again.
*/
- private void resetSessionRestartBackoff()
+ private void resetReplayFailureTracking()
{
if (!remotePendingChanges.hasFailingChanges())
{
consecutiveSessionRestarts.set(0);
+ foldedReplayRetryWarnings.set(0);
}
}
@@ -3630,6 +3673,15 @@
* wait would have a replica which gives up on a change now and then ask for every
* change of an outage as fast as the replication server can send them, which is what
* {@link #consecutiveSessionRestarts} is there to prevent.
+ * <p>
+ * The deliveries folded into no warning are not: a warning logged over the next failure
+ * - a day later, on the default budget - would otherwise read as counting them. What is
+ * forgotten is a count, of deliveries which each had their trace line: the line which
+ * says the change is being skipped reports every delivery of this change when its
+ * budget was spent, and nothing reports the folded deliveries of the changes which were
+ * replayed meanwhile. They are forgotten when this was the last change failing, as they
+ * are when a change is replayed: while another change is still failing, they are
+ * deliveries of the outage the next warning is about.
*
* @param csn the CSN of the change which could not be replayed
* @param cause the message describing why it could not be replayed
@@ -3640,6 +3692,10 @@
{
numFailedReplayedUpdates.incrementAndGet();
sendUnreplayedChangeAlert(cause);
+ if (!remotePendingChanges.hasFailingChanges())
+ {
+ foldedReplayRetryWarnings.set(0);
+ }
}
// Otherwise the change is not listed as pending anymore - the domain was disabled
// while it was being replayed - so it has not been skipped: the replication server
@@ -3679,7 +3735,89 @@
@VisibleForTesting
public void resetUnreplayedChangeAlertThrottle()
{
- lastUnreplayedChangeAlertTime.set(UNREPLAYED_CHANGE_ALERT_NEVER_SENT);
+ lastUnreplayedChangeAlertTime.set(REPLAY_FAILURE_NEVER_REPORTED);
+ }
+
+ /**
+ * Warns that this replica could not replay a change and is asking for it again.
+ * <p>
+ * The warning is not logged again for {@link #REPLAY_RETRY_WARNING_INTERVAL_IN_MS},
+ * for the reason the alert above is not sent again and for one more: a change which
+ * keeps failing is delivered again every {@link #MAX_REPLAY_RETRY_DELAY_IN_MS} at the
+ * slowest, so one line per delivery is the same warning every ten seconds, for as long
+ * as the give-up budget of the change lasts - a budget the administrator sets, and
+ * which can be unlimited (issue #942).
+ * <p>
+ * The deliveries which are not warned about are counted, so that the line which is
+ * logged next says how many of them it stands for, and they are traced for whoever
+ * turns the replication debug logging on.
+ *
+ * @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
+ */
+ private void logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure failure)
+ {
+ final long now = monotonicNowInMs();
+ final long lastLogged = lastReplayRetryWarningTime.get();
+ if (now - lastLogged >= replayRetryWarningIntervalInMs
+ && lastReplayRetryWarningTime.compareAndSet(lastLogged, now))
+ {
+ logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts(),
+ failure.getFailingForMs(), foldedReplayRetryWarnings.getAndSet(0));
+ }
+ else
+ {
+ /*
+ * A failure which loses the race against the thread which is logging right now is
+ * counted for the line being written when its increment lands before that thread
+ * takes the count, and for the next line otherwise. It is counted once either way:
+ * what the count says is how many deliveries went without a warning of their own,
+ * and which of two lines a minute apart says it says nothing else.
+ */
+ foldedReplayRetryWarnings.incrementAndGet();
+ logger.trace("Could not replay change %s in domain %s: delivery %d, failing for %d ms",
+ csn, getBaseDN(), failure.getAttempts(), failure.getFailingForMs());
+ }
+ }
+
+ /**
+ * Lets the next change whose replay fails be warned about straight away.
+ * <p>
+ * Only there for the tests which check the warning: they must not be at the mercy of
+ * the warning another test logged less than
+ * {@link #REPLAY_RETRY_WARNING_INTERVAL_IN_MS} ago. The deliveries folded into no
+ * warning are left alone: when they are forgotten is this domain's to decide, and the
+ * tests check that it does.
+ */
+ @VisibleForTesting
+ public void resetReplayRetryWarningThrottle()
+ {
+ lastReplayRetryWarningTime.set(REPLAY_FAILURE_NEVER_REPORTED);
+ }
+
+ /**
+ * Returns how long this domain does not warn again about a change it asks for again.
+ *
+ * @return the interval in milliseconds
+ */
+ @VisibleForTesting
+ public long getReplayRetryWarningInterval()
+ {
+ return replayRetryWarningIntervalInMs;
+ }
+
+ /**
+ * Sets how long this domain does not warn again about a change it asks for again.
+ * <p>
+ * Only there for the tests which check the warning: they can not wait out
+ * {@link #REPLAY_RETRY_WARNING_INTERVAL_IN_MS} between two of them.
+ *
+ * @param intervalInMs the interval in milliseconds
+ */
+ @VisibleForTesting
+ public void setReplayRetryWarningInterval(long intervalInMs)
+ {
+ replayRetryWarningIntervalInMs = intervalInMs;
}
/**
@@ -3801,9 +3939,12 @@
* Not on the road out of a JVM which has run out of memory: building this line asks
* it for the memory it has just refused, and the ack of the delivery already says
* that the change was not applied. The constant that ack carries exists for the same
- * reason.
+ * reason. The throttle is left alone as well - the trace line a folded delivery is
+ * written to asks for that memory too, and this delivery is not one which goes
+ * unlogged: the error ends the replay thread, and the uncaught exception handler of
+ * DirectoryThread writes the line and raises the alert for it.
*/
- logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts());
+ logReplayRetryWarning(csn, failure);
}
/*
* This change is not owned by anyone anymore, so the session has to be restarted for
@@ -4121,8 +4262,9 @@
}
/**
- * Returns how many times in a row the session was restarted without a change being
- * replayed in between: the count the backoff of the next restart is computed from.
+ * Returns how many times in a row the session was restarted while this replica kept
+ * failing - see {@link #consecutiveSessionRestarts} for what does and does not reset it:
+ * the count the backoff of the next restart is computed from.
* <p>
* Only there for the tests, which read it to see a restart reach its backoff rather than
* wait a delay out and hope it began: the count is bumped on the way into the wait, so
@@ -4917,10 +5059,15 @@
/*
* The recovery from a failed replay is over as well: the change it was asking for
* is gone with the pending changes, so a leftover request would have a replay thread
- * stop and start the session once for a delivery which can not come.
+ * stop and start the session once for a delivery which can not come. The deliveries
+ * folded into no warning go with it, or the first warning over the data loaded back
+ * would read as counting the deliveries of a change which is not listed anymore. A
+ * delivery whose failure is being recorded right now is folded after this, and is
+ * forgotten by enable() instead.
*/
sessionRestarts.clear();
consecutiveSessionRestarts.set(0);
+ foldedReplayRetryWarnings.set(0);
}
// Woken outside the lock it does not hold: a restart which is waiting has nothing
// left to wait for, the session it would start being one this domain is not serving.
@@ -5085,8 +5232,20 @@
* change which is gone with the pending changes. Every request standing here is that
* one: the domain has been disabled since anything could ask, and the session started
* below asks for everything the ServerState loaded below does not cover.
+ *
+ * The deliveries folded into no warning are forgotten here for the same window: a
+ * thread which is recording a failed replay reads the flag, then folds the delivery
+ * - recoverFromReplayFailure() logs before it asks - then asks, and runs what it
+ * asked for itself, which restartSession() turns down on a domain which owns its
+ * session. So on that road it is the fold rather than the request which outlives
+ * disable()'s clear, and the first warning over the data loaded back would count a
+ * delivery of a change which went with the pending changes. What reads the count is
+ * a warning, and none is logged between disable() and here short of that thread's
+ * own, so a test tells this zeroing and disable()'s apart by nothing: they stand or
+ * fall together.
*/
sessionRestarts.clear();
+ foldedReplayRetryWarnings.set(0);
try
{
loadDataState();
@@ -5597,11 +5756,13 @@
* listed a change meanwhile: the listener thread is the one running this import,
* and the replay threads gave up every attempt while the flag was set. The restart
* a replay thread may have asked for before the total update owned the session
- * goes with them: the caller starts the session again from the reloaded state.
+ * goes with them, as do the deliveries folded into no warning: the caller starts
+ * the session again from the reloaded state.
*/
remotePendingChanges.clear();
sessionRestarts.clear();
consecutiveSessionRestarts.set(0);
+ foldedReplayRetryWarnings.set(0);
importingData = false;
}
}
--
Gitblit v1.10.0