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

Valery Kharseko
11 hours ago 776339a8c63c0bd8c83cfe8a619f6146794358e8
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
@@ -338,6 +338,17 @@
   * in flight unreplayable, and one alert per change would be a storm.
   */
  private static final long UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MS = 60000;
  /**
   * What the ack of a delivery whose replay ran out of memory says did not apply the
   * change.
   * <p>
   * A constant rather than a message built from the change, because this is read on the
   * road out of an {@code OutOfMemoryError}: formatting one asks the JVM for the memory
   * it has just refused, and what {@code processUpdateDone()} needs of this string is
   * that there is one - it sets {@code hasReplayError} on the ack and never sends the
   * text, which is why no ordinal is spent on it either.
   */
  private static final String REPLAY_RAN_OUT_OF_MEMORY = "the replay of this change ran out of memory";
  /** 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. */
@@ -2564,6 +2575,11 @@
  /**
   * Create and replay a synchronized Operation from an UpdateMsg.
   * <p>
   * The change is given back on the way out of a replay which was unwound before it
   * reached one of the roads which give it back: a change which stays owned by a thread
   * which is not replaying it anymore is refused as a duplicate on every later delivery,
   * so this domain's ServerState would never move past it (issue #922).
   *
   * @param msg
   *          The UpdateMsg to be replayed.
@@ -2572,6 +2588,153 @@
   */
  void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown)
  {
    try
    {
      replayChangeAndTheChangesWaitingForIt(msg, replayThreadShutdown);
    }
    catch (Throwable t)
    {
      /*
       * The roads which run to their end give the change back themselves, so what is left
       * to give back here is a change whose replay was unwound over them: by an
       * OutOfMemoryError, which is left to end this thread, or by a throw from what the
       * replay runs once the ack of the delivery has been published - the give-back of the
       * change and the hand-out of the changes which were waiting for it are on that road.
       *
       * It takes the road of a failed replay rather than being handed back on the spot, so
       * that a change which keeps unwinding the replays it is given to is counted as
       * failing and is eventually given up on: handing it back bare would have this domain
       * ask for it, and restart its session for it, for as long as the server is up.
       *
       * The changes this thread parked as waiting for another change are left alone: they
       * are handed to whichever thread clears the change they are waiting for, and that
       * thread takes them over.
       *
       * Which change this thread owns is read before anything is done with it, and that
       * read takes no lock and allocates nothing: everything below is gated on the answer,
       * so a lookup which threw in its turn - on the road out of a JVM which has just
       * refused an allocation - would leave the change listed, uncommitted and owned by a
       * thread which is about to end, which is the state this whole issue is about.
       */
      CSN owned = null;
      try
      {
        owned = remotePendingChanges.getChangeOwnedByCurrentThread();
        if (owned != null)
        {
          if (replayThreadShutdown.get() || shutdown.get() || disabled)
          {
            /*
             * The replay was not failing, it was being abandoned: this thread is stopping
             * because their number is being changed, or the domain is going away or is
             * being imported into. The change is handed back without being counted against
             * its give-up budget, which is the road a replay abandoned that way takes.
             */
            abandonReplay(owned);
          }
          else
          {
            /*
             * A JVM which has run out of memory is told apart here rather than reported on
             * its own road: the change is given back counted, the way every other unwound
             * replay gives it back, but the line which says it is being asked for again is
             * not built - that asks the JVM for the memory it has just refused - and the
             * session is restarted without sitting through the backoff, since the thread
             * which is doing it is on its way out.
             *
             * A change whose budget is spent is still reported and still raises its alert
             * on this road: it is the one line which says this replica has diverged, and an
             * operator who is not told would be left with a replica which is silently
             * behind. That is the deliberate exception to the rule above.
             */
            recoverFromReplayFailure(owned, replayThreadShutdown, t instanceof OutOfMemoryError);
          }
        }
      }
      catch (Throwable recoveryFailure)
      {
        /*
         * The give-back is the road which hands the change over, so a throw out of it - an
         * allocation which fails in its turn, where what unwound the replay was a JVM out
         * of memory - would leave the change owned by this thread after all. Hand it back
         * bare, without the failure count that road did not reach, and restart the session
         * so that it is delivered again: this is the last resort, the throwable is rethrown
         * whatever happens here, and the thread this runs on may well be ending on it - so
         * a request left for the next recovery of this domain to pick up is a request which
         * may never be run.
         */
        if (owned != null)
        {
          remotePendingChanges.replayFailed(owned);
          sessionRestartRequested.set(true);
          /*
           * Reported and restarted under guards of their own, and in that order: the report
           * is the line an operator acts on, and the restart is what has the change
           * delivered again - so a report which can not be formatted, on a road an
           * OutOfMemoryError leads to, must not cost the restart.
           */
          try
          {
            logger.error(ERR_REPLAY_GIVE_BACK_FAILED, owned, getBaseDN(),
                stackTraceToSingleLineString(recoveryFailure));
          }
          catch (Throwable reportFailure)
          {
            suppress(recoveryFailure, reportFailure);
          }
          try
          {
            runRequestedSessionRestarts(false);
          }
          catch (Throwable restartFailure)
          {
            /*
             * Nothing is left to try: the change is listed, uncommitted and unowned, so any
             * later session restart of this domain delivers it again. This goes with the
             * throwable which is rethrown below rather than being reported on its own.
             */
            suppress(recoveryFailure, restartFailure);
          }
        }
        // The error which unwound the replay is the one reported, whatever the give-back
        // ran into on top of it.
        suppress(t, recoveryFailure);
      }
      throw t;
    }
  }
  /**
   * Records the second throwable as one the first suppressed, unless the two are one and
   * the same.
   * <p>
   * A JVM which has run out of memory hands out the error it prepared before it ran out as
   * often as it is asked for one, so the roads out of a replay can carry the same instance
   * twice - and a throwable can not suppress itself.
   *
   * @param thrown the throwable which is reported
   * @param alsoThrown what was met on the way out of it
   */
  private static void suppress(Throwable thrown, Throwable alsoThrown)
  {
    if (thrown != alsoThrown)
    {
      thrown.addSuppressed(alsoThrown);
    }
  }
  /**
   * Replays the change of the provided message, then the changes which were waiting for
   * it, for as long as there are some.
   *
   * @param msg
   *          The UpdateMsg to be replayed.
   * @param replayThreadShutdown
   *          whether the replay thread was asked to stop
   */
  private void replayChangeAndTheChangesWaitingForIt(
      LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown)
  {
    // Try replay the operation, then flush (replaying) any pending operation
    // whose dependency has been replayed until no more left.
    do
@@ -2582,6 +2745,13 @@
      boolean replayAbandoned = false;
      String replayErrorMsg = null;
      CSN csn = null;
      /*
       * Read once, before anything of this delivery has run, so that the report of an ack
       * which could not be published names the change even when reading it off the message
       * is what threw: that report is written on the way out of a road which is already
       * failing, and it must not be the throw which unwinds the replay.
       */
      final CSN delivered = msg.getCSN();
      try
      {
        // The next operation for which to attempt replay.
@@ -2922,7 +3092,57 @@
          replayErrorMsg = message.toString();
          replayFailed = true;
        }
      } finally
      }
      catch (OutOfMemoryError e)
      {
        /*
         * The JVM is out of memory, which is not something to carry on replaying from: the
         * error is left to unwind the replay thread, which ends on it (issue #923). It is
         * not turned into a report of its own here either - the stack trace of an error
         * which can be raised anywhere says little, and the uncaught exception handler of
         * DirectoryThread writes the one line this is worth, with an alert.
         *
         * The other errors of the JVM take the road below. A StackOverflowError is met by
         * the thread which recursed and is gone once the stack has unwound, and the entry
         * being replayed is what raises it rather than the state of this server: ending a
         * thread on it would have one change this replica can not replay cost it a replay
         * thread per delivery, and nothing creates a replay thread to replace one which
         * ends.
         *
         * The change is given back, counted as failing and asked for again on the way out
         * (issue #922). The one thing done here is to make the ack this delivery publishes
         * below say that the change was not applied: a replica which is asking for a change
         * again must not have told an assured write that it is in the data here. It is a
         * constant rather than a message built from the change, because building one asks
         * the JVM for the memory it has just refused.
         */
        replayErrorMsg = REPLAY_RAN_OUT_OF_MEMORY;
        throw e;
      }
      catch (Error e)
      {
        /*
         * An Error out of the replay - a LinkageError met where a plugin or a backend class
         * is loaded, an AssertionError - unwinds every road this replay has out of here, the
         * ones which give the change back among them. The change is not in the data, so it
         * is reported and given back here, on the road every other failed replay takes: the
         * ack below says it was not applied, the failure counts against the give-up budget
         * of the change, the session is restarted for it to be delivered again (issue #922)
         * and the changes which were waiting for this one are replayed rather than left
         * waiting for a thread to hand them out.
         *
         * It is not given up on where no operation could be built from the message, which is
         * what an Exception at that point means: an Error says that this server could not run
         * the replay, not that the message is one no delivery could ever build an operation
         * from.
         */
        final LocalizableMessage message = ERR_ERROR_REPLAYING_CHANGE.get(
            msg.getCSN(), getBaseDN(), stackTraceToSingleLineString(e));
        logger.error(message);
        replayErrorMsg = message.toString();
        replayFailed = true;
      }
      finally
      {
        if (!dependency)
        {
@@ -2934,7 +3154,88 @@
           * 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);
          try
          {
            processUpdateDone(msg, replayErrorMsg);
          }
          catch (OutOfMemoryError e)
          {
            /*
             * The one throw from here which is not caught, on the same terms as the arm
             * above: a JVM which has run out of memory is not something to carry on
             * replaying from, the error is left to end this replay thread, and the uncaught
             * exception handler of DirectoryThread raises the alert #923 is about. Swallowed
             * here instead, it would have this thread go on to the roads below and to the
             * next change of the loop on an exhausted heap, with nothing reported anywhere.
             *
             * What the give-back on the way out of replay() then finds depends on the road
             * the replay took to get here. A replay which failed, or which was unwinding on
             * an OutOfMemoryError of its own, still owns its change: it is given back
             * counted, and the thread ends on this error rather than on the one it stepped
             * over. A replay which committed owns nothing anymore - commit() cleared the
             * owner, and the index the give-back reads, in the same step - so the give-back
             * is a no-op, and rightly so: a change which is in the data is not one to ask
             * for again. What that road steps over is getNextUpdate() below, so the changes
             * parked behind the committed change wait for the next replay of this domain to
             * hand them out. That is the trade #923 asks for: a thread which met this error
             * is not to carry on, not even for them.
             */
            throw e;
          }
          catch (Throwable ackFailure)
          {
            /*
             * Publishing the ack says nothing about whether the change was applied, so a
             * throw here must not be a road out of the replay. It would step over the
             * give-back of the change and, where the change was committed, over the
             * getNextUpdate() below - the one drain of the changes this thread parked as
             * waiting for it - leaving them waiting for a thread which is not replaying
             * anything anymore.
             *
             * A master which is waiting for the ack waits out its assured timeout either
             * way - and there is one only for an assured write in safe-read mode, which is
             * the one delivery a replica acknowledges. processUpdateDone() runs for every
             * delivery all the same: it accounts for the delivery in the receive window and
             * in the processed-updates counter, and a throw from there is a throw from this
             * bookkeeping, with no ack owed to anybody. Either way it is reported and the
             * replay carries on to the road the change itself decided: applied, failed and
             * asked for again, or given up on.
             *
             * Every step of processUpdateDone() catches what it can meet - the broker keeps
             * a failure to publish to itself and retries it - so what reaches here is what
             * no code of the replication protocol expected: an Error, and the tests drive
             * it as one.
             */
            try
            {
              logger.error(ERR_ACK_NOT_PUBLISHED, delivered, getBaseDN(),
                  stackTraceToSingleLineString(ackFailure));
            }
            catch (Throwable reportFailure)
            {
              /*
               * Guarded like the report of a give-back which failed: this one runs on the
               * same kind of road - the stack of the throwable is walked to build the line -
               * and a report which can not be built must not become the throw which unwinds
               * the replay past the give-back and past getNextUpdate().
               *
               * The line is tried once more with the name of the error alone, which walks no
               * stack. Nothing rethrows what was caught here, so an error recorded as
               * suppressed on it would be recorded nowhere, and this is the road on which an
               * operator has the least to go on. A second refusal leaves nothing to say it
               * with, and the replay carries on with what the change decided.
               */
              try
              {
                logger.error(ERR_ACK_NOT_PUBLISHED, delivered, getBaseDN(),
                    ackFailure.getClass().getName());
              }
              catch (Throwable secondReportFailure)
              {
                // Nothing is left to say it with.
              }
            }
          }
        }
      }
