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

Valery Kharseko
2 days ago 8377b61404033c479afa38424211314e84c2105a
Replace per-operation backend read lock with a scalable shared-access gate (#680)
4 files modified
1 files added
318 ■■■■ changed files
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/BackendImpl.java 60 ●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/EntryContainer.java 112 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/RootContainer.java 4 ●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/StripedCounter.java 95 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/test/java/org/opends/server/backends/pluggable/PluggableBackendImplTestCase.java 47 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/BackendImpl.java
@@ -29,7 +29,6 @@
import java.util.Set;
import java.util.SortedSet;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicInteger;
import org.forgerock.i18n.LocalizableException;
import org.forgerock.i18n.LocalizableMessage;
@@ -93,9 +92,14 @@
  /** The root container to use for this backend. */
  private RootContainer rootContainer;
  // FIXME: this is broken. Replace with read-write lock.
  /** A count of the total operation threads currently in the backend. */
  private final AtomicInteger threadTotalCount = new AtomicInteger(0);
  /**
   * A count of the total operation threads currently in the backend. Bumped
   * twice per operation by all worker threads, so it uses a striped counter
   * to avoid contending on a single cache line; it is only read when waiting
   * for the backend to become quiescent, which is why it is not a LongAdder —
   * see {@link StripedCounter}.
   */
  private final StripedCounter threadTotalCount = new StripedCounter();
  /** The base DNs defined for this backend instance. */
  private Set<DN> baseDNs;
@@ -146,14 +150,14 @@
      throw new DirectoryException(
          noEntryContainerResultCode, ERR_BACKEND_ENTRY_DOESNT_EXIST.get(entryDN, getBackendID()));
    }
    threadTotalCount.getAndIncrement();
    threadTotalCount.increment();
    return ec;
  }
  /** End a Backend API method that accesses the EntryContainer. */
  private void accessEnd()
  {
    threadTotalCount.getAndDecrement();
    threadTotalCount.decrement();
  }
  /**
@@ -163,7 +167,7 @@
   */
  private void waitUntilQuiescent()
  {
    while (threadTotalCount.get() > 0)
    while (threadTotalCount.sum() > 0)
    {
      // Still have threads accessing the storage so sleep a little
      try
@@ -268,7 +272,7 @@
    }
    // Make sure the thread counts are zero for next initialization.
    threadTotalCount.set(0);
    threadTotalCount.reset();
    // Log an informational message.
    logger.info(NOTE_BACKEND_OFFLINE, cfg.getBackendId());
@@ -356,7 +360,7 @@
      throw de;
    }
    container.sharedLock.lock();
    container.beginSharedAccess();
    try
    {
      return ConditionResult.valueOf(container.hasSubordinates(entryDN));
@@ -367,7 +371,7 @@
    }
    finally
    {
      container.sharedLock.unlock();
      container.endSharedAccess();
      accessEnd();
    }
  }
@@ -378,7 +382,7 @@
    checkNotNull(baseDN, "baseDN must not be null");
    final EntryContainer ec = accessBegin(null, baseDN);
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
      return ec.getNumberOfEntriesInBaseDN();
@@ -390,7 +394,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
@@ -417,7 +421,7 @@
      throw de;
    }
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
      return ec.getNumberOfChildren(parentDN);
