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