@@ -3191,6 +3492,29 @@
   */
  private boolean recoverFromReplayFailure(CSN csn, AtomicBoolean replayThreadShutdown)
  {
    return recoverFromReplayFailure(csn, replayThreadShutdown, false);
  }
  /**
   * Recovers from a change which could not be replayed, telling apart the replay which was
   * unwound by a JVM out of memory.
   *
   * @param csn
   *          the CSN of the change which could not be replayed
   * @param replayThreadShutdown
   *          whether the replay thread was asked to stop
   * @param outOfMemory
   *          whether what unwound the replay was the JVM running out of memory, in which
   *          case the line which says the change is being asked for again is not built -
   *          the report of a change which is given up on is, since it is what says this
   *          replica has diverged - and the session is restarted without the backoff
   * @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, boolean outOfMemory)
  {
    /*
     * 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()
@@ -3254,7 +3578,16 @@
      return true;
    }
    logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts());
    if (!outOfMemory)
    {
      /*
       * Not on the road out of a JVM which has run out of memory: building this line asks
       * it for the memory it has just refused, and the ack of the delivery already says
       * that the change was not applied. The constant that ack carries exists for the same
       * reason.
       */
      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
@@ -3267,9 +3600,12 @@
     * 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.
     * backoff on its way out, though: the backend is not what is going away. Neither does
     * the thread an OutOfMemoryError is ending, for the same reason - and the restart is
     * run rather than left to be asked for again, because that thread will not be there to
     * run it, and a change nobody asks for again holds this domain's ServerState back.
     */
    runRequestedSessionRestarts(!replayThreadShutdown.get());
    runRequestedSessionRestarts(!replayThreadShutdown.get() && !outOfMemory);
    return true;
  }
