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/RemotePendingChanges.java
@@ -23,6 +23,8 @@
import java.util.SortedSet;
import java.util.TreeMap;
import java.util.TreeSet;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ConcurrentSkipListSet;
import java.util.concurrent.locks.ReentrantLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
@@ -73,6 +75,32 @@
   */
  private final ConcurrentSkipListSet<PendingChange> activeAndDependentChanges = new ConcurrentSkipListSet<>();
  /**
   * The change each replay thread is replaying right now, read by the give-back on the way
   * out of a replay which was unwound.
   * <p>
   * It is an index of what {@link PendingChange#isOwnedBy(Thread)} already says rather than
   * a second copy of it: it is written in the same locked step wherever a change is taken
   * over or given back, and every road which acts on what it answers checks the ownership
   * of the change again. What it buys is the read. That read is made on a road an
   * {@code OutOfMemoryError} leads to, and looking the change up by walking the pending
   * changes takes two locks and allocates an iterator - an allocation a JVM which has just
   * refused one may well refuse again, and the change would then be left listed,
   * uncommitted and owned by a thread which is not replaying it anymore, which is the wedge
   * this issue is about (issue #922).
   * <p>
   * A thread is entered here when it takes a change over and removed when it gives it back,
   * applies it, or parks it as waiting for another change - the parked ones are handed to
   * whichever thread clears what they wait for, so they are not this one's to give back.
   * <p>
   * The entry of a thread is written by that thread and by nobody else, and that - not the
   * lock - is what keeps the writes apart: the park in {@link #addDependency(PendingChange)}
   * clears it under the read lock, where every other writer holds the write lock, and it
   * races nothing for it. {@link #clear()} is the one writer of every entry, and it holds
   * both locks.
   */
  private final ConcurrentMap<Thread, CSN> changeBeingReplayed = new ConcurrentHashMap<>();
  private final ReentrantReadWriteLock pendingChangesLock = new ReentrantReadWriteLock(true);
  private final ReentrantReadWriteLock.ReadLock pendingChangesReadLock = pendingChangesLock.readLock();
  private final ReentrantReadWriteLock.WriteLock pendingChangesWriteLock = pendingChangesLock.writeLock();
