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")