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/RemotePendingChanges.java
@@ -17,7 +17,11 @@
 */
package org.opends.server.replication.plugin;
import static java.util.Collections.*;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List;
import java.util.NoSuchElementException;
import java.util.SortedMap;
import java.util.SortedSet;
@@ -90,8 +94,10 @@
   * 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.
   * applies it, or parks it as waiting for another change - a parked change is not the one
   * this thread is replaying, and giving it back is
   * {@link #releaseParkedChangesOwnedByCurrentThread()}, which reads the changes which are
   * waiting rather than this index (issue #954).
   * <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)}
@@ -545,9 +551,10 @@
   * 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).
   * change. The parked ones are left out: they are not the change this thread is replaying,
   * and giving one back is more than dropping its owner - it has to be unparked in the same
   * step, or it would be handed out by two roads at once, which is what
   * {@link #releaseParkedChangesOwnedByCurrentThread()} does (issues #922 and #954).
   * <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,
@@ -570,6 +577,95 @@
  }
  /**
   * Gives back the changes the calling thread parked as waiting for another change, and
   * takes them out of the changes which are waiting in the same step.
   * <p>
   * A parked change stays owned by the thread which parked it while that thread goes on
   * to the changes which follow: {@link #getNextUpdate()} is what hands it out again, to
   * whichever replay thread clears the change it was waiting for, and that thread takes it
   * over. A replay which is unwound leaves the thread which parked it without that road -
   * it takes the next delivery off the replay queue instead - so the change would be left
   * owned by a thread which is never coming back to it, and every redelivery of a change a
   * replay thread owns is refused as a duplicate (issue #954).
   * <p>
   * Unparking a change and giving it back is one step, under both locks, so that only one
   * road can hand it out: a change which was released while it is still listed as waiting
   * would be handed to the thread {@link #getNextUpdate()} gives it to and to the thread
   * which takes over the delivery which follows - the double replay the ownership is there
   * to prevent (OPENDJ-1115).
   * <p>
   * The changes stay listed and uncommitted, and stay among the changes the newer ones are
   * checked against, the way a change whose replay failed does: they are not in the data,
   * so they hold this domain's ServerState back and the changes which follow them keep
   * waiting for them.
   * <p>
   * The changes another thread parked are left alone: a change is given back by the thread
   * which owns it and by nobody else (issue #922). That thread may be inside the dependency
   * checks which parked it - they park a change once per dependency it has - so a change
   * released under it would be listed as waiting again a moment later, and handed out while
   * the delivery which took it over is being replayed.
   *
   * @return the CSNs of the changes it gave back, oldest first; empty when this thread owns
   *         no parked change - the changes a thread parked stay its own, whichever replay
   *         parked them, until {@link #getNextUpdate()} hands them to the thread which
   *         cleared what they wait for or they are given back here
   */
  List<CSN> releaseParkedChangesOwnedByCurrentThread()
  {
    final Thread current = Thread.currentThread();
    /*
     * The second lock is taken inside the try of the first: taking a lock which is held
     * by another thread allocates the node this one waits on, and this runs on the road
     * out of a JVM which has just refused an allocation. A throw out of the second lock
     * would otherwise unwind past the first with that one held by a thread which is
     * ending, and every road which lists, commits or gives back a change would wait for
     * it for good.
     */
    pendingChangesWriteLock.lock();
    try
    {
      dependentChangesLock.lock();
      try
      {
        if (dependentChanges.isEmpty())
        {
          // Nothing is waiting, which is the state every replay but a handful leaves behind.
          return emptyList();
        }
        /*
         * Sized for every change which is waiting, so that the one allocation of this
         * method past the locks is made before anything is taken out of the set. The rule
         * getNextUpdate() states for itself holds here: an allocation which fails once a
         * change has been unparked and released - and this runs on the road out of a JVM
         * which has just refused one - would have moved that change out of the hands which
         * hand it out again, with the caller never told that it did.
         */
        final List<CSN> released = new ArrayList<>(dependentChanges.size());
        final Iterator<PendingChange> it = dependentChanges.iterator();
        while (it.hasNext())
        {
          final PendingChange change = it.next();
          if (change.isOwnedBy(current))
          {
            it.remove();
            change.setOwner(null);
            released.add(change.getCSN());
          }
        }
        return released;
      }
      finally
      {
        dependentChangesLock.unlock();
      }
    }
    finally
    {
      pendingChangesWriteLock.unlock();
    }
  }
  /**
   * 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
@@ -672,8 +768,10 @@
       * 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).
       * over rather than leave it to nobody - and the give-back of the change a replay was
       * unwound on leaves it alone (issue #922). What hands a parked change back is
       * releaseParkedChangesOwnedByCurrentThread(), which unparks it in the same step so
       * that the two roads can not hand it out at once (issue #954).
       */
      changeBeingReplayed.remove(Thread.currentThread(), dependentChange.getCSN());
    }