@@ -3292,7 +3628,41 @@
      {
        while (sessionRestartRequested.getAndSet(false))
        {
          restartSession(wait);
          boolean restarted = false;
          try
          {
            restartSession(wait);
            restarted = true;
          }
          finally
          {
            if (!restarted)
            {
              /*
               * The request is put back where it was taken from. The flag is read and
               * cleared before the restart runs, so a restart which ends abruptly - the
               * session is stopped first, and starting it again creates a listener thread,
               * which the operating system can refuse - would otherwise leave this domain
               * with no session and with nothing left to ask for one.
               *
               * What a request left standing buys is bounded, and the bound is worth
               * stating. Its two readers are the roads out of a failed and of an abandoned
               * replay of this domain, and with no listener thread nothing is delivered
               * anymore: the replays left to run are the changes already taken off the
               * session - the ones waiting in the replay queue, and the ones parked as
               * dependencies. One of those failing finds the request standing and runs the
               * restart, which starts from a clean state, since disableService() drops the
               * listener thread which was never started. Once they are spent, the domain
               * stays down until it is disabled and enabled back, or the server is
               * restarted. That is said where it can be heard: a refused thread is an
               * OutOfMemoryError, and one which leaves recoverFromReplayFailure() or
               * abandonReplay() ends the replay thread it is met on, so the uncaught
               * exception handler of DirectoryThread writes the line and raises the alert,
               * with the start of the listener thread in the trace.
               */
              sessionRestartRequested.set(true);
            }
          }
        }
      }
      finally