@@ -209,9 +237,25 @@
  /**
   * Mark an update message as committed.
   * <p>
   * A change another replay thread owns is not this thread's to record: it is reported as
   * a change which is not here, the way one which is not listed anymore is. A thread
   * decides that a change is in the data, or that it is to be given up on, and records
   * that decision a turn of this lock later - long enough for the change to have been
   * handed back, delivered again and taken over in between. Recording it then would
   * advance the ServerState over a change which is not in the data yet (issue #889) and
   * have the thread which is applying it right now fail to commit.
   * <p>
   * The give-back on the way out of an unwound replay is what makes this reachable: it
   * runs wherever the replay was left, so the checks which keep it from taking a change
   * away from the thread which owns it now belong on every road which reads ownership,
   * not only on {@link #replayFailed(CSN)} (issue #922).
   *
   * @param csn
   *          The CSN of the update message that must be set as committed.
   * @throws NoSuchElementException
   *          if there is no change with that CSN for this thread to record: it is not
   *          listed as pending anymore, or another replay thread owns it
   */
  public void commit(CSN csn)
  {
@@ -219,12 +263,14 @@
    try
    {
      PendingChange curChange = pendingChanges.get(csn);
      if (curChange == null)
      if (curChange == null
          || (curChange.isOwned() && !curChange.isOwnedBy(Thread.currentThread())))
      {
        throw new NoSuchElementException();
      }
      curChange.setCommitted(true);
      curChange.setOwned(false);
      curChange.setOwner(null);
      changeBeingReplayed.remove(Thread.currentThread(), csn);
      activeAndDependentChanges.remove(curChange);
      final Iterator<PendingChange> it = pendingChanges.values().iterator();
@@ -271,6 +317,12 @@
   * The changes another replay thread is applying right now are left alone: they are
   * about to commit, and forgetting them would have their {@code commit()} fail, the
   * ServerState stay behind them and the replication server replay them a second time.
   * <p>
   * Only the thread which owns the change gives it back. A release which arrives from
   * another one is a release of a change which has been taken over since - the failure of
   * a delivery is reported after the change it carried was handed to the delivery which
   * follows it - and taking the change away from the thread which is replaying it right
   * now is the double replay the ownership is there to prevent (issue #922).
   *
   * @param csn the CSN of the change whose replay failed
   */
@@ -280,9 +332,10 @@
    try
    {
      final PendingChange change = pendingChanges.get(csn);
      if (change != null && !change.isCommitted())
      if (change != null && !change.isCommitted() && change.isOwnedBy(Thread.currentThread()))
      {
        change.setOwned(false);
        change.setOwner(null);
        changeBeingReplayed.remove(Thread.currentThread(), csn);
      }
    }
    finally
@@ -339,9 +392,13 @@
   *          the CSN of the change whose replay failed
   * @param nowMs
   *          when it failed, on a clock which only moves forward
   * @return the failures of the change, or {@code null} when it is not listed as an
   *         uncommitted change anymore, which happens when the domain was disabled while
   *         it was being replayed: there is no change left here to give up on
   * @return the failures of the change, or {@code null} when there is no change left here
   *         for this thread to give up on: it is not listed as an uncommitted change
   *         anymore, which happens when the domain was disabled while it was being
   *         replayed, or another replay thread owns it, which is that change having been
   *         delivered again and taken over while this thread was on its way to reporting
   *         on it. Spending the give-up budget of a change another thread is applying
   *         would have this replica skip a change which is being written (issue #922)
   */
  public ReplayFailure recordReplayFailure(CSN csn, long nowMs)
  {
@@ -349,7 +406,8 @@
    try
    {
      final PendingChange change = pendingChanges.get(csn);
      if (change == null || change.isCommitted())
      if (change == null || change.isCommitted()
          || (change.isOwned() && !change.isOwnedBy(Thread.currentThread())))
      {
        return null;
      }
@@ -434,6 +492,7 @@
      pendingChanges.clear();
      dependentChanges.clear();
      activeAndDependentChanges.clear();
      changeBeingReplayed.clear();
      failingChanges = 0;
    }
    finally
@@ -465,8 +524,16 @@
      {
        return false;
      }
      change.setOwned(true);
      /*
       * Listed as being replayed before it is owned, and not the other way round: the
       * caller enters the replay - where the give-back on the way out lives - once this
       * returns, so nothing which allocates must run between the owner being stamped and
       * that. A change listed here without an owner is the state a failed replay leaves
       * behind, and the dependency checks which read this set do not read the owner.
       */
      activeAndDependentChanges.add(change);
      changeBeingReplayed.put(Thread.currentThread(), change.getCSN());
      change.setOwner(Thread.currentThread());
      return true;
    }
    finally
@@ -475,7 +542,40 @@
    }
  }
  /**
   * Returns the CSN of the change the calling thread is replaying, when it still owns one.
   * <p>
   * A thread owns the change it is replaying and the ones it parked as waiting for another
   * change. The parked ones are left out: they are handed to whichever thread clears the
   * change they are waiting for, and that thread takes them over, so giving one back here
   * would have the same change handed to two threads (issue #922).
   * <p>
   * It is a plain read of {@link #changeBeingReplayed}: no lock is taken and nothing is
   * allocated. This is what the give-back on the way out of an unwound replay asks first,
   * and it runs on a road an {@code OutOfMemoryError} leads to - a lookup which allocated
   * could be refused in its turn, and a give-back which does not know which change to give
   * back leaves it listed, uncommitted and owned by a thread which is not replaying it
   * anymore, which is the wedge this issue is about.
   * <p>
   * An answer which is out of date is safe: every road which acts on it - {@code commit()},
   * {@link #recordReplayFailure(CSN, long)} and {@link #replayFailed(CSN)} - checks the
   * ownership of the change again under the write lock, and is a no-op for a change this
   * thread does not own anymore.
   *
   * @return the CSN of the change this thread is replaying, or {@code null} when it does
   *         not own one anymore - the road its replay took gave it back, or it was applied
   */
  CSN getChangeOwnedByCurrentThread()
  {
    return changeBeingReplayed.get(Thread.currentThread());
  }
  /**
   * Get the first update in the list that have some dependencies cleared.
   * <p>
   * The change is handed to the calling thread, which owns it from then on: it is
   * replayed by whichever replay thread cleared the change it was waiting for rather than
   * by the one which parked it, and a change is given back by the thread which owns it
   * and by nobody else (issue #922).
   *
   * @return The LDAPUpdateMsg to be handled.
   */
@@ -485,22 +585,65 @@
    dependentChangesLock.lock();
    try
    {
      if (!dependentChanges.isEmpty() && !pendingChanges.isEmpty())
      if (!hasChangeToHandOut())
      {
        PendingChange firstDependentChange = dependentChanges.first();
        if (pendingChanges.firstKey().isNewerThanOrEqualTo(firstDependentChange.getCSN()))
        {
          dependentChanges.remove(firstDependentChange);
          return firstDependentChange.getLDAPUpdateMsg();
        }
        /*
         * Nothing is waiting, or what waits is still held back by the changes before it.
         * This is called at the end of every replay, by every replay thread, so the answer
         * is looked for under the read lock: taking the write lock here would have a
         * backlog of waiting changes serialize the replay of the changes which have none.
         */
        return null;
      }
      return null;
    }
    finally
    {
      dependentChangesLock.unlock();
      pendingChangesReadLock.unlock();
    }
    /*
     * There is one to hand out, and handing it out writes its owner, which is written
     * under the write lock as the rest of the state of a change is. It is looked for again
     * under that lock: another replay thread may have been handed it in between.
     */
    pendingChangesWriteLock.lock();
    dependentChangesLock.lock();
    try
    {
      if (hasChangeToHandOut())
      {
        final PendingChange firstDependentChange = dependentChanges.first();
        /*
         * Entered as the change this thread is replaying before it is taken out of the
         * ones which are waiting, for the same reason markInProgress() enters it before it
         * stamps the owner: an allocation which fails here must leave the change where it
         * was rather than take it out of the hands which would hand it out again.
         */
        changeBeingReplayed.put(Thread.currentThread(), firstDependentChange.getCSN());
        dependentChanges.remove(firstDependentChange);
        firstDependentChange.setOwner(Thread.currentThread());
        return firstDependentChange.getLDAPUpdateMsg();
      }
      return null;
    }
    finally
    {
      dependentChangesLock.unlock();
      pendingChangesWriteLock.unlock();
    }
  }
  /**
   * Returns whether the first change waiting for another one can be replayed now, that is
   * whether every change before it has left the pending changes.
   */
  @GuardedBy("pendingChangesLock, dependentChangesLock")
  private boolean hasChangeToHandOut()
  {
    return !dependentChanges.isEmpty()
        && !pendingChanges.isEmpty()
        && pendingChanges.firstKey().isNewerThanOrEqualTo(dependentChanges.first().getCSN());
  }
  /**
@@ -524,6 +667,15 @@
      {
        dependentChanges.add(dependentChange);
      }
      /*
       * Whichever of the two it was, this thread is not replaying that change anymore: a
       * parked one is handed to the thread which clears what it waits for, and one which is
       * not listed here anymore is gone with the pending changes of a domain which was
       * disabled. The owner stays as it is - it is what has getNextUpdate() hand the change
       * over rather than leave it to nobody - and the give-back on the way out of an
       * unwound replay leaves it alone (issue #922).
       */
      changeBeingReplayed.remove(Thread.currentThread(), dependentChange.getCSN());
    }
    finally
    {