mirror of https://github.com/OpenIdentityPlatform/OpenDJ.git

Valery Kharseko
17 hours ago 661dc06886df4738b9add2206a39d74d7dedf476
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
@@ -57,6 +57,7 @@
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import net.jcip.annotations.GuardedBy;
@@ -143,6 +144,7 @@
import org.opends.server.types.ExistingFileBehavior;
import org.opends.server.types.LDIFExportConfig;
import org.opends.server.types.LDIFImportConfig;
import org.opends.server.types.LockManager;
import org.opends.server.types.Modification;
import org.opends.server.types.Operation;
import org.opends.server.types.OperationType;
@@ -389,6 +391,75 @@
  @GuardedBy("serviceStateLock")
  private long sessionGeneration;
  /**
   * Held while a replay thread applies a change of this domain, and taken exclusively by
   * this domain on its way down.
   * <p>
   * A change which reached the backend has to be recorded in the ServerState which is saved
   * when the domain is disabled or shut down, or it ends up in the data and in no
   * ServerState: the replication server sends it again and a change which is already
   * applied is replayed a second time (issue #908). The flags which stop the replay -
   * {@link #disabled}, {@link #shutdown} - are read under this lock as well, so a domain on
   * its way down sets one of them and then takes this lock: what it waits for is the
   * changes which were already being applied, and no attempt starts after that.
   * <p>
   * What this closes is the window between an operation reaching the backend and its
   * {@code commit()}, which is the one issue #908 reports. It is not every road by which a
   * change can be in the data and in no saved ServerState: {@code commit()} only advances
   * the state over the changes which are committed from the head of the pending list, so a
   * change applied while an older one is still to be replayed is forgotten by the
   * {@code clear()} on the way down and sent again - the barrier of issue #889 seen from
   * the other side, and a thing to fix where that barrier is rather than here. A change
   * this replica made itself is a third road, and the one the ServerState recovery in
   * {@code PersistentServerState.loadState()} already repairs, since it only looks for the
   * CSNs of this server.
   * <p>
   * Deliberately not the fair kind. The read lock is taken for every attempt made on the
   * backend, and a queued writer blocks the readers which come after it even without
   * fairness; the readers which do barge past it are the ones which see the flag and leave
   * without applying anything. The flag is read once before the lock for that reason too -
   * the replay threads are a pool shared by every domain, and a thread which parks on the
   * lock of a domain going down is a thread no other domain gets its changes replayed by.
   * <p>
   * The lock is released at the end of every attempt, so a domain going down waits for one
   * attempt and the short backoff which follows it inside the loop, rather than for all
   * the attempts a delivery is given. The backoff between two attempts is a Thread.sleep()
   * of tens of milliseconds; the one between two session restarts, which is counted in
   * seconds, is outside every lock and stays there.
   */
  private final ReentrantReadWriteLock replayLock = new ReentrantReadWriteLock();
  private final ReentrantReadWriteLock.ReadLock replayReadLock = replayLock.readLock();
  private final ReentrantReadWriteLock.WriteLock replayWriteLock = replayLock.writeLock();
  /**
   * How long this domain waits for the replay threads which are applying one of its changes
   * before it saves its ServerState and goes down.
   * <p>
   * Derived from the ceiling the server itself puts on an operation which is waiting for an
   * entry rather than picked: {@link LockManager} gives the subtree lock and the entry lock
   * {@link LockManager#DEFAULT_LOCK_TIMEOUT} each, so a replayed change whose target is held
   * by a concurrent local operation - a client deleting the subtree above it, say - is
   * inside its attempt for twice that before it gives up with BUSY. A bound under that
   * ceiling would be spent by ordinary lock contention, and the change which is applied
   * after it would be the one this whole barrier exists to keep out of that window: no
   * import, no index rebuild and no wedged backend needed.
   * <p>
   * A ceiling rather than a guarantee: an operation also waits for the subtree lock of every
   * entry above its target, one timeout each, so a deep contended chain outlasts this. What
   * it buys is that the wait is not lost to the contention a serving backend has anyway.
   * <p>
   * It is paid in three places - held under serviceStateLock, inside BackendConfigManager's
   * write lock when a backend is being deregistered, and once per domain by a server going
   * down - which is why it is bounded at all. It is only ever spent in full by a replay
   * which is genuinely stuck: the wait ends the moment the attempt does.
   */
  private static final long REPLAY_DRAIN_TIMEOUT_IN_MS =
      2 * LockManager.DEFAULT_LOCK_TIMEOUT_UNITS.toMillis(LockManager.DEFAULT_LOCK_TIMEOUT) + 1000;
  /**
   * How long this domain waits for the replay of its changes on its way down. Only the
   * tests, which can not hold a replay thread for {@link #REPLAY_DRAIN_TIMEOUT_IN_MS}, set
   * another value.
   */
  private volatile long replayDrainTimeoutInMs = REPLAY_DRAIN_TIMEOUT_IN_MS;
  /**
   * Stands for "the alert about a change this replica gave up on was never sent". The
   * time it is compared with only moves forward from an origin which is arbitrary, so
   * zero is not far enough in the past to say it.
@@ -556,7 +627,27 @@
          Thread.currentThread().interrupt();
        }
      }
      state.save();
      /*
       * A disabled domain saved its ServerState and cleared it from memory, and an import
       * or a restore is about to replace the data: saving here would write that empty state
       * over the saved one, which is a REPLACE of ds-sync-state with no value at all - the
       * replica would come back with no ServerState rather than with the one it saved on
       * its way down. The disabled flag does not catch the total update this replica is the
       * target of: preBackendImport() sets ignoreBackendInitializationEvent, so disable() is
       * not called on that road and only the import says the data is being replaced.
       *
       * The direction is asked for rather than ieRunning(), which the save in the loop above
       * settles for: an export leaves the data and the ServerState of this domain alone, and
       * the replay of this domain keeps running for the whole of it - a remote-requested
       * export is dispatched to a thread pool for that very reason. Since the save in the
       * loop is skipped for either direction, this is the only one which persists the
       * changes replayed since the export began, and there is no repair for them afterwards:
       * checkAndUpdateServerState() only repairs the CSNs of this server.
       */
      if (!disabled && !importInProgress())
      {
        state.save();
      }
      done = true;
    }
