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 07328c2e54 [#13363] fix(authz): serialize JCasbin policy reload with 
reads (#13364)
07328c2e54 is described below

commit 07328c2e54cc92590008231593850edf1928a518
Author: jarred0214 <[email protected]>
AuthorDate: Wed Sep 23 16:18:21 2026 +0800

    [#13363] fix(authz): serialize JCasbin policy reload with reads (#13364)
    
    ### What changes were proposed in this pull request?
    
    This PR adds an authorizer-level `ReentrantReadWriteLock` around the
    in-memory JCasbin enforcer state used by `JcasbinAuthorizer`.
    
    The lock makes role policy reload atomically visible to authorization
    reads:
    
    - Authorization reads and policy inspections use the shared read lock.
    - Role policy replacement, invalidation, loaded-role cache cleanup, and
    grouping bind/prune use the write lock.
    - Metadata id resolution and DB queries remain outside the write lock.
    
    It also adds a regression test that simulates concurrent authorization
    requests arriving while a role policy reload is in progress.
    
    ### Why are the changes needed?
    
    `SyncedEnforcer` makes individual JCasbin method calls thread-safe, but
    role policy reload is a multi-call sequence:
    
    ```text
    remove old role policies
    add new role policies
    mark role as loaded
    ```
    
    A concurrent `enforce()` can observe the transient state after the old
    policies are removed and before the new policies are added, causing a
    false denial. This was observed as intermittent `loadTable`
    authorization failures where retrying shortly afterwards succeeded.
    
    This PR ensures that concurrent authorization reads wait for the reload
    sequence to complete instead of reading the intermediate empty policy
    state.
    
    Fixes #13363
    
    ### Does this PR introduce any user-facing change?
    
    No. It does not change authorization semantics or grant/revoke behavior.
    It only makes the existing in-memory policy state transition atomic from
    the perspective of authorization reads.
    
    ### How was this patch tested?
    
    ```bash
    
JAVA_HOME=/Users/hujie3/.gradle/jdks/amazon_com_inc_-17-x86_64-os_x/amazon-corretto-17.jdk/Contents/Home
 ./gradlew :server-common:test --tests 
org.apache.gravitino.server.authorization.jcasbin.TestJcasbinAuthorizer
    ```
    
    Result:
    
    ```text
    BUILD SUCCESSFUL
    ```
---
 .../authorization/jcasbin/JcasbinAuthorizer.java   | 210 ++++++++++++++-------
 .../jcasbin/TestJcasbinAuthorizer.java             | 154 ++++++++++++++-
 2 files changed, 286 insertions(+), 78 deletions(-)

diff --git 
a/server-common/src/main/java/org/apache/gravitino/server/authorization/jcasbin/JcasbinAuthorizer.java
 
b/server-common/src/main/java/org/apache/gravitino/server/authorization/jcasbin/JcasbinAuthorizer.java
index dcf11a9ba7..06f17d4366 100644
--- 
a/server-common/src/main/java/org/apache/gravitino/server/authorization/jcasbin/JcasbinAuthorizer.java
+++ 
b/server-common/src/main/java/org/apache/gravitino/server/authorization/jcasbin/JcasbinAuthorizer.java
@@ -35,7 +35,7 @@ import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.locks.ReentrantLock;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 import org.apache.commons.io.IOUtils;
 import org.apache.commons.lang3.StringUtils;
@@ -139,23 +139,19 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
   private static final long PARTIAL_ROLE_LOAD_RETRY_MS = 10_000L;
 
   /**
-   * Serializes every mutation of role permission policies, including {@link
-   * #invalidateRolePolicies}, {@link #replaceRolePolicies}, and the policy 
writes in {@link
-   * #applyRolePolicies}. Both {@code SyncedEnforcer} calls are individually 
atomic, but the {@code
-   * clear -> re-add} sequence is not, and neither is it ordered against the 
{@link
-   * JcasbinLoadedRolesCache} removal listener, which calls {@link 
#clearRolePoliciesOnCacheRemoval}
-   * from whichever thread happens to drain Caffeine's maintenance queue. 
Without this lock an
-   * eviction firing between another thread's policy writes and its {@link 
#loadedRoles} update
-   * erases the rows that thread just wrote while leaving the marker saying 
they are loaded — a
-   * state no subsequent version check can detect or repair.
+   * Guards every access to the in-memory JCasbin enforcer state. Policy 
reload mutates that state
+   * with a clear-then-add sequence, which must be atomic not only against 
other writers but also
+   * against authorization reads; otherwise a concurrent {@code enforce()} can 
observe the temporary
+   * empty policy set and deny a request that should be allowed.
    *
-   * <p>The lock guards in-memory enforcer mutations only: metadata ids are 
resolved before it is
-   * taken (see {@link #resolveRolePolicies}), so no DB round-trip ever runs 
inside the critical
-   * section. It is reentrant because a {@link #loadedRoles} write performed 
under the lock can
-   * itself trigger an eviction, and therefore a nested {@link 
#clearRolePoliciesOnCacheRemoval}
-   * call.
+   * <p>Writers include {@link #invalidateRolePolicies}, {@link 
#replaceRolePolicies}, {@link
+   * #bindUserRoles}, stale grouping-row pruning, and the {@link 
JcasbinLoadedRolesCache} removal
+   * listener. Readers include all {@code enforce()}, grouping, and 
policy-inspection calls. The
+   * lock only guards in-memory enforcer access: metadata ids are resolved 
before write-lock
+   * acquisition (see {@link #resolveRolePolicies}), so no DB round-trip runs 
inside the critical
+   * section.
    */
-  private final ReentrantLock rolePolicyLock = new ReentrantLock();
+  private final ReentrantReadWriteLock rolePolicyLock = new 
ReentrantReadWriteLock();
 
   /** Jcasbin enforcer is used for metadata authorization. */
   private Enforcer allowEnforcer;
@@ -436,29 +432,35 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
 
     Set<String> privilegeNames = 
privileges.stream().map(Enum::name).collect(Collectors.toSet());
     String userIdStr = String.valueOf(userId);
-    // This is an existence query, not a per-object check: it answers "does 
any deny on these
-    // privileges exist for the user's roles, at any scope?" The standard 
enforce path needs a
-    // concrete metadataId, so reusing it would mean iterating every listed 
object and defeat the
-    // short-circuit. Filtering the deny enforcer's policies by role keeps the 
scan bounded by the
-    // user's role/policy count, never by the number of listed objects. The 
match is intentionally
-    // scope-agnostic (no metadataType filter): a parent-scope deny hides the 
whole subtree and an
-    // object-scope deny hides one object, and both must disable the 
short-circuit.
-    for (String roleId : denyEnforcer.getRolesForUser(userIdStr)) {
-      // getFilteredNamedPolicy returns every "p" row (p = sub, metadataType, 
metadataId, act, eft)
-      // whose field at POLICY_SUBJECT_FIELD_INDEX (sub) equals roleId, i.e. 
all rules carried by
-      // this role. denyEnforcer is a dedicated enforcer that is only ever 
loaded with privileges
-      // whose condition is DENY (see loadPolicyByRoleEntity), so every row 
here represents a deny
-      // regardless of its stored eft string. Each returned row is the list of 
those five fields,
-      // so we read field POLICY_ACTION_FIELD_INDEX (act) to compare the 
denied privilege.
-      for (List<String> policy :
-          denyEnforcer.getFilteredNamedPolicy("p", POLICY_SUBJECT_FIELD_INDEX, 
roleId)) {
-        if (policy.size() > POLICY_ACTION_FIELD_INDEX
-            && privilegeNames.contains(policy.get(POLICY_ACTION_FIELD_INDEX))) 
{
-          return true;
+    rolePolicyLock.readLock().lock();
+    try {
+      // This is an existence query, not a per-object check: it answers "does 
any deny on these
+      // privileges exist for the user's roles, at any scope?" The standard 
enforce path needs a
+      // concrete metadataId, so reusing it would mean iterating every listed 
object and defeat the
+      // short-circuit. Filtering the deny enforcer's policies by role keeps 
the scan bounded by the
+      // user's role/policy count, never by the number of listed objects. The 
match is intentionally
+      // scope-agnostic (no metadataType filter): a parent-scope deny hides 
the whole subtree and an
+      // object-scope deny hides one object, and both must disable the 
short-circuit.
+      for (String roleId : denyEnforcer.getRolesForUser(userIdStr)) {
+        // getFilteredNamedPolicy returns every "p" row (p = sub, 
metadataType, metadataId, act,
+        // eft)
+        // whose field at POLICY_SUBJECT_FIELD_INDEX (sub) equals roleId, i.e. 
all rules carried by
+        // this role. denyEnforcer is a dedicated enforcer that is only ever 
loaded with privileges
+        // whose condition is DENY (see loadPolicyByRoleEntity), so every row 
here represents a deny
+        // regardless of its stored eft string. Each returned row is the list 
of those five fields,
+        // so we read field POLICY_ACTION_FIELD_INDEX (act) to compare the 
denied privilege.
+        for (List<String> policy :
+            denyEnforcer.getFilteredNamedPolicy("p", 
POLICY_SUBJECT_FIELD_INDEX, roleId)) {
+          if (policy.size() > POLICY_ACTION_FIELD_INDEX
+              && 
privilegeNames.contains(policy.get(POLICY_ACTION_FIELD_INDEX))) {
+            return true;
+          }
         }
       }
+      return false;
+    } finally {
+      rolePolicyLock.readLock().unlock();
     }
-    return false;
   }
 
   @Override
@@ -934,19 +936,22 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
       // against just those roles (enforceNarrowed). ALL or an absent header 
falls through to the
       // normal check over every role the caller holds.
       ActiveRoles activeRoles = requestContext.getActiveRoles();
-      if (narrowByActiveRoles && !activeRoles.isAll()) {
-        boolean allowed =
-            enforceNarrowed(
-                userId, metadataType, metadataIdStr, privilege, activeRoles, 
requestContext);
-        if (!allowed) {
-          diagnoseDenial(userId, metadataType, metadataIdStr, privilege);
+      boolean allowed;
+      rolePolicyLock.readLock().lock();
+      try {
+        if (narrowByActiveRoles && !activeRoles.isAll()) {
+          allowed =
+              enforceNarrowed(
+                  userId, metadataType, metadataIdStr, privilege, activeRoles, 
requestContext);
+        } else {
+          allowed =
+              enforcer.enforce(String.valueOf(userId), metadataType, 
metadataIdStr, privilege);
         }
-        return allowed;
+      } finally {
+        rolePolicyLock.readLock().unlock();
       }
 
-      boolean allowed =
-          enforcer.enforce(String.valueOf(userId), metadataType, 
metadataIdStr, privilege);
-      if (!allowed && narrowByActiveRoles) {
+      if (narrowByActiveRoles && !allowed) {
         diagnoseDenial(userId, metadataType, metadataIdStr, privilege);
       }
       return allowed;
@@ -1207,10 +1212,28 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
             desiredRoleIds.add(String.valueOf(id));
           }
           String userIdStr = String.valueOf(userId);
-          for (String currentRole : allowEnforcer.getRolesForUser(userIdStr)) {
-            if (!desiredRoleIds.contains(currentRole)) {
-              allowEnforcer.deleteRoleForUser(userIdStr, currentRole);
-              denyEnforcer.deleteRoleForUser(userIdStr, currentRole);
+          List<String> staleRoleIds = new ArrayList<>();
+          rolePolicyLock.readLock().lock();
+          try {
+            Set<String> currentRoleIds = new 
HashSet<>(allowEnforcer.getRolesForUser(userIdStr));
+            currentRoleIds.addAll(denyEnforcer.getRolesForUser(userIdStr));
+            for (String currentRole : currentRoleIds) {
+              if (!desiredRoleIds.contains(currentRole)) {
+                staleRoleIds.add(currentRole);
+              }
+            }
+          } finally {
+            rolePolicyLock.readLock().unlock();
+          }
+          if (!staleRoleIds.isEmpty()) {
+            rolePolicyLock.writeLock().lock();
+            try {
+              for (String currentRole : staleRoleIds) {
+                allowEnforcer.deleteRoleForUser(userIdStr, currentRole);
+                denyEnforcer.deleteRoleForUser(userIdStr, currentRole);
+              }
+            } finally {
+              rolePolicyLock.writeLock().unlock();
             }
           }
 
@@ -1460,7 +1483,7 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
    */
   private boolean replaceRolePolicies(
       long roleId, long dbUpdatedAt, ResolvedRolePolicies resolved) {
-    rolePolicyLock.lock();
+    rolePolicyLock.writeLock().lock();
     try {
       Optional<Long> latestLoadedAt = loadedRoles.getIfPresent(roleId);
       if (latestLoadedAt.isPresent() && latestLoadedAt.get() >= dbUpdatedAt) {
@@ -1490,13 +1513,13 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
       }
       return true;
     } finally {
-      rolePolicyLock.unlock();
+      rolePolicyLock.writeLock().unlock();
     }
   }
 
   /** Clears a role's policies and removes its loaded marker as one serialized 
operation. */
   private void invalidateRolePolicies(long roleId) {
-    rolePolicyLock.lock();
+    rolePolicyLock.writeLock().lock();
     try {
       // An explicit invalidation is stronger than the retry throttle. Clear 
it under the same lock
       // before removing the loaded marker so the removal callback cannot 
mistake the old partial
@@ -1509,7 +1532,7 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
         clearRolePoliciesWithoutLock(roleId);
       }
     } finally {
-      rolePolicyLock.unlock();
+      rolePolicyLock.writeLock().unlock();
     }
   }
 
@@ -1522,7 +1545,7 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
    * removal event and must not clear the newly installed policies.
    */
   private void clearRolePoliciesOnCacheRemoval(long roleId) {
-    rolePolicyLock.lock();
+    rolePolicyLock.writeLock().lock();
     try {
       if (loadedRoles.getIfPresent(roleId).isPresent()
           || partialRoleLoadBackoff.getIfPresent(roleId).isPresent()) {
@@ -1533,7 +1556,7 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
       }
       clearRolePoliciesWithoutLock(roleId);
     } finally {
-      rolePolicyLock.unlock();
+      rolePolicyLock.writeLock().unlock();
     }
   }
 
@@ -1544,9 +1567,38 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
   }
 
   private void bindUserRoles(long userId, List<Long> roleIds) {
-    for (Long roleId : roleIds) {
-      allowEnforcer.addRoleForUser(String.valueOf(userId), 
String.valueOf(roleId));
-      denyEnforcer.addRoleForUser(String.valueOf(userId), 
String.valueOf(roleId));
+    if (roleIds.isEmpty()) {
+      return;
+    }
+
+    String userIdStr = String.valueOf(userId);
+    List<Long> missingRoleIds = new ArrayList<>();
+    rolePolicyLock.readLock().lock();
+    try {
+      Set<String> allowRoleIds = new 
HashSet<>(allowEnforcer.getRolesForUser(userIdStr));
+      Set<String> denyRoleIds = new 
HashSet<>(denyEnforcer.getRolesForUser(userIdStr));
+      for (Long roleId : roleIds) {
+        String roleIdStr = String.valueOf(roleId);
+        if (!allowRoleIds.contains(roleIdStr) || 
!denyRoleIds.contains(roleIdStr)) {
+          missingRoleIds.add(roleId);
+        }
+      }
+    } finally {
+      rolePolicyLock.readLock().unlock();
+    }
+    if (missingRoleIds.isEmpty()) {
+      return;
+    }
+
+    rolePolicyLock.writeLock().lock();
+    try {
+      for (Long roleId : missingRoleIds) {
+        String roleIdStr = String.valueOf(roleId);
+        allowEnforcer.addRoleForUser(userIdStr, roleIdStr);
+        denyEnforcer.addRoleForUser(userIdStr, roleIdStr);
+      }
+    } finally {
+      rolePolicyLock.writeLock().unlock();
     }
   }
 
@@ -1642,22 +1694,36 @@ public class JcasbinAuthorizer implements 
GravitinoAuthorizer {
     }
     try {
       String userIdStr = String.valueOf(userId);
-      List<String> boundRoles = allowEnforcer.getRolesForUser(userIdStr);
-      if (boundRoles.isEmpty()) {
-        LOG.debug(
-            "Denied [{}, {}, {}, {}]: no role is bound to the user in the 
allow enforcer",
-            userIdStr,
-            metadataType,
-            metadataIdStr,
-            privilege);
-        return;
+      Map<String, Integer> rolePolicyCounts = new HashMap<>();
+      rolePolicyLock.readLock().lock();
+      try {
+        List<String> boundRoles = allowEnforcer.getRolesForUser(userIdStr);
+        if (boundRoles.isEmpty()) {
+          LOG.debug(
+              "Denied [{}, {}, {}, {}]: no role is bound to the user in the 
allow enforcer",
+              userIdStr,
+              metadataType,
+              metadataIdStr,
+              privilege);
+          return;
+        }
+
+        for (String roleIdStr : boundRoles) {
+          int policyCount =
+              allowEnforcer
+                  .getFilteredNamedPolicy("p", POLICY_SUBJECT_FIELD_INDEX, 
roleIdStr)
+                  .size();
+          rolePolicyCounts.put(roleIdStr, policyCount);
+        }
+      } finally {
+        rolePolicyLock.readLock().unlock();
       }
 
       List<String> rolesWithoutPolicies = new ArrayList<>();
-      List<String> roleStates = new ArrayList<>(boundRoles.size());
-      for (String roleIdStr : boundRoles) {
-        int policyCount =
-            allowEnforcer.getFilteredNamedPolicy("p", 
POLICY_SUBJECT_FIELD_INDEX, roleIdStr).size();
+      List<String> roleStates = new ArrayList<>(rolePolicyCounts.size());
+      for (Map.Entry<String, Integer> roleState : rolePolicyCounts.entrySet()) 
{
+        String roleIdStr = roleState.getKey();
+        int policyCount = roleState.getValue();
         Optional<Long> loadedAt = 
loadedRoles.getIfPresent(Long.parseLong(roleIdStr));
         roleStates.add(
             roleIdStr
diff --git 
a/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinAuthorizer.java
 
b/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinAuthorizer.java
index f47e925bdd..06d5e2d4a2 100644
--- 
a/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinAuthorizer.java
+++ 
b/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinAuthorizer.java
@@ -25,6 +25,7 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyInt;
@@ -46,6 +47,7 @@ import java.io.IOException;
 import java.lang.reflect.Field;
 import java.lang.reflect.Method;
 import java.security.Principal;
+import java.time.Duration;
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.HashMap;
@@ -55,10 +57,13 @@ import java.util.Objects;
 import java.util.Optional;
 import java.util.Set;
 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.AtomicLong;
 import java.util.concurrent.atomic.AtomicReference;
-import java.util.concurrent.locks.ReentrantLock;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.function.Function;
 import java.util.stream.Collectors;
 import org.apache.gravitino.Entity;
@@ -705,7 +710,7 @@ public class TestJcasbinAuthorizer {
   public void testStaleRemovalDoesNotClearReloadedPolicies() throws Exception {
     Enforcer allowEnforcer = getAllowEnforcer(jcasbinAuthorizer);
     GravitinoCache<Long, Long> loadedRoles = 
getLoadedRolesCache(jcasbinAuthorizer);
-    ReentrantLock rolePolicyLock = getRolePolicyLock(jcasbinAuthorizer);
+    ReentrantReadWriteLock rolePolicyLock = 
getRolePolicyLock(jcasbinAuthorizer);
     String roleIdStr = String.valueOf(ALLOW_ROLE_ID);
     String[] policyRow =
         new String[] {
@@ -720,7 +725,7 @@ public class TestJcasbinAuthorizer {
 
     CountDownLatch invalidationStarted = new CountDownLatch(1);
     AtomicReference<Throwable> failure = new AtomicReference<>();
-    rolePolicyLock.lock();
+    rolePolicyLock.writeLock().lock();
     Thread invalidator =
         new Thread(
             () -> {
@@ -751,7 +756,7 @@ public class TestJcasbinAuthorizer {
       allowEnforcer.addPolicy(policyRow);
       loadedRoles.put(ALLOW_ROLE_ID, 2L);
     } finally {
-      rolePolicyLock.unlock();
+      rolePolicyLock.writeLock().unlock();
     }
 
     invalidator.join(5000L);
@@ -823,6 +828,106 @@ public class TestJcasbinAuthorizer {
         "the completed load's policies must survive a late partial result");
   }
 
+  @Test
+  public void testConcurrentAuthorizationWaitsForPolicyReload() throws 
Exception {
+    Principal currentPrincipal = PrincipalUtils.getCurrentPrincipal();
+    MetadataObject catalog = MetadataObjects.of(null, "testCatalog", 
MetadataObject.Type.CATALOG);
+    RoleEntity allowRole =
+        mockRoleInStore(ALLOW_ROLE_ID, "allowRole", 
ImmutableList.of(getAllowSecurableObject()));
+    mockDirectUserRoles(allowRole);
+
+    Enforcer allowEnforcer = getAllowEnforcer(jcasbinAuthorizer);
+    String roleIdStr = String.valueOf(ALLOW_ROLE_ID);
+
+    assertTrue(
+        jcasbinAuthorizer.authorize(
+            currentPrincipal, METALAKE, catalog, USE_CATALOG, new 
AuthorizationRequestContext()));
+    List<List<String>> policyRows = allowEnforcer.getFilteredPolicy(0, 
roleIdStr);
+    assertFalse(policyRows.isEmpty());
+
+    ReentrantReadWriteLock rolePolicyLock = 
getRolePolicyLock(jcasbinAuthorizer);
+    rolePolicyLock.writeLock().lock();
+    ExecutorService executor = Executors.newFixedThreadPool(16);
+    List<Future<Boolean>> futures = new ArrayList<>();
+    try {
+      allowEnforcer.removeFilteredPolicy(0, roleIdStr);
+      CountDownLatch start = new CountDownLatch(1);
+      CountDownLatch started = new CountDownLatch(16);
+      for (int i = 0; i < 16; i++) {
+        futures.add(
+            executor.submit(
+                () -> {
+                  assertTrue(start.await(5, TimeUnit.SECONDS));
+                  started.countDown();
+                  return invokeAuthorizeByJcasbin(
+                      jcasbinAuthorizer,
+                      USER_ID,
+                      METALAKE,
+                      catalog,
+                      CATALOG_ID,
+                      USE_CATALOG,
+                      new AuthorizationRequestContext());
+                }));
+      }
+
+      start.countDown();
+      assertTrue(started.await(5, TimeUnit.SECONDS));
+      for (List<String> policyRow : policyRows) {
+        allowEnforcer.addPolicy(policyRow);
+      }
+    } finally {
+      rolePolicyLock.writeLock().unlock();
+    }
+
+    try {
+      for (Future<Boolean> future : futures) {
+        assertTrue(future.get(5, TimeUnit.SECONDS));
+      }
+    } finally {
+      executor.shutdown();
+    }
+    assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+  }
+
+  @Test
+  public void 
testLoadedRoleExpirationCleanerCanTakeWriteLockAfterDiagnosticRead()
+      throws Exception {
+    ReentrantReadWriteLock rolePolicyLock = 
getRolePolicyLock(jcasbinAuthorizer);
+    JcasbinLoadedRolesCache expiringLoadedRoles =
+        new JcasbinLoadedRolesCache(
+            1,
+            10,
+            roleId -> {
+              rolePolicyLock.writeLock().lock();
+              try {
+                // Simulate the authorizer cleanup path entered by Caffeine's 
synchronous removal
+                // listener. The test fails by timeout if a caller still holds 
readLock while
+                // touching loadedRoles.
+              } finally {
+                rolePolicyLock.writeLock().unlock();
+              }
+            });
+
+    expiringLoadedRoles.put(ALLOW_ROLE_ID, 1L);
+    Thread.sleep(10L);
+
+    assertTimeoutPreemptively(
+        Duration.ofSeconds(2),
+        () -> {
+          Map<String, Integer> rolePolicyCounts = new HashMap<>();
+          rolePolicyLock.readLock().lock();
+          try {
+            rolePolicyCounts.put(String.valueOf(ALLOW_ROLE_ID), 0);
+          } finally {
+            rolePolicyLock.readLock().unlock();
+          }
+
+          for (String roleIdStr : rolePolicyCounts.keySet()) {
+            expiringLoadedRoles.getIfPresent(Long.parseLong(roleIdStr));
+          }
+        });
+  }
+
   /** Reflectively invoke the private versionCheckAndLoadRoles. */
   private static void invokeVersionCheckAndLoadRoles(
       JcasbinAuthorizer authorizer,
@@ -840,6 +945,42 @@ public class TestJcasbinAuthorizer {
     m.invoke(authorizer, metalake, roleIds, requestContext);
   }
 
+  /** Reflectively invoke the allow-side JCasbin authorization step. */
+  private static boolean invokeAuthorizeByJcasbin(
+      JcasbinAuthorizer authorizer,
+      long userId,
+      String metalake,
+      MetadataObject metadataObject,
+      Long metadataId,
+      Privilege.Name privilege,
+      AuthorizationRequestContext requestContext)
+      throws Exception {
+    Field field = 
JcasbinAuthorizer.class.getDeclaredField("allowInternalAuthorizer");
+    field.setAccessible(true);
+    Object allowInternalAuthorizer = field.get(authorizer);
+    Method method =
+        allowInternalAuthorizer
+            .getClass()
+            .getDeclaredMethod(
+                "authorizeByJcasbin",
+                long.class,
+                String.class,
+                MetadataObject.class,
+                Long.class,
+                String.class,
+                AuthorizationRequestContext.class);
+    method.setAccessible(true);
+    return (Boolean)
+        method.invoke(
+            allowInternalAuthorizer,
+            userId,
+            metalake,
+            metadataObject,
+            metadataId,
+            privilege.name(),
+            requestContext);
+  }
+
   private static boolean invokeReplaceRolePolicies(
       JcasbinAuthorizer authorizer,
       long roleId,
@@ -2653,10 +2794,11 @@ public class TestJcasbinAuthorizer {
     return (GravitinoCache<Long, Long>) field.get(authorizer);
   }
 
-  private static ReentrantLock getRolePolicyLock(JcasbinAuthorizer authorizer) 
throws Exception {
+  private static ReentrantReadWriteLock getRolePolicyLock(JcasbinAuthorizer 
authorizer)
+      throws Exception {
     Field field = JcasbinAuthorizer.class.getDeclaredField("rolePolicyLock");
     field.setAccessible(true);
-    return (ReentrantLock) field.get(authorizer);
+    return (ReentrantReadWriteLock) field.get(authorizer);
   }
 
   @SuppressWarnings("unchecked")

Reply via email to