This is an automated email from the ASF dual-hosted git repository.

yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new cd72eac711 [#13399] fix(core): make SegmentedLock.withGlobalLock 
exclusive against in-flight segment ops (#13400)
cd72eac711 is described below

commit cd72eac7116dcafb56a660476e64bddabf4e5477
Author: YangJie <[email protected]>
AuthorDate: Tue Sep 22 08:07:33 2026 -0400

    [#13399] fix(core): make SegmentedLock.withGlobalLock exclusive against 
in-flight segment ops (#13400)
    
    ### What changes were proposed in this pull request?
    
    Replace the entry latch in `SegmentedLock` with a
    `ReentrantReadWriteLock` gate. Every `withLock` variant now holds the
    read lock across its full critical section (segment-lock acquire plus
    the action), and `withGlobalLock` takes the write lock, so the global
    action waits for all in-flight segment operations to finish before it
    runs. The concurrent-global `IllegalStateException` and `isClearing()`
    semantics are preserved.
    
    ### Why are the changes needed?
    
    `withGlobalLock` promised exclusive access to all segments but acquired
    no segment lock, so a segment operation already inside its critical
    section ran concurrently with the global action. In
    `CaffeineEntityCache`, a `doPut` racing `clear()` could write its index
    entry into the tree `clear()` was swapping out, leaving an entity in
    `cacheData` but absent from the active `cacheIndex` and served stale
    until TTL.
    
    Fix: #13399
    
    ### Does this PR introduce _any_ user-facing change?
    
    No.
    
    ### How was this patch tested?
    
    Added `TestSegmentedLock.testGlobalClearingWaitsForInFlightOperations`,
    a deterministic test that pins that a global action does not proceed
    until an in-flight segment operation has finished. It fails against the
    pre-fix code.
---
 .../org/apache/gravitino/cache/SegmentedLock.java  | 120 ++++++++++++---------
 .../apache/gravitino/cache/TestSegmentedLock.java  |  60 +++++++++++
 2 files changed, 129 insertions(+), 51 deletions(-)

diff --git a/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java 
b/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java
index cfbb8d649e..b6b2d5a8dc 100644
--- a/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java
+++ b/core/src/main/java/org/apache/gravitino/cache/SegmentedLock.java
@@ -21,9 +21,9 @@ package org.apache.gravitino.cache;
 
 import com.google.common.annotations.VisibleForTesting;
 import com.google.common.util.concurrent.Striped;
-import java.util.concurrent.CountDownLatch;
-import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 /**
  * Segmented lock for improved concurrency. Divides locks into segments to 
reduce contention.
@@ -34,8 +34,18 @@ public class SegmentedLock {
 
   private final Striped<Lock> stripedLocks;
 
-  /** CountDownLatch for global operations - null when no operation is in 
progress */
-  private final AtomicReference<CountDownLatch> globalOperationLatch = new 
AtomicReference<>();
+  /**
+   * Gates segment operations against global operations: segment operations 
hold the read lock
+   * across their whole critical section, global operations hold the write 
lock, so a global
+   * operation excludes every segment operation, including ones already in 
flight when it starts.
+   * The lock is fair so that a steady stream of segment operations cannot 
indefinitely starve a
+   * waiting global operation; reentrant read reacquisition is still permitted 
while a writer is
+   * queued, so the nested cache paths that reacquire the read lock do not 
deadlock.
+   */
+  private final ReentrantReadWriteLock globalGate = new 
ReentrantReadWriteLock(true);
+
+  /** True while a global operation is in progress, used to reject concurrent 
global operations. */
+  private final AtomicBoolean clearing = new AtomicBoolean(false);
 
   /**
    * Creates a SegmentedLock with the specified number of segments. Guava's 
Striped automatically
@@ -82,14 +92,19 @@ public class SegmentedLock {
    * @throws RuntimeException if interrupted
    */
   public void withLock(Object key, Runnable action) {
-    waitForGlobalComplete();
-    Lock lock = getSegmentLock(key);
+    Lock readLock = globalGate.readLock();
     try {
-      lock.lockInterruptibly();
+      readLock.lockInterruptibly();
       try {
-        action.run();
+        Lock lock = getSegmentLock(key);
+        lock.lockInterruptibly();
+        try {
+          action.run();
+        } finally {
+          lock.unlock();
+        }
       } finally {
-        lock.unlock();
+        readLock.unlock();
       }
     } catch (InterruptedException e) {
       Thread.currentThread().interrupt();
@@ -108,14 +123,19 @@ public class SegmentedLock {
    * @throws RuntimeException if interrupted
    */
   public <T> T withLock(Object key, java.util.function.Supplier<T> action) {
-    waitForGlobalComplete();
-    Lock lock = getSegmentLock(key);
+    Lock readLock = globalGate.readLock();
     try {
-      lock.lockInterruptibly();
+      readLock.lockInterruptibly();
       try {
-        return action.get();
+        Lock lock = getSegmentLock(key);
+        lock.lockInterruptibly();
+        try {
+          return action.get();
+        } finally {
+          lock.unlock();
+        }
       } finally {
-        lock.unlock();
+        readLock.unlock();
       }
     } catch (InterruptedException e) {
       Thread.currentThread().interrupt();
@@ -136,14 +156,19 @@ public class SegmentedLock {
    */
   public <T, E extends Exception> T withLockAndThrow(
       Object key, EntityCache.ThrowingSupplier<T, E> action) throws E {
-    waitForGlobalComplete();
-    Lock lock = getSegmentLock(key);
+    Lock readLock = globalGate.readLock();
     try {
-      lock.lockInterruptibly();
+      readLock.lockInterruptibly();
       try {
-        return action.get();
+        Lock lock = getSegmentLock(key);
+        lock.lockInterruptibly();
+        try {
+          return action.get();
+        } finally {
+          lock.unlock();
+        }
       } finally {
-        lock.unlock();
+        readLock.unlock();
       }
     } catch (InterruptedException e) {
       Thread.currentThread().interrupt();
@@ -162,14 +187,19 @@ public class SegmentedLock {
    */
   public <E extends Exception> void withLockAndThrow(
       Object key, EntityCache.ThrowingRunnable<E> action) throws E {
-    waitForGlobalComplete();
-    Lock lock = getSegmentLock(key);
+    Lock readLock = globalGate.readLock();
     try {
-      lock.lockInterruptibly();
+      readLock.lockInterruptibly();
       try {
-        action.run();
+        Lock lock = getSegmentLock(key);
+        lock.lockInterruptibly();
+        try {
+          action.run();
+        } finally {
+          lock.unlock();
+        }
       } finally {
-        lock.unlock();
+        readLock.unlock();
       }
     } catch (InterruptedException e) {
       Thread.currentThread().interrupt();
@@ -180,7 +210,7 @@ public class SegmentedLock {
   /** Checks if a global operation is currently in progress. */
   @VisibleForTesting
   public boolean isClearing() {
-    return globalOperationLatch.get() != null;
+    return clearing.get();
   }
 
   /**
@@ -196,42 +226,30 @@ public class SegmentedLock {
    * Executes a global clearing operation with exclusive access to all 
segments. This method sets
    * the clearing flag and ensures no other operations can proceed until the 
clearing is complete.
    *
+   * <p>Exclusivity holds against every segment operation, including ones that 
were already in
+   * flight when this method is called: segment operations hold the {@code 
globalGate} read lock
+   * across their whole critical section, and the write lock acquired here 
waits for all of them.
+   *
+   * <p>Must not be called from inside a {@code withLock} action on the same 
instance: the read lock
+   * cannot upgrade to the write lock, so such a call would deadlock.
+   *
    * @param action The clearing action to execute
    */
   public void withGlobalLock(Runnable action) {
-    // Create a new CountDownLatch for this operation
-    CountDownLatch latch = new CountDownLatch(1);
-
-    // Atomically set the latch, fail if another operation is already in 
progress
-    if (!globalOperationLatch.compareAndSet(null, latch)) {
+    // Mark the global operation in progress, fail if another one is already 
running
+    if (!clearing.compareAndSet(false, true)) {
       throw new IllegalStateException("Global operation already in progress");
     }
 
     try {
-      synchronized (this) {
+      globalGate.writeLock().lock();
+      try {
         action.run();
+      } finally {
+        globalGate.writeLock().unlock();
       }
     } finally {
-      // Clear state first, then signal completion
-      globalOperationLatch.set(null);
-      latch.countDown();
-    }
-  }
-
-  /**
-   * Waits for any ongoing global operation to complete. This method is called 
by regular operations
-   * to ensure they don't interfere with global operations.
-   */
-  private void waitForGlobalComplete() {
-    CountDownLatch latch = globalOperationLatch.get();
-    if (latch != null) {
-      try {
-        latch.await();
-      } catch (InterruptedException e) {
-        Thread.currentThread().interrupt();
-        throw new RuntimeException(
-            "Thread was interrupted while waiting for global operation to 
complete", e);
-      }
+      clearing.set(false);
     }
   }
 }
diff --git 
a/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java 
b/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java
index e377c7751f..edce98096a 100644
--- a/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java
+++ b/core/src/test/java/org/apache/gravitino/cache/TestSegmentedLock.java
@@ -29,6 +29,7 @@ import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.Timeout;
@@ -380,4 +381,63 @@ public class TestSegmentedLock {
               });
         });
   }
+
+  @Test
+  @Timeout(30)
+  void testGlobalClearingWaitsForInFlightOperations() throws 
InterruptedException {
+    SegmentedLock lock = new SegmentedLock(4);
+    CountDownLatch insideSegmentOp = new CountDownLatch(1);
+    CountDownLatch releaseSegmentOp = new CountDownLatch(1);
+    CountDownLatch globalActionDone = new CountDownLatch(1);
+    AtomicBoolean overlapped = new AtomicBoolean(false);
+    AtomicBoolean segmentOpActive = new AtomicBoolean(false);
+
+    Thread segmentOpThread =
+        new Thread(
+            () ->
+                lock.withLock(
+                    "key1",
+                    () -> {
+                      segmentOpActive.set(true);
+                      insideSegmentOp.countDown();
+                      try {
+                        releaseSegmentOp.await();
+                      } catch (InterruptedException e) {
+                        Thread.currentThread().interrupt();
+                      }
+                      segmentOpActive.set(false);
+                    }));
+    segmentOpThread.start();
+    assertTrue(insideSegmentOp.await(5, TimeUnit.SECONDS), "segment operation 
never started");
+
+    Thread globalThread =
+        new Thread(
+            () ->
+                lock.withGlobalLock(
+                    () -> {
+                      // The global action must never overlap the segment 
operation's critical
+                      // section. segmentOpActive stays true from the first to 
the last statement of
+                      // the segment action, and the write lock is granted 
only after the read lock
+                      // is released, so a correct gate always observes it 
false here.
+                      if (segmentOpActive.get()) {
+                        overlapped.set(true);
+                      }
+                      globalActionDone.countDown();
+                    }));
+    globalThread.start();
+
+    // While the segment operation is in flight, the global action must not 
complete:
+    // withGlobalLock promises exclusive access to all segments, so it must 
block
+    // until the in-flight operation releases its critical section. A 
completion
+    // here means the global action ran concurrently with it.
+    assertFalse(
+        globalActionDone.await(1, TimeUnit.SECONDS),
+        "global action ran concurrently with an in-flight segment operation");
+    releaseSegmentOp.countDown();
+    assertTrue(globalActionDone.await(5, TimeUnit.SECONDS), "global action 
never completed");
+    assertFalse(overlapped.get(), "global action overlapped with the in-flight 
segment operation");
+
+    segmentOpThread.join(5000);
+    globalThread.join(5000);
+  }
 }

Reply via email to