@@ -2355,6 +2446,15 @@
  {
    if (shutdown.compareAndSet(false, true))
    {
      /*
       * Wait for the changes which are being applied before the ServerState is flushed for
       * the last time: the flush thread stopped below is the last thing which saves it, so a
       * change which reaches the backend after that save is in the data and in no
       * ServerState (issue #908). The flag was just set, so this waits for the attempts
       * which had started already and no new one begins.
       */
      awaitReplayDrained();
      final RSUpdater rsUpdater = this.rsUpdater.get();
      if (rsUpdater != null)
      {
@@ -2475,7 +2575,205 @@
        int retryCount = IN_PLACE_REPLAY_ATTEMPTS;
        while (!dependency && !replayDone && retryCount-- > 0)
        {
          if (replayThreadShutdown.get() || shutdown.get() || disabled)
          /*
           * The flag which says this domain is going down is read before the lock as well
           * as under it. The replay threads are a pool shared by every domain of this
           * server, so a thread which took a change of a domain which is going down should
           * not queue behind the wait for that domain: the changes of every other domain
           * are behind it in the same pool.
           *
           * The read under the lock is the one which decides; the one above it is a
           * scheduling optimisation for the common case and nothing more. It cannot keep
           * this thread out of the queue: the flag can be set and the writer can queue
           * between the two reads, and a reader which arrives behind a queued writer blocks
           * even on a lock which is not the fair kind.
           */
          boolean goingDown = replayThreadShutdown.get() || shutdown.get() || disabled;
          if (!goingDown)
          {
            /*
             * Every attempt made on the backend is under this lock, and so is the decision
             * to make one: a domain on its way down takes it exclusively once it has set
             * the flag read here, so a change which reaches the backend is recorded in the
             * ServerState which is saved on the way down, or is not applied at all
             * (issue #908).
             */
            replayReadLock.lock();
            try
            {
              goingDown = replayThreadShutdown.get() || shutdown.get() || disabled;
              if (!goingDown)
              {
                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);
                op.setSynchronizationOperation(true);
                // Always add the ManageDSAIT control so that updates to referrals
                // are processed locally.
                op.addRequestControl(new LDAPControl(OID_MANAGE_DSAIT_CONTROL));
                // Warning: specific processing ahead. See OPENDJ-2792
                if (op instanceof ModifyOperation)
                {
                  ModifyOperation modifyOperation = (ModifyOperation) op;
                  if (modifyOperation.getEntryDN().equals(SET_PERMISSIVE_MODIFY_FOR_DN))
                  {
                    op.addRequestControl(new LDAPControl(OID_PERMISSIVE_MODIFY_CONTROL));
                  }
                }
                csn = OperationContext.getCSN(op);
                op.run();
                ResultCode result = op.getResultCode();
                if (result != ResultCode.SUCCESS)
                {
                  if (result == ResultCode.NO_OPERATION)
                  {
                    // Pre-operation conflict resolution detected that the operation
                    // 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)
                  {
                    /*
                     * We probably could not get a lock (OPENDJ-885). Give the server
                     * another chance to process this operation immediately.
                     */
                    Thread.yield();
                    continue;
                  }
                  else if (isServerFailure(result, serverErrorResultCode))
                  {
                    /*
                     * It can happen when a rebuild is performed or the backend is
                     * 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
                  {
                    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 (!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
                {
                  replayDone = true;
                }
              }
            }
            finally
            {
              replayReadLock.unlock();
            }
          }
          if (goingDown)
          {
            /*
             * Either this replay thread or this domain is going away, or the domain is
@@ -2500,168 +2798,6 @@
            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);
          op.setSynchronizationOperation(true);
          // Always add the ManageDSAIT control so that updates to referrals
          // are processed locally.
          op.addRequestControl(new LDAPControl(OID_MANAGE_DSAIT_CONTROL));
          // Warning: specific processing ahead. See OPENDJ-2792
          if (op instanceof ModifyOperation)
          {
            ModifyOperation modifyOperation = (ModifyOperation) op;
            if (modifyOperation.getEntryDN().equals(SET_PERMISSIVE_MODIFY_FOR_DN))
            {
              op.addRequestControl(new LDAPControl(OID_PERMISSIVE_MODIFY_CONTROL));
            }
          }
          csn = OperationContext.getCSN(op);
          op.run();
          ResultCode result = op.getResultCode();
          if (result != ResultCode.SUCCESS)
          {
            if (result == ResultCode.NO_OPERATION)
            {
              // Pre-operation conflict resolution detected that the operation
              // 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)
            {
              /*
               * We probably could not get a lock (OPENDJ-885). Give the server
               * another chance to process this operation immediately.
               */
              Thread.yield();
              continue;
            }
            else if (isServerFailure(result, serverErrorResultCode))
            {
              /*
               * It can happen when a rebuild is performed or the backend is
               * 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
            {
              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 (!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
          {
            replayDone = true;
          }
        }
        if (!replayDone && !dependency)
@@ -3855,11 +3991,29 @@
  {
    synchronized (serviceStateLock)
    {
      state.save();
      state.clearInMemory();
      /*
       * The replay is stopped before the ServerState is saved, and not the other way round:
       * a change a replay thread is applying has to be either recorded in the state which
       * is about to be saved or not applied at all, or it ends up in the data and in no
       * ServerState (issue #908). The flag keeps the attempts which have not started from
       * starting - it is read under the same lock as the attempt it guards - and the wait
       * below is for the ones which had started already.
       *
       * All of it stays under serviceStateLock, so that this and enable() remain the
       * mutually exclusive pair they have always been: an enable() which ran in the middle
       * of this would clear the flag and bring a session up, and this would then go on to
       * cut that session and clear a ServerState which the flush thread - reading a flag
       * which says the domain is enabled - would write back empty. That is what bounds
       * REPLAY_DRAIN_TIMEOUT_IN_MS: the wait is held under a lock which a session restart,
       * a configuration change and the shutdown of this domain take, and it runs inside
       * BackendConfigManager's write lock when a backend is being deregistered.
       */
      disabled = true;
      disableService(); // This will cut the session and wake up the listener
      sessionGeneration++;
      awaitReplayDrained();
      state.save();
      state.clearInMemory();
      /*
       * 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
@@ -3880,6 +4034,90 @@
  }
  /**
   * Waits for the replay threads which are applying a change of this domain to be done
   * with it.
   * <p>
   * Called once {@link #disabled} or {@link #shutdown} has been set, which is what bounds
   * the wait: a replay thread reads those under {@link #replayReadLock}, the lock this
   * takes exclusively, so no attempt starts once this returns and what it waits for is the
   * attempts which were running already. The lock is released before returning for the
   * same reason - what keeps the replay out is the flag, not the lock.
   */
  private void awaitReplayDrained()
  {
    boolean drained = false;
    boolean interrupted = false;
    try
    {
      drained = replayWriteLock.tryLock(replayDrainTimeoutInMs, TimeUnit.MILLISECONDS);
    }
    catch (InterruptedException e)
    {
      /*
       * Give up waiting, and put the interrupt back rather than swallow it: whoever
       * interrupted this thread - the server going down, a thread pool taking its threads
       * away - is still waiting for it to stop, and this is not the last thing it does.
       * The cost is on the shutdown road, whose wait for the last ServerState flush is a
       * Thread.sleep() which ends on its first call once the flag is set: the flush thread
       * still runs that save, this one just stops waiting for it.
       */
      interrupted = true;
      Thread.currentThread().interrupt();
    }
    if (drained)
    {
      replayWriteLock.unlock();
      return;
    }
    /*
     * The change which is being applied may reach the backend without being recorded in the
     * ServerState which is saved next, so the replication server sends it again and it is
     * replayed a second time. Better than holding an administrative task - an import, a
     * restore, a backend being taken offline - for as long as a backend which stopped
     * answering takes to answer.
     *
     * An interrupted wait is reported as what it is: it says nothing about how long the
     * replay of this domain takes, and the timeout it never spent would have an operator
     * reading a backend which is slow into it.
     */
    if (interrupted)
    {
      logger.warn(WARN_REPLAY_DRAIN_INTERRUPTED, getBaseDN());
    }
    else
    {
      logger.warn(WARN_REPLAY_NOT_DRAINED, getBaseDN(), replayDrainTimeoutInMs);
    }
  }
  /**
   * Returns how long this domain waits for the replay threads which are applying one of
   * its changes before it saves its ServerState and goes down.
   *
   * @return the timeout in milliseconds
   */
  @VisibleForTesting
  public long getReplayDrainTimeout()
  {
    return replayDrainTimeoutInMs;
  }
  /**
   * Sets how long this domain waits for the replay threads which are applying one of its
   * changes before it saves its ServerState and goes down.
   * <p>
   * Only there for the tests which check what a domain does when that wait runs out: they
   * can not hold a replay thread for {@link #REPLAY_DRAIN_TIMEOUT_IN_MS}.
   *
   * @param timeoutInMs the timeout in milliseconds
   */
  @VisibleForTesting
  public void setReplayDrainTimeout(long timeoutInMs)
  {
    replayDrainTimeoutInMs = timeoutInMs;
  }
  /**
   * Do what necessary when the data have changed : load state, load
   * generation Id.
   * If there is no such information check if there is a