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

Valery Kharseko
2 days ago 80481f756d71bd58b4bda627758dcd774e9d5dcd
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
@@ -2760,20 +2760,35 @@
       * 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.
       * The changes this thread parked as waiting for another change are given back too,
       * and before the road above runs: that road restarts the session, and a change which
       * is still owned when the replication server sends it again over it is turned down as
       * a duplicate - the one delivery which could have taken it over (issue #954).
       *
       * 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.
       * Which change this thread owns is read first of all, 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. It is read before the parked changes
       * are given back rather than after, because that give-back allocates and can throw on
       * the same road, and the last resort below can only hand back a change it was told
       * about.
       */
      CSN owned = null;
      try
      {
        owned = remotePendingChanges.getChangeOwnedByCurrentThread();
        /*
         * Asked for without the backoff on the two roads recoverFromReplayFailure() asks
         * for it without on: a thread which is stopping, and one which an OutOfMemoryError
         * is ending - the backend is not what is going away. What is asked for is not what
         * is run: a request is never answered by less than it asked for, so a restart
         * another road asked for with the backoff, or one given back with it, keeps its
         * wait whichever thread runs it.
         */
        final boolean parkedGivenBack = giveBackParkedChanges(
            replayThreadShutdown.get() || t instanceof OutOfMemoryError
                ? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF);
        if (owned != null)
        {
          if (replayThreadShutdown.get() || shutdown.get() || disabled)
@@ -2808,6 +2823,27 @@
            recoverFromReplayFailure(owned, replayThreadShutdown, t instanceof OutOfMemoryError);
          }
        }
        if (parkedGivenBack && !replayThreadShutdown.get() && !sessionHasAnOwner())
        {
          /*
           * The road the change this thread was replaying took may have run the restart the
           * give-back asked for - they ask for the same one - and it may have had none to
           * run: this thread owned no change, or the change it owned was given up on. What
           * is still requested is run here rather than left standing: the state
           * checkpointer would run it within its tick, so this is the latency of the
           * delivery the changes which were handed back wait for, and nothing more - a
           * request this thread leaves is not lost.
           *
           * A thread which is stopping leaves it standing, the way abandonReplay() does:
           * the state checkpointer runs one restart for every change the threads of the
           * pool hand back on their way out, rather than each of them running one while
           * the configuration change which is stopping them waits. A domain whose session
           * has an owner is left alone the way the give-back left it: nothing was asked
           * for on that road, and a request another thread left standing is not this
           * one's to spend on a restart which is refused where it runs.
           */
          runRequestedSessionRestarts();
        }
      }
      catch (Throwable recoveryFailure)
      {
@@ -2819,7 +2855,15 @@
         * 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. A
         * restart which can not run here leaves its request standing, and the state
         * checkpointer of this domain runs it.
         * checkpointer of this domain runs it: one which threw asks for itself again on
         * its way out, and one a domain whose session has an owner would refuse is not
         * run at all rather than spent on the refusal.
         *
         * The restart is asked for once the change is released and not before, the way
         * every road which releases one asks: asked for first, it could be taken and run
         * by another thread while this one still owned the change, and the delivery the
         * new session brought would be turned down as the duplicate of a change a replay
         * thread owns, with nothing left standing to ask for it again.
         */
        if (owned != null)
        {
@@ -2840,20 +2884,42 @@
          {
            suppress(recoveryFailure, reportFailure);
          }
          try
        }
        try
        {
          /*
           * Outside the guard above: two roads reach here with a request standing and no
           * change of this thread's to hand back, and both are the parked changes' - the
           * give-back which released them asks for the restart before it reports them, and
           * a throw out of the report - the JVM which unwound this replay is out of memory
           * - leaves the request standing; and a restart the parked road ran and which
           * threw has asked for one again on its way out. The changes it released are
           * listed, uncommitted and unowned, so the request is what brings them back, and
           * this thread is the one there to run it (issue #954). A give-back which threw
           * before it released anything left the parked changes as they were, owned by this
           * thread and handed out by getNextUpdate() to whichever thread clears what they
           * wait for: nothing here can do better for those.
           *
           * Not run on a domain whose session has an owner, the way no road of a failed
           * replay runs it there: the restart is refused where it runs and the request
           * would be spent on the refusal, while a request left standing is run by the
           * state checkpointer once the owner is gone - or forgotten with the pending
           * changes it was made for, when the owner forgets them on its way out.
           */
          if (!sessionHasAnOwner())
          {
            runRequestedSessionRestarts();
          }
          catch (Throwable restartFailure)
          {
            /*
             * Nothing is left to try here: the change is listed, uncommitted and unowned,
             * and the restart which threw has asked for one again, so the session restart
             * the state checkpointer runs delivers it again. This goes with the throwable
             * which is rethrown below rather than being reported on its own.
             */
            suppress(recoveryFailure, restartFailure);
          }
        }
        catch (Throwable restartFailure)
        {
          /*
           * Nothing is left to try here: the changes are listed, uncommitted and unowned,
           * and the restart which threw has asked for one again, so the session restart
           * the state checkpointer runs delivers them 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.
@@ -3387,12 +3453,14 @@
             * 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.
             * owner, and the index the give-back reads, in the same step - so the change it
             * was replaying is not given back, 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, which hands out the changes parked behind the committed change: the
             * ones this thread parked are given back on the way out of replay() and the
             * session is restarted for them (issue #954), the ones other threads parked 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;
          }
@@ -3968,6 +4036,25 @@
  }
  /**
   * Gives back the changes a replay thread which is stopping parked in this domain.
   * <p>
   * Called by that thread on its way out (issue #986). The restart which brings them back
   * is asked for and left standing, the way the thread leaves the request it makes for the
   * change it abandons: the threads of the pool are stopped one after the other and joined,
   * and each running a restart on its way out would have the configuration change which is
   * stopping them wait for one restart per thread. The state checkpointer of this domain
   * runs one restart for the lot within its tick, and holds it while a total update runs
   * over the session, in either direction - a change delivered again before the pool which
   * replaces this thread is up waits in the replay queue for it. Asked for without the
   * backoff: what went away is a replay thread, not the backend, and these changes were
   * never applied here.
   */
  void giveBackChangesParkedByStoppingThread()
  {
    giveBackParkedChanges(SessionRestart.NOW);
  }
  /**
   * Restarts the session as long as changes which could not be replayed are waiting to be
   * delivered again.
   */
@@ -4076,6 +4163,84 @@
  }
  /**
   * Gives back the changes this replay thread parked as waiting for another change, on the
   * way out of a replay which was unwound.
   * <p>
   * A parked change is handed out again by {@code getNextUpdate()} alone, which every
   * replay loop of this domain runs once it is done with a change: a parked change is
   * replayed by whichever thread clears the change it was waiting for. A thread whose
   * replay was unwound is not on that road anymore - it takes the next delivery off the
   * replay queue - so a change it parked would be left owned by a thread which is not
   * coming back to it, while every redelivery of it is refused as a duplicate. On a domain
   * which then goes quiet that change is where this replica's ServerState, and every change
   * behind it from every master, stops (issue #954).
   * <p>
   * They are handed back without a failure being counted against them: they were never
   * applied here, so the give-up budget which decides when this replica skips a change it
   * can not apply is not this delivery's to spend, the way it is not for a change abandoned
   * by a replay thread which is stopping.
   * <p>
   * The delivery which carried one published no ack - the ack of a parked change is
   * published by the delivery which replays it - so it is counted as processed here, the
   * way a delivery which is dropped rather than replayed is: that count is of the
   * deliveries this replica took off the session, and these are over. The window they hold
   * is not given back either, and does not need to be: the session they came over is about
   * to be restarted, and a session which starts is given its receive window anew.
   * <p>
   * On a domain whose session has an owner - the domain itself, going away, or a total
   * update into it, from the moment it is asked for - they are released and nothing more,
   * the way {@code abandonReplay()} hands a change back on that road (see
   * {@link #sessionHasAnOwner()}): there is no session of this thread's to restart, and
   * the restart it would ask for is refused where it runs. The domain forgets its pending
   * changes on its way down, the import forgets them at its end, and a change released
   * for a total update which never begins stays listed until the next failed replay of
   * this domain restarts the session, which has the replication server send it again. A
   * line which says the replication server sends the change again would not hold on any
   * of these - a server which is shutting down abandons every change in flight, and none
   * of them is delivered again before it is started back.
   * <p>
   * A replay thread which is stopping calls this through
   * {@link #giveBackChangesParkedByStoppingThread()}, for every domain of this server: what
   * it parked would be left owned by a thread which does not exist anymore, and every
   * redelivery of a change a replay thread owns is refused as a duplicate (issue #986). The
   * request it makes here is left standing for the state checkpointer, the way that thread
   * leaves the request it makes for the change it abandons.
   *
   * @param restart what the session restart is asked for as: with the backoff a failing
   *          backend is owed, or without it on a thread which is stopping or which an
   *          OutOfMemoryError is ending
   * @return whether any change was handed back: a change which nobody owns is one only a
   *         new delivery brings back, so the caller runs the restart asked for them - on
   *         a thread which is not stopping, and on a domain whose session has no owner
   */
  private boolean giveBackParkedChanges(SessionRestart restart)
  {
    final List<CSN> parked = remotePendingChanges.releaseParkedChangesOwnedByCurrentThread();
    if (parked.isEmpty())
    {
      return false;
    }
    if (sessionHasAnOwner())
    {
      // The domain owns its session, or a total update does: both forget the pending
      // changes, and neither leaves a session for this thread to restart.
      return true;
    }
    /*
     * Asked for before the changes are reported: a throw out of the report - the JVM which
     * unwound this replay is out of memory - must not lose the restart which is what brings
     * them back.
     */
    sessionRestarts.request(restart);
    for (CSN csn : parked)
    {
      incProcessedUpdates();
      logger.info(NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK, csn, getBaseDN());
    }
    return true;
  }
  /**
   * Gives a change back to the replication server when this replay thread stops before it
   * could apply it.
   * <p>