@@ -428,7 +432,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
@@ -437,7 +441,7 @@
  public boolean entryExists(final DN entryDN) throws DirectoryException
  {
    EntryContainer ec = accessBegin(null, entryDN);
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
      return ec.entryExists(entryDN);
@@ -448,7 +452,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
@@ -457,7 +461,7 @@
  public Entry getEntry(DN entryDN) throws DirectoryException
  {
    EntryContainer ec = accessBegin(null, entryDN);
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
      return ec.getEntry(entryDN);
@@ -468,7 +472,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
@@ -478,7 +482,7 @@
  {
    EntryContainer ec = accessBegin(addOperation, entry.getName());
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
      ec.addEntry(entry, addOperation);
@@ -489,7 +493,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
@@ -500,7 +504,7 @@
  {
    EntryContainer ec = accessBegin(deleteOperation, entryDN);
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
      ec.deleteEntry(entryDN, deleteOperation);
@@ -511,7 +515,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
@@ -522,7 +526,7 @@
  {
    EntryContainer ec = accessBegin(modifyOperation, newEntry.getName());
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
@@ -534,7 +538,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
@@ -554,7 +558,7 @@
      throw new DirectoryException(ResultCode.UNWILLING_TO_PERFORM, WARN_FUNCTION_NOT_SUPPORTED.get());
    }
    currentContainer.sharedLock.lock();
    currentContainer.beginSharedAccess();
    try
    {
      currentContainer.renameEntry(currentDN, entry, modifyDNOperation);
@@ -565,7 +569,7 @@
    }
    finally
    {
      currentContainer.sharedLock.unlock();
      currentContainer.endSharedAccess();
      accessEnd();
    }
  }
@@ -577,7 +581,7 @@
    // is concerned: report it as such instead of the UNDEFINED result code used internally.
    EntryContainer ec = accessBegin(searchOperation, searchOperation.getBaseDN(), ResultCode.NO_SUCH_OBJECT);
    ec.sharedLock.lock();
    ec.beginSharedAccess();
    try
    {
@@ -589,7 +593,7 @@
    }
    finally
    {
      ec.sharedLock.unlock();
      ec.endSharedAccess();
      accessEnd();
    }
  }
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/EntryContainer.java
@@ -236,7 +236,7 @@
    {
      final ConfigChangeResult ccr = new ConfigChangeResult();
      exclusiveLock.lock();
      EntryContainer.this.lock();
      try
      {
        storage.write(new WriteOperation()
@@ -256,7 +256,7 @@
      }
      finally
      {
        exclusiveLock.unlock();
        EntryContainer.this.unlock();
      }
      return ccr;
@@ -318,7 +318,7 @@
    public ConfigChangeResult applyConfigurationDelete(final BackendVLVIndexCfg cfg)
    {
      final ConfigChangeResult ccr = new ConfigChangeResult();
      exclusiveLock.lock();
      EntryContainer.this.lock();
      try
      {
        storage.write(new WriteOperation()
@@ -337,7 +337,7 @@
      }
      finally
      {
        exclusiveLock.unlock();
        EntryContainer.this.unlock();
      }
      return ccr;
    }
@@ -348,6 +348,73 @@
  final Lock sharedLock = lock.readLock();
  final Lock exclusiveLock = lock.writeLock();
  /**
   * Striped count of in-flight lock-free shared accesses. Every operation
   * (search, bind, compare, modify, ...) enters the entry container through
   * {@link #beginSharedAccess()}, so acquiring even the read side of the
   * ReentrantReadWriteLock becomes a cross-core hotspot under load: each
   * acquire and release CAS-es the single lock state word. The hot paths
   * register through this striped counter instead and only fall back to
   * waiting when an exclusive locker has closed the gate; exclusive lockers
   * (rare structural changes: index removal, configuration changes, close)
   * close the gate through {@link #lock()} and drain in-flight accesses. The
   * drain relies on {@link StripedCounter#sum()} never under-counting to a
   * false zero, which is why this is not a LongAdder — see StripedCounter.
   */
  private final StripedCounter sharedAccessCount = new StripedCounter();
  /** True while an exclusive locker has closed the gate for lock-free shared access. */
  private volatile boolean exclusiveAccessPending;
  /** Monitor used to park shared accessors while the gate is closed. */
  private final Object sharedAccessMonitor = new Object();
  /**
   * Begins a lock-free shared access to this entry container. Must be paired
   * with {@link #endSharedAccess()} in a finally block on the same thread.
   * Equivalent to acquiring {@link #sharedLock}, but scales with the number
   * of cores.
   */
  void beginSharedAccess()
  {
    boolean interrupted = false;
    for (;;)
    {
      sharedAccessCount.increment();
      if (!exclusiveAccessPending)
      {
        break;
      }
      // An exclusive locker is active or draining: back out and wait.
      sharedAccessCount.decrement();
      synchronized (sharedAccessMonitor)
      {
        while (exclusiveAccessPending)
        {
          try
          {
            sharedAccessMonitor.wait();
          }
          catch (InterruptedException e)
          {
            interrupted = true;
          }
        }
      }
    }
    if (interrupted)
    {
      Thread.currentThread().interrupt();
    }
  }
  /**
   * Ends a lock-free shared access to this entry container. Must be called by
   * the thread that did the paired {@link #beginSharedAccess()}.
   */
  void endSharedAccess()
  {
    sharedAccessCount.decrement();
  }
  EntryContainer(DN baseDN, String backendID, PluggableBackendCfg config, Storage storage, RootContainer rootContainer,
      ServerContext serverContext) throws ConfigException
  {
@@ -2412,7 +2479,7 @@
  {
    final ConfigChangeResult ccr = new ConfigChangeResult();
    exclusiveLock.lock();
    EntryContainer.this.lock();
    try
    {
      storage.write(new WriteOperation()
@@ -2436,7 +2503,7 @@
    }
    finally
    {
      exclusiveLock.unlock();
      EntryContainer.this.unlock();
    }
    return ccr;
@@ -2731,15 +2798,44 @@
    searchOp.addResponseControl(new VLVResponseControl(targetPosition, contentCount, vlvResultCode));
  }
  /** Get the exclusive lock. */
  /**
   * Get the exclusive lock: acquires the write lock (excluding legacy
   * sharedLock readers), closes the gate for lock-free shared accessors and
   * drains the in-flight ones.
   */
  void lock()
  {
    exclusiveLock.lock();
    exclusiveAccessPending = true;
    // The drain must not be abandoned on interrupt: returning early would let
    // the exclusive caller run concurrently with in-flight shared accesses.
    // Exclusive lockers are rare, so sleep-polling is an acceptable trade-off.
    boolean interrupted = false;
    while (sharedAccessCount.sum() != 0)
    {
      try
      {
        Thread.sleep(1);
      }
      catch (InterruptedException e)
      {
        interrupted = true;
      }
    }
    if (interrupted)
    {
      Thread.currentThread().interrupt();
    }
  }
  /** Unlock the exclusive lock. */
  /** Unlock the exclusive lock and reopen the gate for lock-free shared accessors. */
  void unlock()
  {
    exclusiveAccessPending = false;
    synchronized (sharedAccessMonitor)
    {
      sharedAccessMonitor.notifyAll();
    }
    exclusiveLock.unlock();
  }
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/RootContainer.java
@@ -324,14 +324,14 @@
    for (DN baseDN : entryContainers.keySet())
    {
      EntryContainer ec = unregisterEntryContainer(baseDN);
      ec.exclusiveLock.lock();
      ec.lock();
      try
      {
        ec.close();
      }
      finally
      {
        ec.exclusiveLock.unlock();
        ec.unlock();
      }
    }
    config.removePluggableChangeListener(this);
opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/StripedCounter.java
New file
@@ -0,0 +1,95 @@
/*
 * The contents of this file are subject to the terms of the Common Development and
 * Distribution License (the License). You may not use this file except in compliance with the
 * License.
 *
 * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the
 * specific language governing permission and limitations under the License.
 *
 * When distributing Covered Software, include this CDDL Header Notice in each file and include
 * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL
 * Header, with the fields enclosed by brackets [] replaced by your own identifying
 * information: "Portions Copyright [year] [name of copyright owner]".
 *
 * Copyright 2026 3A Systems, LLC
 */
package org.opends.server.backends.pluggable;
import java.util.concurrent.atomic.AtomicLongArray;
/**
 * A striped counter of in-flight accesses whose non-atomic {@link #sum()} scan
 * is safe for quiescence detection: it may transiently over-estimate, but can
 * never return zero while an access is still in flight.
 * <p>
 * The increment and the matching decrement of one logical access must be
 * performed by the same thread; the stripe is a pure function of the thread,
 * so both land in the same slot. A scan reads each slot once, i.e. observes a
 * prefix of each slot's modification history, and within one slot a decrement
 * can never be observed without the increment that preceded it, so every
 * per-slot subtotal is non-negative. {@link java.util.concurrent.atomic.LongAdder}
 * does not provide this: the two halves of a pair may land in different cells
 * (probe rehash after CAS contention, cell table growth), letting a scan
 * observe the decrement while missing the increment and under-count to a
 * false zero.
 */
final class StripedCounter
{
  /** 16 longs = 128 bytes between slots, to keep them on distinct cache lines. */
  private static final int SPACING = 16;
  private static final int STRIPES = nextPowerOfTwo(Runtime.getRuntime().availableProcessors());
  private final AtomicLongArray counts = new AtomicLongArray(STRIPES * SPACING);
  private static int nextPowerOfTwo(int n)
  {
    int p = 1;
    while (p < n)
    {
      p <<= 1;
    }
    return p;
  }
  private static int slot()
  {
    final long id = Thread.currentThread().getId();
    return (((int) ((id * 0x9E3779B97F4A7C15L) >>> 32)) & (STRIPES - 1)) * SPACING;
  }
  void increment()
  {
    counts.getAndIncrement(slot());
  }
  /** Must be called by the same thread that did the paired {@link #increment()}. */
  void decrement()
  {
    counts.getAndDecrement(slot());
  }
  /**
   * Returns the current count. Concurrent updates may cause over-estimation,
   * but a paired increment/decrement is never observed half-way in the
   * decrement-only direction, so the result is zero only if every access
   * whose increment is visible has completed.
   */
  long sum()
  {
    long s = 0;
    for (int i = 0; i < counts.length(); i += SPACING)
    {
      s += counts.get(i);
    }
    return s;
  }
  /** Resets the count to zero. Only safe when no accesses are in flight. */
  void reset()
  {
    for (int i = 0; i < counts.length(); i += SPACING)
    {
      counts.set(i, 0);
    }
  }
}
opendj-server-legacy/src/test/java/org/opends/server/backends/pluggable/PluggableBackendImplTestCase.java
@@ -31,9 +31,12 @@
import java.nio.charset.StandardCharsets;
import java.util.*;
import java.util.concurrent.Callable;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import com.google.common.io.Resources;
import org.forgerock.opendj.ldap.*;
@@ -629,6 +632,50 @@
    return topEntries.size() + entries.size() + workEntries.size();
  }
  @Test(timeOut = 30000)
  public void testExclusiveLockDrainsSharedAccessDespiteInterrupt() throws Exception
  {
    final EntryContainer ec = backend.getRootContainer().getEntryContainer(testBaseDN);
    final CountDownLatch lockAcquired = new CountDownLatch(1);
    final AtomicBoolean interruptPreserved = new AtomicBoolean();
    ec.beginSharedAccess();
    final Thread exclusiveLocker = new Thread("Test exclusive locker")
    {
      @Override
      public void run()
      {
        ec.lock();
        try
        {
          interruptPreserved.set(Thread.currentThread().isInterrupted());
          lockAcquired.countDown();
        }
        finally
        {
          ec.unlock();
        }
      }
    };
    try
    {
      exclusiveLocker.start();
      exclusiveLocker.interrupt();
      assertFalse(lockAcquired.await(200, TimeUnit.MILLISECONDS),
          "lock() returned while a shared access was still in flight");
    }
    finally
    {
      ec.endSharedAccess();
    }
    assertTrue(lockAcquired.await(10, TimeUnit.SECONDS),
        "lock() did not complete after the shared access ended");
    assertTrue(interruptPreserved.get(), "lock() must preserve the caller's interrupt status");
    exclusiveLocker.join(10000);
  }
  @Test
  public void testHasSubordinates() throws Exception
  {