From 92d88ca699cd8090a26b92cbe46789d2b848195f Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Mon, 07 Sep 2026 09:30:36 +0000
Subject: [PATCH] [#889] Keep a change the replay could not apply out of the ServerState (#892)
---
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java | 1156 ++++++++++++++++++++++++++++++++++++++++++++++++---------
1 files changed, 970 insertions(+), 186 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 0cd82bb..7fa68c0 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
@@ -55,8 +55,10 @@
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
-import java.util.zip.DataFormatException;
+
+import net.jcip.annotations.GuardedBy;
import org.forgerock.i18n.LocalizableMessage;
import org.forgerock.i18n.LocalizedIllegalArgumentException;
@@ -79,6 +81,7 @@
import org.forgerock.opendj.server.config.meta.ReplicationDomainCfgDefn.IsolationPolicy;
import org.forgerock.opendj.server.config.server.ExternalChangelogDomainCfg;
import org.forgerock.opendj.server.config.server.ReplicationDomainCfg;
+import org.forgerock.util.annotations.VisibleForTesting;
import org.opends.server.api.AlertGenerator;
import org.opends.server.api.DirectoryThread;
import org.opends.server.api.LocalBackend;
@@ -138,7 +141,6 @@
import org.opends.server.types.DirectoryException;
import org.opends.server.types.Entry;
import org.opends.server.types.ExistingFileBehavior;
-import org.opends.server.types.LDAPException;
import org.opends.server.types.LDIFExportConfig;
import org.opends.server.types.LDIFImportConfig;
import org.opends.server.types.Modification;
@@ -310,6 +312,120 @@
new AtomicInteger();
/** The number of updates replayed successfully by the replication. */
private final AtomicInteger numReplayedPostOpCalled = new AtomicInteger();
+ /**
+ * How many times the replay of a change is attempted straight away, before the session
+ * to the replication server is restarted and the change is asked for again: a lock or
+ * a storage which is busy for a moment (OPENDJ-885) is waited out here.
+ */
+ public static final int IN_PLACE_REPLAY_ATTEMPTS = 10;
+ /**
+ * How long the replay of a change is retried before this replica gives up on it and
+ * moves on to the changes which follow it.
+ * <p>
+ * The budget is a duration rather than a number of attempts because what
+ * {@link #isServerFailure(ResultCode, ResultCode)} reports is measured in minutes: a
+ * backend which is being rebuilt, imported into or restored (OPENDJ-49) serves nothing
+ * while it works, and a handful of attempts would have this replica give up on every
+ * change of a maintenance window it only had to wait out.
+ */
+ private static final long REPLAY_GIVE_UP_DELAY_IN_MS = 300000;
+ /**
+ * How long the session is left down before the change is asked for again, multiplied
+ * by the number of attempts already made: a backend which keeps failing must not be
+ * hammered with a session restart per failed change.
+ */
+ private static final long REPLAY_RETRY_DELAY_IN_MS = 1000;
+ /**
+ * The longest the session is left down between two attempts. The wait runs on one of
+ * the replay threads, which are shared by every domain of this server, so it is kept
+ * short enough not to starve the domains which are healthy.
+ */
+ private static final long MAX_REPLAY_RETRY_DELAY_IN_MS = 10000;
+ /**
+ * How long the alert telling that this replica diverges is not sent again. A single
+ * cause - a schema which does not match, a backend which is gone - makes every change
+ * in flight unreplayable, and one alert per change would be a storm.
+ */
+ private static final long UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MS = 60000;
+ /** The number of updates this replica gave up replaying. */
+ private final AtomicInteger numFailedReplayedUpdates = new AtomicInteger();
+ /** Set while a replay thread is restarting the session after a failed replay. */
+ private final AtomicBoolean replayFailureRecovery = new AtomicBoolean();
+ /**
+ * Set when a change whose replay failed has to be delivered again, and cleared by the
+ * replay thread which restarts the session for it. A change released while the session
+ * was being restarted has to be asked for over yet another session: the delivery this
+ * one makes is turned down as a duplicate while a replay thread still owns it.
+ */
+ private final AtomicBoolean sessionRestartRequested = new AtomicBoolean();
+ /**
+ * 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.
+ */
+ private final AtomicInteger consecutiveSessionRestarts = new AtomicInteger();
+ /**
+ * How long the replay of a change is retried before this replica gives up on it. Only
+ * the tests, which can not wait out {@link #REPLAY_GIVE_UP_DELAY_IN_MS}, set another
+ * value.
+ */
+ private volatile long replayGiveUpDelayInMs = REPLAY_GIVE_UP_DELAY_IN_MS;
+ /**
+ * Serialises the session of this domain being stopped and started again: the replay
+ * thread which restarts it after a failed replay must not race the domain being
+ * disabled for an import or a restore, or it would bring a broker and a listener
+ * thread back up on a domain which is supposed to be down.
+ * <p>
+ * Holding it costs something, and knowingly: {@code enableService()} connects to the
+ * replication servers under this lock, so a shutdown, an import or a configuration
+ * change which arrives while a replay thread is bringing the session back waits for
+ * that connect - up to the configured connection timeout when the replication servers
+ * are unreachable, which is the same outage that failed the replay. Every one of those
+ * stops the session as its first act, so what they wait for is a session which is about
+ * to be stopped again. The wait between the stop and the start is deliberately left
+ * outside the lock, so the waiting is bounded by a connect rather than by the backoff.
+ */
+ private final Object serviceStateLock = new Object();
+ /**
+ * Bumped every time the session of this domain is stopped or started under
+ * {@link #serviceStateLock}. A replay thread which stopped the session only starts it
+ * back if this still is the session it stopped: a configuration change, or the end of
+ * an import, may have started another one while it was waiting for the backend to
+ * recover.
+ * <p>
+ * It does not count the sessions {@code changeConfig()} and {@code readAssuredConfig()}
+ * stop and start, which they do without knowing about it: they run under the lock, so a
+ * replay thread never observes one of theirs, but a session it stopped may well have
+ * been replaced by one of theirs while it was waiting. That is why the guard in
+ * {@link #restartSession(boolean)} reads {@code isListenerShuttingDown()} as well - a
+ * session started outside this counter leaves it untouched, and only the listener says
+ * that one is running.
+ */
+ @GuardedBy("serviceStateLock")
+ private long sessionGeneration;
+ /**
+ * 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.
+ */
+ private static final long UNREPLAYED_CHANGE_ALERT_NEVER_SENT = 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);
+ /**
+ * The result codes conflict resolution knows how to solve. The result code the server
+ * puts on an internal error is configurable and is not validated as a result code, so
+ * it could be set to one of these: it must never take a change away from
+ * {@code solveNamingConflict()}, which is the only thing which can solve them.
+ */
+ private static final Set<ResultCode> CONFLICT_RESULT_CODES = Collections.unmodifiableSet(
+ newHashSet(
+ ResultCode.NO_SUCH_OBJECT, ResultCode.ENTRY_ALREADY_EXISTS,
+ ResultCode.NOT_ALLOWED_ON_RDN, ResultCode.NOT_ALLOWED_ON_NONLEAF,
+ // solveNamingConflict(ModifyDNOperation) solves these two as well
+ ResultCode.UNWILLING_TO_PERFORM, ResultCode.OBJECTCLASS_VIOLATION));
private final PersistentServerState state;
private volatile boolean generationIdSavedStatus;
@@ -682,36 +798,49 @@
// Disable service if configuration changed
final boolean needRestart = needReconnection && allowReconnection;
- if (needRestart)
+ /*
+ * The session is stopped, the configuration it depends on is changed and the session
+ * is started again under the lock which the replay thread restarting the session after
+ * a failed replay holds too: a session brought up in the middle of this would be
+ * reading a fractional configuration which is half way through being changed. The
+ * pair has to be atomic, which the lock inside disableService()/enableService() does
+ * not make it.
+ */
+ synchronized (serviceStateLock)
{
- disableService();
- }
- // Set new configuration
- int newFractionalMode = newFractionalConfig.fractionalConfigToInt();
- fractionalConfig.setFractional(newFractionalMode !=
- FractionalConfig.NOT_FRACTIONAL);
- if (fractionalConfig.isFractional())
- {
- // Set new fractional configuration values
- fractionalConfig.setFractionalExclusive(
- newFractionalMode == FractionalConfig.EXCLUSIVE_FRACTIONAL);
- fractionalConfig.setFractionalSpecificClassesAttributes(
- newFractionalConfig.getFractionalSpecificClassesAttributes());
- fractionalConfig.setFractionalAllClassesAttributes(
- newFractionalConfig.fractionalAllClassesAttributes);
- } else
- {
- // Reset default values
- fractionalConfig.setFractionalExclusive(true);
- fractionalConfig.setFractionalSpecificClassesAttributes(
- new HashMap<String, Set<String>>());
- fractionalConfig.setFractionalAllClassesAttributes(new HashSet<String>());
- }
+ if (needRestart)
+ {
+ disableService();
+ sessionGeneration++;
+ }
+ // Set new configuration
+ int newFractionalMode = newFractionalConfig.fractionalConfigToInt();
+ fractionalConfig.setFractional(newFractionalMode !=
+ FractionalConfig.NOT_FRACTIONAL);
+ if (fractionalConfig.isFractional())
+ {
+ // Set new fractional configuration values
+ fractionalConfig.setFractionalExclusive(
+ newFractionalMode == FractionalConfig.EXCLUSIVE_FRACTIONAL);
+ fractionalConfig.setFractionalSpecificClassesAttributes(
+ newFractionalConfig.getFractionalSpecificClassesAttributes());
+ fractionalConfig.setFractionalAllClassesAttributes(
+ newFractionalConfig.fractionalAllClassesAttributes);
+ } else
+ {
+ // Reset default values
+ fractionalConfig.setFractionalExclusive(true);
+ fractionalConfig.setFractionalSpecificClassesAttributes(
+ new HashMap<String, Set<String>>());
+ fractionalConfig.setFractionalAllClassesAttributes(new HashSet<String>());
+ }
- // Reconnect if required
- if (needRestart)
- {
- enableService();
+ // Reconnect if required
+ if (needRestart)
+ {
+ enableService();
+ sessionGeneration++;
+ }
}
}
@@ -2010,6 +2139,7 @@
logger.error(ERR_OPERATION_NOT_FOUND_IN_PENDING, op, curCSN);
return;
}
+ resetSessionRestartBackoff();
}
else
{
@@ -2249,8 +2379,15 @@
.deregisterLocalBackendInitializationListener(this);
DirectoryServer.deregisterShutdownListener(this);
- // stop the ReplicationDomain
- disableService();
+ // stop the ReplicationDomain, under the lock which the replay thread restarting the
+ // session after a failed replay holds: it must not bring a broker and a listener
+ // thread back up on a domain whose alert generator, flush thread and RSUpdater have
+ // just been taken away.
+ synchronized (serviceStateLock)
+ {
+ disableService();
+ sessionGeneration++;
+ }
}
// wait for completion of the ServerStateFlush thread.
@@ -2268,11 +2405,33 @@
/**
* Marks the specified message as the one currently processed by a replay thread.
+ *
* @param msg the message being processed
+ * @return {@code false} if the change is not pending anymore, which happens when the
+ * session was restarted after a failed replay while this message was waiting
+ * in the replay queue: it is sent again, so this copy must not be replayed
*/
- void markInProgress(LDAPUpdateMsg msg)
+ boolean markInProgress(LDAPUpdateMsg msg)
{
- remotePendingChanges.markInProgress(msg);
+ if (remotePendingChanges.markInProgress(msg))
+ {
+ return true;
+ }
+ /*
+ * This delivery is not replayed, but it was taken off the replay queue all the same:
+ * count it as processed, or replication-processed-updates would drift away from what
+ * the session delivered - that attribute counts the deliveries this replica took off
+ * the session, not the ones which reached the replay. No ack is owed for it either
+ * way, though not always for the same reason. When the change was taken over, this
+ * copy came over a session which has been torn down since and the delivery which took
+ * over from it carries the ack. When a disabled domain forgot the change, no delivery
+ * takes over and none is acked: the session that one came over is gone as well, so an
+ * ack would reach nobody, and a server waiting on an assured write times it out the
+ * way it does for every change in flight when a domain is taken out of the topology
+ * for an import or a restore.
+ */
+ incProcessedUpdates();
+ return false;
}
/**
@@ -2280,10 +2439,10 @@
*
* @param msg
* The UpdateMsg to be replayed.
- * @param shutdown
- * whether the server initiated shutdown
+ * @param replayThreadShutdown
+ * whether the replay thread was asked to stop
*/
- void replay(LDAPUpdateMsg msg, AtomicBoolean shutdown)
+ void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown)
{
// Try replay the operation, then flush (replaying) any pending operation
// whose dependency has been replayed until no more left.
@@ -2291,6 +2450,8 @@
{
Operation op = null; // the last operation on which replay was attempted
boolean dependency = false;
+ boolean replayFailed = false;
+ boolean replayAbandoned = false;
String replayErrorMsg = null;
CSN csn = null;
try
@@ -2301,16 +2462,62 @@
// "op" is already initialized to the next Operation because of the
// error handling paths.
Operation nextOp = op = msg.createOperation(conn);
+ /*
+ * The code this server puts on an internal error is a configuration knob: it is
+ * read once here so that every attempt of this delivery, and the verdict which
+ * follows them, are judged against the same one. Read inside the try - the ack of
+ * a delivery is published in the finally below, whatever the delivery ran into -
+ * and once the operation is built, so that a failure of the read is a failure of
+ * an attempt rather than a message no operation could be built from.
+ */
+ final ResultCode serverErrorResultCode =
+ getServerContext().getCoreConfigManager().getServerErrorResultCode();
dependency = remotePendingChanges.checkDependencies(op, msg);
boolean replayDone = false;
- int retryCount = 10;
+ boolean firstAttempt = true;
+ int retryCount = IN_PLACE_REPLAY_ATTEMPTS;
while (!dependency && !replayDone && retryCount-- > 0)
{
- if (shutdown.get())
+ if (replayThreadShutdown.get() || shutdown.get() || disabled)
{
- // shutdown initiated, let's leave
- return;
+ /*
+ * Either this replay thread or this domain is going away, or the domain is
+ * being imported into or restored, so let's leave. The change was never
+ * applied, so the ack says so - an assured write must not be told that a
+ * change this replica is asking for again is in the data here - and the change
+ * is given back to the replication server rather than left listed as being
+ * replayed by a thread which is gone. Handing it back stops and starts the
+ * session, so it waits until the ack has been published on the session this
+ * delivery came over.
+ *
+ * A disabled domain saved its ServerState, cleared it from memory and forgot
+ * its pending changes: a thread which kept applying changes into the backend
+ * being imported into would have every one of its commits fail on a map which
+ * is empty, one ERR_OPERATION_NOT_FOUND_IN_PENDING per change in flight, and
+ * would be writing into a backend the import owns. abandonReplay() knows there
+ * is nothing left to hand back in that case.
+ */
+ replayErrorMsg =
+ NOTE_REPLAY_ABANDONED_CHANGE.get(msg.getCSN(), getBaseDN()).toString();
+ replayAbandoned = true;
+ replayDone = true;
+ break;
}
+ if (!firstAttempt)
+ {
+ /*
+ * Every attempt runs an operation of its own. An Operation which already ran
+ * carries the request controls and the access log items of that run, so
+ * re-running the same one stacks one ManageDsaIT control - and one access log
+ * record - per attempt. It also picks up the new state of the UpdateMsg when
+ * conflict resolution rewrote it.
+ * Note: When msg is a DeleteMsg, the DeleteOperation is properly created
+ * with subtreeDelete request control when needed.
+ */
+ nextOp = msg.createOperation(conn);
+ }
+ firstAttempt = false;
+
// Try replay the operation
op = nextOp;
op.setInternalOperation(true);
@@ -2343,7 +2550,9 @@
// was a no-op. For example, an add which has already been
// replayed, or a modify DN operation on an entry which has been
// renamed by a more recent modify DN.
+ // The change is in the data: push it to the serverState.
replayDone = true;
+ recordChangeResolved(csn);
}
else if (result == ResultCode.BUSY)
{
@@ -2354,64 +2563,102 @@
Thread.yield();
continue;
}
- else if (result == ResultCode.UNAVAILABLE)
+ else if (isServerFailure(result, serverErrorResultCode))
{
/*
* It can happen when a rebuild is performed or the backend is
- * offline (OPENDJ-49). Give the server another chance to process
- * this operation after some time.
+ * offline (OPENDJ-49), or when the storage failed to serve the
+ * operation. Give the server another chance to process this
+ * operation after some time.
*/
Thread.sleep(50);
continue;
}
- else if (op instanceof ModifyOperation)
- {
- ModifyOperation castOp = (ModifyOperation) op;
- dependency = remotePendingChanges.checkDependencies(castOp);
- ModifyMsg modifyMsg = (ModifyMsg) msg;
- replayDone = !dependency && solveNamingConflict(castOp, modifyMsg);
- }
- else if (op instanceof DeleteOperation)
- {
- DeleteOperation castOp = (DeleteOperation) op;
- dependency = remotePendingChanges.checkDependencies(castOp);
- replayDone = !dependency && solveNamingConflict(castOp, msg);
- }
- else if (op instanceof AddOperation)
- {
- AddOperation castOp = (AddOperation) op;
- AddMsg addMsg = (AddMsg) msg;
- dependency = remotePendingChanges.checkDependencies(castOp);
- replayDone = !dependency && solveNamingConflict(castOp, addMsg);
- }
- else if (op instanceof ModifyDNOperation)
- {
- ModifyDNOperation castOp = (ModifyDNOperation) op;
- ModifyDNMsg modifyDNMsg = (ModifyDNMsg) msg;
- dependency = remotePendingChanges.checkDependencies(modifyDNMsg);
- replayDone = !dependency && solveNamingConflict(castOp, modifyDNMsg);
- }
else
{
- replayDone = true; // unknown type of operation ?!
- }
+ ConflictResolution resolution = ConflictResolution.NOTHING_TO_DO;
+ if (op instanceof ModifyOperation)
+ {
+ ModifyOperation castOp = (ModifyOperation) op;
+ dependency = remotePendingChanges.checkDependencies(castOp);
+ ModifyMsg modifyMsg = (ModifyMsg) msg;
+ resolution = dependency ? resolution : solveNamingConflict(castOp, modifyMsg);
+ }
+ else if (op instanceof DeleteOperation)
+ {
+ DeleteOperation castOp = (DeleteOperation) op;
+ dependency = remotePendingChanges.checkDependencies(castOp);
+ resolution = dependency ? resolution : solveNamingConflict(castOp, msg);
+ }
+ else if (op instanceof AddOperation)
+ {
+ AddOperation castOp = (AddOperation) op;
+ AddMsg addMsg = (AddMsg) msg;
+ dependency = remotePendingChanges.checkDependencies(castOp);
+ resolution = dependency ? resolution : solveNamingConflict(castOp, addMsg);
+ }
+ else if (op instanceof ModifyDNOperation)
+ {
+ ModifyDNOperation castOp = (ModifyDNOperation) op;
+ ModifyDNMsg modifyDNMsg = (ModifyDNMsg) msg;
+ dependency = remotePendingChanges.checkDependencies(modifyDNMsg);
+ resolution = dependency ? resolution : solveNamingConflict(castOp, modifyDNMsg);
+ }
+ // else: unknown type of operation ?! there is nothing to replay
- if (replayDone)
- {
- // the update became a dummy update and the result
- // of the conflict resolution phase is to do nothing.
- // however we still need to push this change to the serverState
- updateError(csn);
- }
- else
- {
- /*
- * Create a new operation reflecting the new state of the UpdateMsg after conflict resolution
- * modified it and try replaying it again. Dependencies might have been replayed by now.
- * Note: When msg is a DeleteMsg, the DeleteOperation is properly
- * created with subtreeDelete request control when needed.
- */
- nextOp = msg.createOperation(conn);
+ if (!dependency)
+ {
+ switch (resolution)
+ {
+ case NOTHING_TO_DO:
+ // the update became a dummy update and the result
+ // of the conflict resolution phase is to do nothing.
+ // however we still need to push this change to the serverState
+ replayDone = true;
+ recordChangeResolved(csn);
+ break;
+
+ case FAILED:
+ if (serverErrorResultCode.equals(result))
+ {
+ /*
+ * The result code is the one this server puts on an internal error and is
+ * one conflict resolution knows how to solve, so the change was left to it
+ * rather than treated as a failure of the server: it had its chance and
+ * could not solve it, so the storage failing is what is left. Give it the
+ * in-place attempts an UNAVAILABLE gets - a storage busy for a moment must
+ * not cost a session restart - and leave the change out of the ServerState
+ * once they are spent, which the failure of the server below the loop
+ * reports and acts on, reading the result of the attempt which spent the
+ * last of them. A change which is not in the data must not advance the
+ * ServerState (issue #889).
+ */
+ Thread.sleep(50);
+ break;
+ }
+ /*
+ * The operation did not fail on a naming conflict and not on the server
+ * either: the change can not be applied on this replica. Skip it so that the
+ * replica keeps replaying the changes which follow, but report the error in
+ * the ack and tell the administrator that the data now diverge.
+ */
+ final LocalizableMessage errorMsg = ERR_ERROR_REPLAYING_OPERATION.get(
+ op, csn, result, op.getErrorMessage());
+ logger.error(errorMsg);
+ replayErrorMsg = errorMsg.toString();
+ replayDone = true;
+ skipUnreplayableChange(csn, errorMsg);
+ break;
+
+ default:
+ /*
+ * Try replaying the change again: the next attempt creates an operation
+ * reflecting the new state of the UpdateMsg after conflict resolution
+ * modified it, and dependencies might have been replayed by now.
+ */
+ break;
+ }
+ }
}
}
else
@@ -2422,47 +2669,125 @@
if (!replayDone && !dependency)
{
- // Continue with the next change but the servers could now become
- // inconsistent.
- // Let the repair tool know about this.
- final LocalizableMessage message = ERR_LOOP_REPLAYING_OPERATION.get(
- op, op.getErrorMessage());
- logger.error(message);
- numUnresolvedNamingConflicts.incrementAndGet();
- replayErrorMsg = message.toString();
- updateError(csn);
+ final ResultCode lastResult = op.getResultCode();
+ /*
+ * BUSY is a lock which could not be taken (OPENDJ-885): the in-place attempts
+ * only yield to the thread which holds it, so a lock held for a while burns
+ * every one of them in no time. It is as transient as a storage which failed,
+ * and the change is just as absent from the data. So is a change whose last
+ * attempt came back with the configured server-error-result-code, whether or
+ * not conflict resolution owns that code: it had its chance, and what is left
+ * is the storage failing to serve the operation. The result of that attempt is
+ * what decides, so that this branch reports the failure it is acting on: an
+ * attempt which ended on something conflict resolution kept rewriting is the
+ * loop below, however the attempts before it ended.
+ */
+ if (isServerFailure(lastResult, serverErrorResultCode)
+ || ResultCode.BUSY.equals(lastResult)
+ || serverErrorResultCode.equals(lastResult))
+ {
+ /*
+ * The server kept failing to apply the change, so the change is not in the data.
+ * Leave it out of the ServerState, otherwise the replication server would never
+ * send it again and this replica would silently diverge while reporting itself
+ * up to date.
+ */
+ final LocalizableMessage message = ERR_ERROR_REPLAYING_OPERATION.get(
+ op, csn, lastResult, op.getErrorMessage());
+ logger.error(message);
+ replayErrorMsg = message.toString();
+ replayFailed = true;
+ }
+ else
+ {
+ // Conflict resolution kept rewriting an operation which kept failing.
+ // Continue with the next change but the servers could now become inconsistent.
+ // Let the repair tool know about this.
+ final LocalizableMessage message = ERR_LOOP_REPLAYING_OPERATION.get(
+ op, op.getErrorMessage());
+ logger.error(message);
+ numUnresolvedNamingConflicts.incrementAndGet();
+ replayErrorMsg = message.toString();
+ skipUnreplayableChange(csn, message);
+ }
}
- } catch (DecodeException | LDAPException | DataFormatException e)
- {
- replayErrorMsg = logDecodingOperationError(msg, e);
} catch (Exception e)
{
- if (csn != null)
+ if (op == null)
{
/*
- * An Exception happened during the replay process.
- * Continue with the next change but the servers will now start
- * to be inconsistent.
- * Let the repair tool know about this.
+ * No operation could be built from this message: createOperation() threw,
+ * whether it said so with a decoding exception or with an unchecked one. There
+ * is nothing to retry and no delivery which would build one any better.
+ *
+ * The decoding exceptions are caught here rather than in a catch of their
+ * own because such a catch would span the whole replay: addConflict(), which
+ * solveNamingConflict() calls once the operation has run, declares one, and a
+ * change whose operation ran must never be given up on where it failed.
*/
- LocalizableMessage message =
- ERR_EXCEPTION_REPLAYING_OPERATION.get(
- stackTraceToSingleLineString(e), op);
+ replayErrorMsg = giveUpOnUndecodableChange(msg, e);
+ }
+ else
+ {
+ /*
+ * An Exception happened during the replay process: the change is not in the
+ * data, so it must not be recorded as replayed.
+ * Let the repair tool know about this.
+ *
+ * The operation was built, so whatever failed is a failure of this attempt
+ * rather than a verdict on every delivery of the change - including a failure
+ * before the CSN of the operation was read, such as the entry DN of a
+ * ModifyMsg which does not parse leaving getEntryDN() null. Giving up on it
+ * where it is reported would record a change which never reached the backend
+ * in the ServerState, which is issue #889 by another route; it is left out of
+ * the ServerState and asked for again instead, and the give-up budget bounds
+ * how long this replica keeps asking.
+ */
+ final LocalizableMessage message =
+ ERR_EXCEPTION_REPLAYING_OPERATION.get(op, stackTraceToSingleLineString(e));
logger.error(message);
replayErrorMsg = message.toString();
- updateError(csn);
- } else
- {
- replayErrorMsg = logDecodingOperationError(msg, e);
+ replayFailed = true;
}
} finally
{
if (!dependency)
{
+ /*
+ * The ack is per delivery, and it says what this delivery did: a change which
+ * failed is not in the data now, whether or not the delivery which follows
+ * manages to apply it. Holding the ack back until the change is resolved would
+ * not tell the truth any better - the session it came over is about to be torn
+ * down, so nothing would reach the server which is waiting for it, and an
+ * assured write would wait out its timeout rather than be told what happened.
+ */
processUpdateDone(msg, replayErrorMsg);
}
}
+ if (replayAbandoned)
+ {
+ /*
+ * The ack has been published, so the change can be handed back now and the session
+ * restarted for it: this thread is on its way out either way.
+ */
+ abandonReplay(msg.getCSN());
+ return;
+ }
+
+ /*
+ * The CSN of the change is read off the message rather than off the operation: the
+ * two are the same - the operation carries the CSN of the message it was built from
+ * - but a failure which happened before the operation was built, or before its CSN
+ * was read, has a change to ask for again all the same.
+ */
+ if (replayFailed && recoverFromReplayFailure(msg.getCSN(), replayThreadShutdown))
+ {
+ // The ack has been published and the change, still owned by the replication
+ // server, is being delivered again: there is nothing left to replay here.
+ return;
+ }
+
// Now replay any pending update that had a dependency and whose
// dependency has been replayed, do that until no more updates of that
// type left...
@@ -2470,11 +2795,37 @@
} while (msg != null);
}
- private String logDecodingOperationError(LDAPUpdateMsg msg, Exception e)
+ /**
+ * Reports a message this replica can not turn into an operation, and gives up on the
+ * change it carries.
+ * <p>
+ * There is no operation to retry and no delivery which would decode any better, so the
+ * change is skipped rather than left out of the ServerState: a change which stays
+ * listed and uncommitted is the barrier which holds this domain's ServerState - and
+ * every change which follows it, from every master - back for good, since nothing asks
+ * for it again and the delivery which would is turned down while a replay thread still
+ * owns it. Skipping it says out loud what a wedged domain would only have implied: this
+ * replica has diverged and must be reinitialized.
+ * <p>
+ * Only a message which no operation could be built from comes here, and it is the
+ * {@code op == null} of its single caller which says so rather than the type of the
+ * exception: a decoding exception is declared past the point where the operation ran
+ * as well - by addConflict() - so a catch which read the type would give up on a
+ * change the backend may well have applied. A failure of the replay of an operation
+ * which was built, whenever it happens, keeps its change out of the ServerState and
+ * has it delivered again instead: that one is a failure of an attempt, not of every
+ * delivery of the change.
+ *
+ * @param msg the message which could not be decoded
+ * @param e the failure to decode it
+ * @return the error to report in the ack of this delivery
+ */
+ private String giveUpOnUndecodableChange(LDAPUpdateMsg msg, Exception e)
{
LocalizableMessage message =
ERR_EXCEPTION_DECODING_OPERATION.get(msg + " " + stackTraceToSingleLineString(e));
logger.error(message);
+ skipUnreplayableChange(msg.getCSN(), message);
return message.toString();
}
@@ -2485,12 +2836,16 @@
* called when error or Exceptions happen during the operation replay.
*
* @param csn the CSN of the operation with error.
+ * @return {@code false} if the change was not listed in the pending changes anymore,
+ * so that it has not been recorded as replayed: the replication server sends
+ * it again.
*/
- private void updateError(CSN csn)
+ private boolean updateError(CSN csn)
{
try
{
remotePendingChanges.commit(csn);
+ return true;
}
catch (NoSuchElementException e)
{
@@ -2502,10 +2857,389 @@
"LDAPReplicationDomain.updateError: Unable to find remote "
+ "pending change for CSN %s", csn);
}
+ return false;
}
}
/**
+ * Returns whether the provided result code reports a failure of this server rather
+ * than a change which can not be applied: the backend being offline or rebuilt
+ * (OPENDJ-49), or the storage failing to serve the operation.
+ *
+ * @param result the result code of a replayed operation
+ * @param serverErrorResultCode the result code this server puts on an internal error
+ * @return {@code true} if the operation failed on the server itself
+ */
+ private static boolean isServerFailure(ResultCode result, ResultCode serverErrorResultCode)
+ {
+ /*
+ * The result code the server puts on an internal error is configurable and is not
+ * validated as a result code, so it may well be one conflict resolution knows how to
+ * solve: such a setting must not take a change away from solveNamingConflict(), which
+ * is the only thing which can solve them. A change it could not solve either is a
+ * failure of the server all the same, which replay() acts on once conflict resolution
+ * has reported it.
+ */
+ return ResultCode.UNAVAILABLE.equals(result)
+ || (serverErrorResultCode.equals(result) && !CONFLICT_RESULT_CODES.contains(result));
+ }
+
+ /**
+ * Returns a time which only ever moves forward, in milliseconds.
+ * <p>
+ * How long a change has been failing and how long ago the last alert was sent are
+ * durations rather than dates: the wall clock stepping backwards must not have this
+ * replica retry a change for good, and stepping forwards must not have it give up on a
+ * change it had only just started to retry.
+ *
+ * @return the number of milliseconds since an arbitrary origin
+ */
+ private static long monotonicNowInMs()
+ {
+ return TimeUnit.NANOSECONDS.toMillis(System.nanoTime());
+ }
+
+ /**
+ * Records a change which conflict resolution found nothing left to do for: the change
+ * is in the data, so it is recorded as replayed and the failures it went through - and
+ * the session restarts they caused - are history.
+ *
+ * @param csn the CSN of the change
+ */
+ private void recordChangeResolved(CSN csn)
+ {
+ updateError(csn);
+ resetSessionRestartBackoff();
+ }
+
+ /**
+ * Has the change which fails next start the backoff between the session restarts over,
+ * 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
+ * change which was replayed says nothing of the kind - a change which can never be
+ * 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.
+ */
+ private void resetSessionRestartBackoff()
+ {
+ if (!remotePendingChanges.hasFailingChanges())
+ {
+ consecutiveSessionRestarts.set(0);
+ }
+ }
+
+ /**
+ * Records a change which could not be replayed as replayed anyway, so that this replica
+ * keeps replaying the changes which follow it, and warns that the data now diverge.
+ * <p>
+ * The backoff between the session restarts is deliberately left alone: giving up on a
+ * change is not a change being replayed, and whatever made this one unreplayable is
+ * still failing the ones which are in flight with it. Restarting it from its shortest
+ * 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.
+ *
+ * @param csn the CSN of the change which could not be replayed
+ * @param cause the message describing why it could not be replayed
+ */
+ private void skipUnreplayableChange(CSN csn, LocalizableMessage cause)
+ {
+ if (updateError(csn))
+ {
+ numFailedReplayedUpdates.incrementAndGet();
+ sendUnreplayedChangeAlert(cause);
+ }
+ // 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
+ // sends it again and this replica gives up on it then.
+ }
+
+ /**
+ * Tells the administrator that this replica gave up on a change and now diverges from
+ * the rest of the topology.
+ * <p>
+ * Whatever makes a change unreplayable - a schema which does not match, a backend
+ * which is gone - makes every change in flight unreplayable too, so the alert is not
+ * sent again for {@link #UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MS}: each skipped change
+ * is logged, the alert is there to have the administrator look at the log.
+ *
+ * @param cause the message describing why the change could not be replayed
+ */
+ private void sendUnreplayedChangeAlert(LocalizableMessage cause)
+ {
+ final long now = monotonicNowInMs();
+ final long lastSent = lastUnreplayedChangeAlertTime.get();
+ if (now - lastSent >= UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MS
+ && lastUnreplayedChangeAlertTime.compareAndSet(lastSent, now))
+ {
+ DirectoryServer.sendAlertNotification(
+ this, ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE, cause);
+ }
+ }
+
+ /**
+ * Lets the next change this replica gives up on raise its alert straight away.
+ * <p>
+ * Only there for the tests which check the alert: they must not be at the mercy of the
+ * alert another test raised less than
+ * {@link #UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MS} ago.
+ */
+ @VisibleForTesting
+ public void resetUnreplayedChangeAlertThrottle()
+ {
+ lastUnreplayedChangeAlertTime.set(UNREPLAYED_CHANGE_ALERT_NEVER_SENT);
+ }
+
+ /**
+ * Recovers from a change which could not be replayed.
+ * <p>
+ * The change has deliberately been left out of the ServerState, so the replication
+ * server still owns it: restart the session so that it is sent again and replayed on
+ * a backend which has hopefully recovered in the meantime. Give up once its replay has
+ * been failing for {@link #REPLAY_GIVE_UP_DELAY_IN_MS} and record it as replayed, so
+ * that a change which can never be applied here does not stop this replica for good:
+ * the administrator is told that this replica has diverged and must be reinitialized.
+ *
+ * @param csn
+ * the CSN of the change which could not be replayed
+ * @param replayThreadShutdown
+ * whether the replay thread was asked to stop
+ * @return {@code true} when the caller must stop replaying because the session is
+ * being restarted or is going away, {@code false} when it may carry on with
+ * the changes which follow
+ */
+ private boolean recoverFromReplayFailure(CSN csn, AtomicBoolean replayThreadShutdown)
+ {
+ /*
+ * The failure is recorded, and the change given up on, while this thread still owns
+ * it: a change which is listed, uncommitted and unowned is what putRemoteUpdate()
+ * takes over, so releasing it before the decision is made would let another delivery
+ * be replayed by another thread while this one goes on to record the change as
+ * skipped.
+ */
+ final long now = monotonicNowInMs();
+ final RemotePendingChanges.ReplayFailure failure =
+ remotePendingChanges.recordReplayFailure(csn, now);
+ if (failure == null)
+ {
+ /*
+ * There is no uncommitted change left to give up on or to ask for again: the domain
+ * was disabled while this one was being replayed - its ServerState was saved and is
+ * read again from the backend when the domain is enabled back - or the change had
+ * already been recorded when this failure was reported. Carry on with the changes
+ * which were waiting for it; a domain on its way down lists none.
+ */
+ return false;
+ }
+ if (failure.getFailingForMs() >= replayGiveUpDelayInMs)
+ {
+ final LocalizableMessage message = ERR_REPLAY_SKIPPING_CHANGE.get(
+ csn, getBaseDN(), failure.getFailingForMs(), failure.getAttempts());
+ logger.error(message);
+ skipUnreplayableChange(csn, message);
+ return false;
+ }
+
+ /*
+ * The change stays listed as pending and uncommitted - it holds the ServerState back
+ * so that the replication server sends it again - but this thread does not own it
+ * anymore: the next delivery is the one which is replayed, and the copy which may
+ * still wait in the shared replay queue is dropped when a replay thread takes it out,
+ * because markInProgress() only accepts the delivery which is listed as pending.
+ */
+ remotePendingChanges.replayFailed(csn);
+
+ if (shutdown.get() || disabled)
+ {
+ /*
+ * This whole domain is going away or is being imported into: there is no session of
+ * this thread's to restart. Restarting the one which is being stopped would leave a
+ * broker and a listener thread behind on a domain whose alert generator, flush
+ * thread and RSUpdater are already gone.
+ */
+ return true;
+ }
+
+ logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts());
+ /*
+ * This change is not owned by anyone anymore, so the session has to be restarted for
+ * the replication server to deliver it again. Ask for the restart before trying to
+ * run it: a restart which is already under way may have started before this change
+ * was released, and the delivery it asked for would then have been turned down as a
+ * duplicate of a change a replay thread still owned.
+ */
+ sessionRestartRequested.set(true);
+ /*
+ * A replay thread which is stopping - the number of them is being changed - restarts
+ * the session all the same: nothing else would ask for the change it just released,
+ * and the ServerState would stay behind it for good. It does not sit through the
+ * backoff on its way out, though: the backend is not what is going away.
+ */
+ runRequestedSessionRestarts(!replayThreadShutdown.get());
+ return true;
+ }
+
+ /**
+ * Restarts the session as long as changes which could not be replayed are waiting to be
+ * delivered again.
+ *
+ * @param wait whether to leave the backend some time to recover between two restarts
+ */
+ private void runRequestedSessionRestarts(boolean wait)
+ {
+ /*
+ * 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
+ * it to this one.
+ */
+ while (sessionRestartRequested.get() && replayFailureRecovery.compareAndSet(false, true))
+ {
+ try
+ {
+ while (sessionRestartRequested.getAndSet(false))
+ {
+ restartSession(wait);
+ }
+ }
+ finally
+ {
+ replayFailureRecovery.set(false);
+ }
+ }
+ }
+
+ /**
+ * Gives a change back to the replication server when this replay thread stops before it
+ * could apply it.
+ * <p>
+ * The change is not owned by anyone anymore and it was never applied, so the session is
+ * restarted for it to be delivered again: it is left out of the ServerState, and the
+ * changes which follow it are held back until it is replayed.
+ *
+ * @param csn the CSN of the change this thread was replaying
+ */
+ private void abandonReplay(CSN csn)
+ {
+ remotePendingChanges.replayFailed(csn);
+ if (shutdown.get() || disabled)
+ {
+ // The domain owns its session, and it forgets its pending changes on its way down.
+ return;
+ }
+ /*
+ * Logged here rather than where the change is abandoned: a server which is shutting
+ * down abandons every change in flight, and none of them is asked for again before it
+ * is started back - one line per change would say otherwise.
+ */
+ logger.info(NOTE_REPLAY_ABANDONED_CHANGE, csn, getBaseDN());
+ sessionRestartRequested.set(true);
+ runRequestedSessionRestarts(false);
+ }
+
+ /**
+ * Stops the session to the replication server and starts it again, so that the changes
+ * this replica could not replay are delivered again.
+ */
+ private void restartSession(boolean wait)
+ {
+ final long stoppedSession;
+ synchronized (serviceStateLock)
+ {
+ if (shutdown.get() || disabled)
+ {
+ // The domain is going away or is being imported into: it owns its session.
+ return;
+ }
+ disableService();
+ stoppedSession = ++sessionGeneration;
+ }
+ if (wait)
+ {
+ /*
+ * Leave the backend some time to recover rather than ask for the change straight
+ * away: a session restart is not free for the replication server either. The wait
+ * is not held under the lock, or a domain being disabled for an import would wait
+ * it out.
+ */
+ waitBeforeSessionRestart(consecutiveSessionRestarts.incrementAndGet());
+ }
+ synchronized (serviceStateLock)
+ {
+ if (shutdown.get() || disabled
+ || sessionGeneration != stoppedSession || !isListenerShuttingDown())
+ {
+ /*
+ * The domain went away while this thread was waiting, or the session was stopped
+ * and started again by something else - a configuration change, the end of an
+ * import - in the meantime: the session this thread stopped is gone, so it has
+ * nothing left to start. The generation says a session was started under this
+ * lock; the listener says one is running, which is what a restart made outside it
+ * leaves behind.
+ */
+ return;
+ }
+ enableService();
+ sessionGeneration++;
+ }
+ }
+
+ /**
+ * Waits for a while before the session to the replication server is started again, so
+ * that a backend which keeps failing is not asked for every change it can not apply as
+ * fast as the replication server can send them.
+ *
+ * @param restarts how many times in a row the session was restarted already
+ */
+ private void waitBeforeSessionRestart(int restarts)
+ {
+ try
+ {
+ Thread.sleep(Math.min(REPLAY_RETRY_DELAY_IN_MS * restarts, MAX_REPLAY_RETRY_DELAY_IN_MS));
+ }
+ catch (InterruptedException e)
+ {
+ /*
+ * Do not wait, but do start the session again all the same: the session is down
+ * because this thread stopped it, and leaving it down would take this domain out of
+ * the topology until the server is restarted. The interrupt is left set for whoever
+ * asked this thread to stop.
+ */
+ Thread.currentThread().interrupt();
+ }
+ }
+
+ /**
+ * Returns how long the replay of a change is retried before this replica gives up on
+ * it.
+ * <p>
+ * Only there for the tests, which set another value and put this one back.
+ *
+ * @return how long a change is retried, in milliseconds
+ */
+ @VisibleForTesting
+ public long getReplayGiveUpDelay()
+ {
+ return replayGiveUpDelayInMs;
+ }
+
+ /**
+ * Sets how long the replay of a change is retried before this replica gives up on it.
+ * <p>
+ * Only there for the tests, which can not wait out the {@link
+ * #REPLAY_GIVE_UP_DELAY_IN_MS} a backend under maintenance is given.
+ *
+ * @param delayInMs how long a change is retried, in milliseconds
+ */
+ @VisibleForTesting
+ public void setReplayGiveUpDelay(long delayInMs)
+ {
+ this.replayGiveUpDelayInMs = delayInMs;
+ }
+
+ /**
* Generate a new CSN and insert it in the pending list.
*
* @param operation
@@ -2581,14 +3315,25 @@
return null;
}
+ /** Outcome of the conflict resolution attempted after a replayed operation failed. */
+ private enum ConflictResolution
+ {
+ /** The update message was adjusted: the operation must be replayed again. */
+ REPLAY_AGAIN,
+ /** The change is already reflected in the data: there is nothing left to replay. */
+ NOTHING_TO_DO,
+ /** The operation failed for a reason which is not a naming conflict. */
+ FAILED
+ }
+
/**
* Solve a conflict detected when replaying a modify operation.
*
* @param op The operation that triggered the conflict detection.
* @param msg The operation that triggered the conflict detection.
- * @return true if the process is completed, false if it must continue..
+ * @return the outcome of the conflict resolution
*/
- private boolean solveNamingConflict(ModifyOperation op, ModifyMsg msg)
+ private ConflictResolution solveNamingConflict(ModifyOperation op, ModifyMsg msg)
{
ResultCode result = op.getResultCode();
ModifyContext ctx = (ModifyContext) op.getAttachment(SYNCHROCONTEXT);
@@ -2609,14 +3354,14 @@
// replay the modify using the current dn of this entry.
msg.setDN(newDN);
numResolvedNamingConflicts.incrementAndGet();
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
else
{
// This entry does not exist anymore.
// It has probably been deleted, stop the processing of this operation
numResolvedNamingConflicts.incrementAndGet();
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
}
else if (result == ResultCode.NOT_ALLOWED_ON_RDN)
@@ -2631,7 +3376,7 @@
{
// The entry does not exist anymore.
numResolvedNamingConflicts.incrementAndGet();
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
// The modify operation is trying to delete the value that is
@@ -2657,15 +3402,13 @@
}
msg.setMods(mods);
numResolvedNamingConflicts.incrementAndGet();
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
else
{
- // The other type of errors can not be caused by naming conflicts.
- // Log a message for the repair tool.
- logger.error(ERR_ERROR_REPLAYING_OPERATION,
- op, ctx.getCSN(), result, op.getErrorMessage());
- return true;
+ // The other type of errors can not be caused by naming conflicts:
+ // the operation simply failed, replay() reports it.
+ return ConflictResolution.FAILED;
}
}
@@ -2674,9 +3417,9 @@
*
* @param op The operation that triggered the conflict detection.
* @param msg The operation that triggered the conflict detection.
- * @return true if the process is completed, false if it must continue..
+ * @return the outcome of the conflict resolution
*/
- private boolean solveNamingConflict(DeleteOperation op, LDAPUpdateMsg msg)
+ private ConflictResolution solveNamingConflict(DeleteOperation op, LDAPUpdateMsg msg)
{
ResultCode result = op.getResultCode();
DeleteContext ctx = (DeleteContext) op.getAttachment(SYNCHROCONTEXT);
@@ -2695,14 +3438,14 @@
* In any case, there is nothing more to do.
*/
numResolvedNamingConflicts.incrementAndGet();
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
else
{
// This entry has been renamed, replay the delete using its new DN.
msg.setDN(currentDN);
numResolvedNamingConflicts.incrementAndGet();
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
}
else if (result == ResultCode.NOT_ALLOWED_ON_NONLEAF)
@@ -2722,15 +3465,13 @@
numUnresolvedNamingConflicts.incrementAndGet();
}
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
else
{
- // The other type of errors can not be caused by naming conflicts.
- // Log a message for the repair tool.
- logger.error(ERR_ERROR_REPLAYING_OPERATION,
- op, ctx.getCSN(), result, op.getErrorMessage());
- return true;
+ // The other type of errors can not be caused by naming conflicts:
+ // the operation simply failed, replay() reports it.
+ return ConflictResolution.FAILED;
}
}
@@ -2739,10 +3480,10 @@
*
* @param op The operation that triggered the conflict detection.
* @param msg The operation that triggered the conflict detection.
- * @return true if the process is completed, false if it must continue.
+ * @return the outcome of the conflict resolution
* @throws Exception When the operation is not valid.
*/
-private boolean solveNamingConflict(ModifyDNOperation op, LDAPUpdateMsg msg)
+private ConflictResolution solveNamingConflict(ModifyDNOperation op, LDAPUpdateMsg msg)
throws Exception
{
ResultCode result = op.getResultCode();
@@ -2788,7 +3529,7 @@
{
markConflictEntry(op, currentDN, currentDN.parent().child(newRDN));
numUnresolvedNamingConflicts.incrementAndGet();
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
DN newDN = newSuperior.child(newRDN);
@@ -2801,7 +3542,7 @@
// The entry has been deleted, we can safely assume
// that the operation is completed.
numResolvedNamingConflicts.incrementAndGet();
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
// if the newDN and the current DN match then the operation
@@ -2810,7 +3551,7 @@
if (newDN.equals(currentDN))
{
numResolvedNamingConflicts.incrementAndGet();
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
if (result == ResultCode.NO_SUCH_OBJECT
@@ -2825,7 +3566,7 @@
modifyDnMsg.setDN(currentDN);
modifyDnMsg.setNewSuperior(newSuperior.toString());
numResolvedNamingConflicts.incrementAndGet();
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
else if (result == ResultCode.ENTRY_ALREADY_EXISTS)
{
@@ -2841,15 +3582,13 @@
modifyDnMsg.getNewRDN()));
modifyDnMsg.setNewSuperior(newSuperior.toString());
numUnresolvedNamingConflicts.incrementAndGet();
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
else
{
- // The other type of errors can not be caused by naming conflicts.
- // Log a message for the repair tool.
- logger.error(ERR_ERROR_REPLAYING_OPERATION,
- op, ctx.getCSN(), result, op.getErrorMessage());
- return true;
+ // The other type of errors can not be caused by naming conflicts:
+ // the operation simply failed, replay() reports it.
+ return ConflictResolution.FAILED;
}
}
@@ -2858,10 +3597,10 @@
*
* @param op The operation that triggered the conflict detection.
* @param msg The message that triggered the conflict detection.
- * @return true if the process is completed, false if it must continue.
+ * @return the outcome of the conflict resolution
* @throws Exception When the operation is not valid.
*/
- private boolean solveNamingConflict(AddOperation op, AddMsg msg)
+ private ConflictResolution solveNamingConflict(AddOperation op, AddMsg msg)
throws Exception
{
ResultCode result = op.getResultCode();
@@ -2884,7 +3623,7 @@
* message for the repair tool to look at this problem.
* TODO : Log the message
*/
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
DN parentDn = findEntryDN(parentUniqueId);
if (parentDn == null)
@@ -2911,7 +3650,7 @@
msg.setDN(DN.valueOf(msg.getDN().rdn() + "," + parentDn));
numResolvedNamingConflicts.incrementAndGet();
}
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
else if (result == ResultCode.ENTRY_ALREADY_EXISTS)
{
@@ -2927,7 +3666,7 @@
if (findEntryDN(entryUUID) != null)
{
// entry already exist : this is a replay
- return true;
+ return ConflictResolution.NOTHING_TO_DO;
}
else
{
@@ -2936,16 +3675,14 @@
generateConflictRDN(entryUUID, msg.getDN().toString());
msg.setDN(DN.valueOf(conflictRDN));
numUnresolvedNamingConflicts.incrementAndGet();
- return false;
+ return ConflictResolution.REPLAY_AGAIN;
}
}
else
{
- // The other type of errors can not be caused by naming conflicts.
- // log a message for the repair tool.
- logger.error(ERR_ERROR_REPLAYING_OPERATION,
- op, ctx.getCSN(), result, op.getErrorMessage());
- return true;
+ // The other type of errors can not be caused by naming conflicts:
+ // the operation simply failed, replay() reports it.
+ return ConflictResolution.FAILED;
}
}
@@ -3129,10 +3866,30 @@
*/
public void disable()
{
- state.save();
- state.clearInMemory();
- disabled = true;
- disableService(); // This will cut the session and wake up the listener
+ synchronized (serviceStateLock)
+ {
+ state.save();
+ state.clearInMemory();
+ disabled = true;
+ disableService(); // This will cut the session and wake up the listener
+ sessionGeneration++;
+ /*
+ * The ServerState this bookkeeping goes with is now gone from memory and is loaded
+ * again from the backend when the domain is enabled back, so the changes listed as
+ * pending must not outlive it: one which stayed would be discarded as a duplicate
+ * when the replication server sends it again, and nothing would ever replay it or
+ * record it in the ServerState. The listener thread is stopped first, or a change
+ * it lists after this would be the one left behind.
+ */
+ remotePendingChanges.clear();
+ /*
+ * 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.
+ */
+ sessionRestartRequested.set(false);
+ consecutiveSessionRestarts.set(0);
+ }
}
/**
@@ -3163,23 +3920,27 @@
*/
public void enable()
{
- try
+ synchronized (serviceStateLock)
{
- loadDataState();
- }
- catch (Exception e)
- {
- /* TODO should mark that replicationServer service is
- * not available, log an error and retry upon timeout
- * should we stop the modifications ?
- */
- logger.error(ERR_LOADING_GENERATION_ID, getBaseDN(), stackTraceToSingleLineString(e));
- return;
- }
+ try
+ {
+ loadDataState();
+ }
+ catch (Exception e)
+ {
+ /* TODO should mark that replicationServer service is
+ * not available, log an error and retry upon timeout
+ * should we stop the modifications ?
+ */
+ logger.error(ERR_LOADING_GENERATION_ID, getBaseDN(), stackTraceToSingleLineString(e));
+ return;
+ }
- enableService();
+ enableService();
+ sessionGeneration++;
- disabled = false;
+ disabled = false;
+ }
}
/**
@@ -3795,11 +4556,20 @@
ReplicationDomainCfg configuration)
{
this.config = configuration;
- changeConfig(configuration);
+ /*
+ * Each of these stops and starts the session when what it changes calls for it, and
+ * the configuration they change is read as the session comes up: hold the lock the
+ * replay thread restarting the session after a failed replay takes, so that none of
+ * them is interleaved with a session it did not start itself.
+ */
+ synchronized (serviceStateLock)
+ {
+ changeConfig(configuration);
- // Read assured + fractional configuration and each time reconnect if needed
- readAssuredConfig(configuration, true);
- readFractionalConfig(configuration, true);
+ // Read assured + fractional configuration and each time reconnect if needed
+ readAssuredConfig(configuration, true);
+ readFractionalConfig(configuration, true);
+ }
solveConflictFlag = isSolveConflict(configuration);
@@ -3847,6 +4617,8 @@
alerts.put(ALERT_TYPE_REPLICATION_UNRESOLVED_CONFLICT,
ALERT_DESCRIPTION_REPLICATION_UNRESOLVED_CONFLICT);
+ alerts.put(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE,
+ ALERT_DESCRIPTION_REPLICATION_UNREPLAYED_CHANGE);
return alerts;
}
@@ -4254,10 +5026,20 @@
if (!remotePendingChanges.putRemoteUpdate(msg))
{
/*
- * Already received this change so ignore it. This may happen if there
- * are uncommitted changes in the queue and session failover occurs
- * causing a recovery of all changes since the current committed server
- * state. See OPENDJ-1115.
+ * A replay thread already owns this change, so this delivery is a duplicate:
+ * ignore it. This happens when there are uncommitted changes in the queue and
+ * session failover occurs causing a recovery of all changes since the current
+ * committed server state. See OPENDJ-1115.
+ *
+ * The copy which is already listed owns the change: it is the one which records
+ * it in the ServerState once it really has been replayed. Report this delivery
+ * as done - the window and the ack are per delivery - but as handled
+ * asynchronously, so that the listener does not push the CSN to the ServerState
+ * over a change which is still being replayed or is failing (issue #889).
+ *
+ * A change whose replay failed is not owned by anyone anymore, so this is not the
+ * path it takes: putRemoteUpdate() takes this delivery over from the one which
+ * failed and has it replayed again.
*/
if (logger.isTraceEnabled())
{
@@ -4265,7 +5047,8 @@
"LDAPReplicationDomain.processUpdate: ignoring "
+ "duplicate change %s", msg.getCSN());
}
- return true;
+ processUpdateDone(msg, null);
+ return false;
}
// Put update message into the replay queue
@@ -4301,6 +5084,7 @@
{
attributes.add("pending-updates", pendingChanges.size());
attributes.add("replayed-updates-ok", numReplayedPostOpCalled);
+ attributes.add("replayed-updates-failed", numFailedReplayedUpdates);
attributes.add("resolved-modify-conflicts", numResolvedModifyConflicts);
attributes.add("resolved-naming-conflicts", numResolvedNamingConflicts);
attributes.add("unresolved-naming-conflicts", numUnresolvedNamingConflicts);
--
Gitblit v1.10.0