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 077cab445e [#13001] fix(core): Fence role membership writes against
concurrent deletion (#13006)
077cab445e is described below
commit 077cab445e41149621faaff5947f9070e83b78f6
Author: Qi Yu <[email protected]>
AuthorDate: Thu Sep 24 07:32:44 2026 +0800
[#13001] fix(core): Fence role membership writes against concurrent
deletion (#13006)
### What changes were proposed in this pull request?
Fence the observed Metalake and newly referenced Role IDs in the same
transaction as User/Group membership writes. Lock the principal before
Roles and acquire Role locks in stable ID order. Use shared locks on
MySQL/PostgreSQL and the existing exclusive-lock fallback on H2. Cover
update, insert, and overwrite.
### Why are the changes needed?
A principal CAS does not detect Role deletion. Role or Metalake deletion
can finish membership cleanup before a concurrent grant inserts its
relation, leaving an active relation to a deleted endpoint. Metalake
fencing must happen before the principal write because its cascade
cleans membership rows first.
Fix: #13001
### Does this PR introduce _any_ user-facing change?
Grants referencing a deleted, missing, or foreign-metalake Role fail
through the existing IllegalRoleException path instead of writing an
invalid membership. Stale IDs never retarget same-name replacements.
Existing grant/revoke idempotency and overwrite behavior remain; no API
signatures, configuration, or storage formats change. Historical orphan
cleanup is tracked separately in #13003.
### How was this patch tested?
All 27 TestRoleMembershipWrites cases across H2, MySQL, and PostgreSQL
and all 8 authorization tests passed. Coverage includes both deletion
orders, rollback for invalid batches, missing/deleted/foreign/recreated
Roles, insert/overwrite, idempotency, shared-lock concurrency, and
Metalake cascades. Core compilation, Spotless, and core check passed.
All 65 final H2/authorization regression tests passed across the User,
Group, Role, membership-write, and permission-manager suites. git diff
--check passed. No full-repository test run was performed.
---------
Co-authored-by: Jerry Shao <[email protected]>
---
.../gravitino/authorization/PermissionManager.java | 84 ++-
.../storage/relational/mapper/RoleMetaMapper.java | 12 +
.../mapper/RoleMetaSQLProviderFactory.java | 13 +-
.../provider/base/RoleMetaBaseSQLProvider.java | 16 +-
.../postgresql/RoleMetaPostgreSQLProvider.java | 5 +
.../relational/service/GroupMetaService.java | 53 +-
.../relational/service/RoleMetaService.java | 34 +
.../relational/service/UserMetaService.java | 53 +-
.../TestAccessControlManagerForPermissions.java | 176 +++++
.../service/TestRoleMembershipWrites.java | 761 +++++++++++++++++++++
10 files changed, 1096 insertions(+), 111 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/authorization/PermissionManager.java
b/core/src/main/java/org/apache/gravitino/authorization/PermissionManager.java
index fadb3463bc..bc11f46c52 100644
---
a/core/src/main/java/org/apache/gravitino/authorization/PermissionManager.java
+++
b/core/src/main/java/org/apache/gravitino/authorization/PermissionManager.java
@@ -33,6 +33,7 @@ import java.util.Map;
import java.util.Set;
import java.util.function.Function;
import java.util.stream.Collectors;
+import javax.annotation.Nullable;
import org.apache.commons.lang3.StringUtils;
import org.apache.gravitino.Entity;
import org.apache.gravitino.EntityStore;
@@ -84,14 +85,9 @@ class PermissionManager {
UserEntity.class,
Entity.EntityType.USER,
userEntity -> {
- List<RoleEntity> roleEntities = Lists.newArrayList();
- if (userEntity.roleNames() != null) {
- for (String role : userEntity.roleNames()) {
- roleEntities.add(roleManager.getRole(metalake, role));
- }
- }
- List<String> roleNames =
Lists.newArrayList(toRoleNames(roleEntities));
- List<Long> roleIds =
Lists.newArrayList(toRoleIds(roleEntities));
+ checkObservedRoles(userEntity.roleNames(),
userEntity.roleIds());
+ List<String> roleNames = mutableCopy(userEntity.roleNames());
+ List<Long> roleIds = mutableCopy(userEntity.roleIds());
for (RoleEntity roleEntityToGrant : roleEntitiesToGrant) {
if (roleIds.contains(roleEntityToGrant.id())) {
@@ -173,14 +169,9 @@ class PermissionManager {
GroupEntity.class,
Entity.EntityType.GROUP,
groupEntity -> {
- List<RoleEntity> roleEntities = Lists.newArrayList();
- if (groupEntity.roleNames() != null) {
- for (String role : groupEntity.roleNames()) {
- roleEntities.add(roleManager.getRole(metalake, role));
- }
- }
- List<String> roleNames =
Lists.newArrayList(toRoleNames(roleEntities));
- List<Long> roleIds =
Lists.newArrayList(toRoleIds(roleEntities));
+ checkObservedRoles(groupEntity.roleNames(),
groupEntity.roleIds());
+ List<String> roleNames = mutableCopy(groupEntity.roleNames());
+ List<Long> roleIds = mutableCopy(groupEntity.roleIds());
for (RoleEntity roleEntityToGrant : roleEntitiesToGrant) {
if (roleIds.contains(roleEntityToGrant.id())) {
@@ -262,19 +253,16 @@ class PermissionManager {
GroupEntity.class,
Entity.EntityType.GROUP,
groupEntity -> {
- List<RoleEntity> roleEntities = Lists.newArrayList();
- if (groupEntity.roleNames() != null) {
- for (String role : groupEntity.roleNames()) {
- roleEntities.add(roleManager.getRole(metalake, role));
- }
- }
- List<String> roleNames =
Lists.newArrayList(toRoleNames(roleEntities));
- List<Long> roleIds =
Lists.newArrayList(toRoleIds(roleEntities));
+ checkObservedRoles(groupEntity.roleNames(),
groupEntity.roleIds());
+ List<String> roleNames = mutableCopy(groupEntity.roleNames());
+ List<Long> roleIds = mutableCopy(groupEntity.roleIds());
for (RoleEntity roleEntityToRevoke : roleEntitiesToRevoke) {
- roleNames.remove(roleEntityToRevoke.name());
- boolean removed = roleIds.remove(roleEntityToRevoke.id());
- if (!removed) {
+ int index = roleIds.indexOf(roleEntityToRevoke.id());
+ if (index >= 0) {
+ roleNames.remove(index);
+ roleIds.remove(index);
+ } else {
LOG.warn(
"Failed to revoke, role {} does not exist in the group
{} of metalake {}",
roleEntityToRevoke.name(),
@@ -351,20 +339,16 @@ class PermissionManager {
UserEntity.class,
Entity.EntityType.USER,
userEntity -> {
- List<RoleEntity> roleEntities = Lists.newArrayList();
- if (userEntity.roleNames() != null) {
- for (String role : userEntity.roleNames()) {
- roleEntities.add(roleManager.getRole(metalake, role));
- }
- }
-
- List<String> roleNames =
Lists.newArrayList(toRoleNames(roleEntities));
- List<Long> roleIds =
Lists.newArrayList(toRoleIds(roleEntities));
+ checkObservedRoles(userEntity.roleNames(),
userEntity.roleIds());
+ List<String> roleNames = mutableCopy(userEntity.roleNames());
+ List<Long> roleIds = mutableCopy(userEntity.roleIds());
for (RoleEntity roleEntityToRevoke : roleEntitiesToRevoke) {
- roleNames.remove(roleEntityToRevoke.name());
- boolean removed = roleIds.remove(roleEntityToRevoke.id());
- if (!removed) {
+ int index = roleIds.indexOf(roleEntityToRevoke.id());
+ if (index >= 0) {
+ roleNames.remove(index);
+ roleIds.remove(index);
+ } else {
LOG.warn(
"Failed to revoke, role {} doesn't exist in the user
{} of metalake {}",
roleEntityToRevoke.name(),
@@ -870,11 +854,25 @@ class PermissionManager {
}
}
- private List<Long> toRoleIds(List<RoleEntity> roleEntities) {
- return
roleEntities.stream().map(RoleEntity::id).collect(Collectors.toList());
+ // The principal handed to the updater already carries its memberships as a
(name, ID) pair per
+ // role, read from the membership join that UserMetaService#updateUser and
+ // GroupMetaService#updateGroup run inside the update transaction. Resolving
those names again
+ // here would cost one query per existing role inside that transaction, and
it would resolve them
+ // by name, which is what lets a deleted-and-recreated role take an observed
membership over. The
+ // observed IDs are carried forward untouched instead, so a replacement
never inherits a grant and
+ // the transaction issues no extra reads. Only the pairing itself still
needs checking, because
+ // role IDs are an optional entity field and a name alone cannot prove
membership identity.
+ private static void checkObservedRoles(
+ @Nullable List<String> roleNames, @Nullable List<Long> roleIds) {
+ if (roleNames == null || roleNames.isEmpty()) {
+ return;
+ }
+ if (roleIds == null || roleNames.size() != roleIds.size()) {
+ throw new IllegalRoleException("Existing role names and IDs must be
paired");
+ }
}
- private List<String> toRoleNames(List<RoleEntity> roleEntities) {
- return
roleEntities.stream().map(RoleEntity::name).collect(Collectors.toList());
+ private static <T> List<T> mutableCopy(@Nullable List<T> values) {
+ return values == null ? Lists.newArrayList() : Lists.newArrayList(values);
}
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaMapper.java
index 1f03e5df6d..aca5c0d8d2 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaMapper.java
@@ -20,6 +20,7 @@
package org.apache.gravitino.storage.relational.mapper;
import java.util.List;
+import javax.annotation.Nullable;
import org.apache.gravitino.storage.relational.po.RolePO;
import org.apache.gravitino.storage.relational.po.auth.RoleUpdatedAt;
import org.apache.ibatis.annotations.DeleteProvider;
@@ -47,6 +48,17 @@ public interface RoleMetaMapper {
RolePO selectRoleMetaByMetalakeIdAndName(
@Param("metalakeId") Long metalakeId, @Param("roleName") String
roleName);
+ /**
+ * Returns an active role by ID and holds its lock for the current
transaction.
+ *
+ * <p>The lock is shared on MySQL/PostgreSQL and exclusive on H2.
+ *
+ * @return the active role, or null if it does not exist
+ */
+ @Nullable
+ @SelectProvider(type = RoleMetaSQLProviderFactory.class, method =
"selectRoleMetaByIdForShare")
+ RolePO selectRoleMetaByIdForShare(@Param("roleId") Long roleId);
+
/** Returns and locks an active role by ID for the current transaction. */
@SelectProvider(type = RoleMetaSQLProviderFactory.class, method =
"selectRoleMetaByIdForUpdate")
RolePO selectRoleMetaByIdForUpdate(@Param("roleId") Long roleId);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaSQLProviderFactory.java
index 3cc763e1ba..4d990b09ea 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/RoleMetaSQLProviderFactory.java
@@ -48,13 +48,24 @@ public class RoleMetaSQLProviderFactory {
static class RoleMetaMySQLProvider extends RoleMetaBaseSQLProvider {}
- static class RoleMetaH2Provider extends RoleMetaBaseSQLProvider {}
+ static class RoleMetaH2Provider extends RoleMetaBaseSQLProvider {
+ @Override
+ public String selectRoleMetaByIdForShare(Long roleId) {
+ // H2 has no shared row-lock syntax, matching the other parent-fencing
providers.
+ return selectRoleMetaByIdForUpdate(roleId);
+ }
+ }
public static String selectRoleMetaByMetalakeIdAndName(
@Param("metalakeId") Long metalakeId, @Param("roleName") String
roleName) {
return getProvider().selectRoleMetaByMetalakeIdAndName(metalakeId,
roleName);
}
+ /** Returns SQL that selects an active role by ID and locks it for shared
access. */
+ public static String selectRoleMetaByIdForShare(@Param("roleId") Long
roleId) {
+ return getProvider().selectRoleMetaByIdForShare(roleId);
+ }
+
/** Returns SQL that selects and locks an active role by ID. */
public static String selectRoleMetaByIdForUpdate(@Param("roleId") Long
roleId) {
return getProvider().selectRoleMetaByIdForUpdate(roleId);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/RoleMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/RoleMetaBaseSQLProvider.java
index e20fdfa752..17083abc43 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/RoleMetaBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/RoleMetaBaseSQLProvider.java
@@ -43,14 +43,24 @@ public class RoleMetaBaseSQLProvider {
+ " AND deleted_at = 0";
}
- /** Returns SQL that selects and locks an active role by ID. */
- public String selectRoleMetaByIdForUpdate(@Param("roleId") Long roleId) {
+ /** Returns SQL that selects an active role by ID. */
+ protected String selectRoleMetaById(Long roleId) {
return "SELECT role_id as roleId, role_name as roleName, metalake_id as
metalakeId,"
+ " properties, audit_info as auditInfo, current_version as
currentVersion,"
+ " last_version as lastVersion, deleted_at as deletedAt"
+ " FROM "
+ ROLE_TABLE_NAME
- + " WHERE role_id = #{roleId} AND deleted_at = 0 FOR UPDATE";
+ + " WHERE role_id = #{roleId} AND deleted_at = 0";
+ }
+
+ /** Returns SQL that selects and locks an active role by ID. */
+ public String selectRoleMetaByIdForUpdate(@Param("roleId") Long roleId) {
+ return selectRoleMetaById(roleId) + " FOR UPDATE";
+ }
+
+ /** Returns SQL that selects an active role by ID and locks it for shared
access. */
+ public String selectRoleMetaByIdForShare(@Param("roleId") Long roleId) {
+ return selectRoleMetaById(roleId) + " LOCK IN SHARE MODE";
}
public String selectRoleIdByMetalakeIdAndName(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/RoleMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/RoleMetaPostgreSQLProvider.java
index bdfb9253e8..534134bd6c 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/RoleMetaPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/RoleMetaPostgreSQLProvider.java
@@ -26,6 +26,11 @@ import org.apache.gravitino.storage.relational.po.RolePO;
import org.apache.ibatis.annotations.Param;
public class RoleMetaPostgreSQLProvider extends RoleMetaBaseSQLProvider {
+ @Override
+ public String selectRoleMetaByIdForShare(Long roleId) {
+ return selectRoleMetaById(roleId) + " FOR SHARE";
+ }
+
@Override
public String softDeleteRoleMetaByRoleId(
@Param("roleId") Long roleId, @Param("currentVersion") Long
currentVersion) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/GroupMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/GroupMetaService.java
index 028fe024cf..44087c7366 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/GroupMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/GroupMetaService.java
@@ -193,7 +193,10 @@ public class GroupMetaService {
POConverters.initializeGroupRoleRelsPOWithVersion(groupEntity,
roleIds);
SessionUtils.doMultipleWithCommit(
- () -> lockMetalakeForGroupCreate(metalakePO),
+ () ->
+ MetalakeMetaService.getInstance()
+ .lockMetalakeForChildWrite(
+ metalakePO.getMetalakeName(),
metalakePO.getMetalakeId()),
() ->
SessionUtils.doWithoutCommit(
GroupMetaMapper.class,
@@ -204,6 +207,9 @@ public class GroupMetaService {
mapper.insertGroupMeta(groupPO);
}
}),
+ () ->
+ RoleMetaService.getInstance()
+ .lockRolesForMembership(metalakePO.getMetalakeId(), roleIds),
() -> {
SessionUtils.doWithoutCommit(
GroupRoleRelMapper.class,
@@ -303,6 +309,17 @@ public class GroupMetaService {
// metadata-only change, such as the audit info, still has to be written.
try {
SessionUtils.doMultipleWithCommit(
+ () -> {
+ if (!insertRoleIds.isEmpty() || !deleteRoleIds.isEmpty()) {
+ // The cascade writes memberships before principals; this update
does the reverse.
+ // Fence grants and revokes before the principal CAS to avoid
both orphan grants and
+ // a revoke/cascade deadlock. Metadata-only updates write no
membership rows, so they
+ // need no parent lock (which would serialize unrelated updates
on H2).
+ MetalakeMetaService.getInstance()
+ .lockMetalakeForChildWrite(
+ identifier.namespace().level(0),
oldGroupPO.getMetalakeId());
+ }
+ },
() -> {
int updated =
SessionUtils.getWithoutCommit(
@@ -315,6 +332,9 @@ public class GroupMetaService {
throw groupWriteFailure(identifier, oldGroupPO,
GroupLookup.NAME);
}
},
+ () ->
+ RoleMetaService.getInstance()
+ .lockRolesForMembership(oldGroupPO.getMetalakeId(),
insertRoleIds),
() -> {
if (insertRoleIds.isEmpty()) {
return;
@@ -443,37 +463,6 @@ public class GroupMetaService {
return new PagedResult<>(totalCount, groups);
}
- /**
- * Holds the parent metalake row for the rest of the transaction, so the
group cannot be created
- * under a metalake that is going away.
- *
- * <p>The lock is shared, not exclusive: many groups can be created under
the same metalake at the
- * same time. Dropping a metalake takes an exclusive lock on this row, so a
drop and a create
- * cannot overlap. Whoever gets the row first wins, and the loser either
sees the metalake gone or
- * inserts under a metalake that is still there.
- *
- * <p>The name is compared again because the ID alone cannot tell a rename
apart: the caller
- * looked the metalake up by name, so a renamed row means the name in the
request no longer
- * exists.
- *
- * <p>The metalake's version is deliberately not compared, matching {@code
CatalogMetaService}.
- * Holding the row is what makes the create safe. An unrelated metalake edit
that commits in
- * between bumps the version without making this create wrong, so comparing
it would reject the
- * create for no reason.
- */
- private void lockMetalakeForGroupCreate(MetalakePO observedMetalakePO) {
- OccWriteSupport.lockParentForChildWrite(
- observedMetalakePO.getMetalakeName(),
- Entity.EntityType.METALAKE,
- () ->
- SessionUtils.getWithoutCommit(
- MetalakeMetaMapper.class,
- mapper ->
-
mapper.selectMetalakeMetaByIdForShare(observedMetalakePO.getMetalakeId())),
- null,
- current -> Objects.equals(current.getMetalakeName(),
observedMetalakePO.getMetalakeName()));
- }
-
private RuntimeException groupWriteFailure(
NameIdentifier identifier, GroupPO observedGroupPO, GroupLookup lookup) {
// Sessions run at READ_COMMITTED, so a plain read would already see the
latest committed row.
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/RoleMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/RoleMetaService.java
index 08d69202db..769474c192 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/RoleMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/RoleMetaService.java
@@ -24,6 +24,7 @@ import com.google.common.base.Preconditions;
import com.google.common.collect.Lists;
import com.google.common.collect.Sets;
import java.io.IOException;
+import java.util.Collection;
import java.util.Collections;
import java.util.HashSet;
import java.util.List;
@@ -31,6 +32,7 @@ import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
+import java.util.TreeSet;
import java.util.function.Function;
import java.util.stream.Collectors;
import org.apache.gravitino.Entity;
@@ -41,6 +43,7 @@ import org.apache.gravitino.Namespace;
import org.apache.gravitino.authorization.AuthorizationUtils;
import org.apache.gravitino.authorization.SecurableObject;
import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.exceptions.NoSuchRoleException;
import org.apache.gravitino.meta.RoleEntity;
import org.apache.gravitino.meta.UserEntity;
import org.apache.gravitino.metrics.Monitored;
@@ -291,6 +294,37 @@ public class RoleMetaService {
}
}
+ /**
+ * Fences newly referenced roles until the surrounding membership
transaction commits.
+ *
+ * <p>Existing and removed memberships do not need role locks: deletion can
clean existing rows,
+ * and a revoke cannot leave a new relation behind. Only lock the added IDs
to keep the number of
+ * locking reads proportional to the grant, not the principal's full set of
roles.
+ *
+ * <p>Call after writing the principal row and before modifying any
membership rows. This keeps
+ * the principal-before-role order used by metalake cascades. The caller
must also fence the
+ * metalake before the principal write. Shared locks permit independent
grants of the same role
+ * while excluding its deletion; H2 uses exclusive locks instead. Roles are
locked by stable ID in
+ * ascending order, never re-resolved by a reusable name.
+ *
+ * @throws IllegalStateException if called outside a transaction
+ */
+ void lockRolesForMembership(Long metalakeId, Collection<Long> roleIds) {
+ Preconditions.checkState(
+ SessionUtils.isInTransaction(), "Role membership locks require an
active transaction");
+ for (Long roleId : new TreeSet<>(roleIds)) {
+ RolePO role =
+ SessionUtils.getWithoutCommit(
+ RoleMetaMapper.class, mapper ->
mapper.selectRoleMetaByIdForShare(roleId));
+ if (role == null || !Objects.equals(role.getMetalakeId(), metalakeId)) {
+ // PermissionManager maps a missing role to IllegalRoleException. A
generic missing-entity
+ // exception would incorrectly report the principal as missing instead.
+ throw new NoSuchRoleException(
+ "Role with ID %s does not exist in metalake with ID %s", roleId,
metalakeId);
+ }
+ }
+ }
+
private List<SecurableObjectPO> toSecurableObjectPOs(
Set<SecurableObject> deleteObjects, RoleEntity oldRoleEntity, String
metalake) {
List<SecurableObjectPO> securableObjectPOs = Lists.newArrayList();
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/UserMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/UserMetaService.java
index b4fce90d25..b6c2fcd3f1 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/UserMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/UserMetaService.java
@@ -153,7 +153,10 @@ public class UserMetaService {
POConverters.initializeUserRoleRelsPOWithVersion(userEntity,
roleIds);
SessionUtils.doMultipleWithCommit(
- () -> lockMetalakeForUserCreate(metalakePO),
+ () ->
+ MetalakeMetaService.getInstance()
+ .lockMetalakeForChildWrite(
+ metalakePO.getMetalakeName(),
metalakePO.getMetalakeId()),
() ->
SessionUtils.doWithoutCommit(
UserMetaMapper.class,
@@ -164,6 +167,9 @@ public class UserMetaService {
mapper.insertUserMeta(userPO);
}
}),
+ () ->
+ RoleMetaService.getInstance()
+ .lockRolesForMembership(metalakePO.getMetalakeId(), roleIds),
() -> {
SessionUtils.doWithoutCommit(
UserRoleRelMapper.class,
@@ -260,6 +266,17 @@ public class UserMetaService {
// metadata-only change, such as the audit info, still has to be written.
try {
SessionUtils.doMultipleWithCommit(
+ () -> {
+ if (!insertRoleIds.isEmpty() || !deleteRoleIds.isEmpty()) {
+ // The cascade writes memberships before principals; this update
does the reverse.
+ // Fence grants and revokes before the principal CAS to avoid
both orphan grants and
+ // a revoke/cascade deadlock. Metadata-only updates write no
membership rows, so they
+ // need no parent lock (which would serialize unrelated updates
on H2).
+ MetalakeMetaService.getInstance()
+ .lockMetalakeForChildWrite(
+ identifier.namespace().level(0),
oldUserPO.getMetalakeId());
+ }
+ },
() -> {
int updated =
SessionUtils.getWithoutCommit(
@@ -271,6 +288,9 @@ public class UserMetaService {
throw userWriteFailure(identifier, oldUserPO, UserLookup.NAME);
}
},
+ () ->
+ RoleMetaService.getInstance()
+ .lockRolesForMembership(oldUserPO.getMetalakeId(),
insertRoleIds),
() -> {
if (insertRoleIds.isEmpty()) {
return;
@@ -398,37 +418,6 @@ public class UserMetaService {
return new PagedResult<>(totalCount, users);
}
- /**
- * Holds the parent metalake row for the rest of the transaction, so the
user cannot be created
- * under a metalake that is going away.
- *
- * <p>The lock is shared, not exclusive: many users can be created under the
same metalake at the
- * same time. Dropping a metalake takes an exclusive lock on this row, so a
drop and a create
- * cannot overlap. Whoever gets the row first wins, and the loser either
sees the metalake gone or
- * inserts under a metalake that is still there.
- *
- * <p>The name is compared again because the ID alone cannot tell a rename
apart: the caller
- * looked the metalake up by name, so a renamed row means the name in the
request no longer
- * exists.
- *
- * <p>The metalake's version is deliberately not compared, matching {@code
CatalogMetaService}.
- * Holding the row is what makes the create safe. An unrelated metalake edit
that commits in
- * between bumps the version without making this create wrong, so comparing
it would reject the
- * create for no reason.
- */
- private void lockMetalakeForUserCreate(MetalakePO observedMetalakePO) {
- OccWriteSupport.lockParentForChildWrite(
- observedMetalakePO.getMetalakeName(),
- Entity.EntityType.METALAKE,
- () ->
- SessionUtils.getWithoutCommit(
- MetalakeMetaMapper.class,
- mapper ->
-
mapper.selectMetalakeMetaByIdForShare(observedMetalakePO.getMetalakeId())),
- null,
- current -> Objects.equals(current.getMetalakeName(),
observedMetalakePO.getMetalakeName()));
- }
-
private RuntimeException userWriteFailure(
NameIdentifier identifier, UserPO observedUserPO, UserLookup lookup) {
// Sessions run at READ_COMMITTED, so a plain read would already see the
latest committed row.
diff --git
a/core/src/test/java/org/apache/gravitino/authorization/TestAccessControlManagerForPermissions.java
b/core/src/test/java/org/apache/gravitino/authorization/TestAccessControlManagerForPermissions.java
index 228ae3b1ac..66d7773f06 100644
---
a/core/src/test/java/org/apache/gravitino/authorization/TestAccessControlManagerForPermissions.java
+++
b/core/src/test/java/org/apache/gravitino/authorization/TestAccessControlManagerForPermissions.java
@@ -62,6 +62,9 @@ import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.Mockito;
public class TestAccessControlManagerForPermissions {
@@ -197,6 +200,116 @@ public class TestAccessControlManagerForPermissions {
}
}
+ @Test
+ public void testRoleDisappearingDuringGrantIsReportedAsIllegalRole() throws
IOException {
+ EntityStore failingStore = Mockito.mock(EntityStore.class);
+ RoleManager roleManager = Mockito.mock(RoleManager.class);
+ Mockito.when(roleManager.getRole(METALAKE,
roleEntity.name())).thenReturn(roleEntity);
+ NoSuchRoleException missing = new NoSuchRoleException("Role was deleted
during grant");
+ Mockito.doThrow(missing).when(failingStore).update(any(), any(), any(),
any());
+ PermissionManager manager = new PermissionManager(failingStore,
roleManager);
+
+ IllegalRoleException userFailure =
+ Assertions.assertThrows(
+ IllegalRoleException.class,
+ () -> manager.grantRolesToUser(METALAKE,
List.of(roleEntity.name()), USER));
+ Assertions.assertSame(missing, userFailure.getCause());
+ IllegalRoleException groupFailure =
+ Assertions.assertThrows(
+ IllegalRoleException.class,
+ () -> manager.grantRolesToGroup(METALAKE,
List.of(roleEntity.name()), GROUP));
+ Assertions.assertSame(missing, groupFailure.getCause());
+ }
+
+ @ParameterizedTest
+ @CsvSource({"false, false", "false, true", "true, false", "true, true"})
+ void testMembershipUpdateCarriesObservedRolesForward(boolean group, boolean
grant)
+ throws IOException {
+ try (EntityStore store = new TestMemoryEntityStore.InMemoryEntityStore()) {
+ store.initialize(config);
+ RoleEntity retained = membershipRole(10L, "retained");
+ RoleEntity target = membershipRole(20L, "target");
+ putMembershipPrincipal(store, group, List.of(retained, target));
+ RoleManager roles = Mockito.mock(RoleManager.class);
+ Mockito.when(roles.getRole(METALAKE, target.name())).thenReturn(target);
+ // "retained" has been deleted and recreated under the same name. The
updater must never
+ // resolve an existing membership by name, so the replacement cannot
inherit the grant and
+ // the update costs no lookup per existing role.
+ Mockito.when(roles.getRole(METALAKE, retained.name()))
+ .thenReturn(membershipRole(11L, retained.name()));
+ PermissionManager manager = new PermissionManager(store, roles);
+ reset(authorizationPlugin);
+
+ changeMembership(manager, group, grant, List.of(target.name()));
+
+ Entity updated = membershipPrincipal(store, group);
+ List<String> names =
+ group ? ((GroupEntity) updated).roleNames() : ((UserEntity)
updated).roleNames();
+ List<Long> ids = group ? ((GroupEntity) updated).roleIds() :
((UserEntity) updated).roleIds();
+ // Granting an already-held role is a no-op, revoking it drops only that
pair.
+ Assertions.assertEquals(
+ grant ? List.of(retained.name(), target.name()) :
List.of(retained.name()), names);
+ Assertions.assertEquals(
+ grant ? List.of(retained.id(), target.id()) :
List.of(retained.id()), ids);
+ Mockito.verify(roles, Mockito.never()).getRole(METALAKE,
retained.name());
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(booleans = {false, true})
+ void testMembershipUpdateRequiresPairedRoleIds(boolean missingIds) throws
IOException {
+ try (EntityStore store = new TestMemoryEntityStore.InMemoryEntityStore()) {
+ store.initialize(config);
+ // Role IDs are optional entity fields, but a name alone cannot prove
membership identity.
+ UserEntity observed =
+ UserEntity.builder()
+ .withId(100L)
+ .withName(USER)
+ .withNamespace(AuthorizationUtils.ofUserNamespace(METALAKE))
+ .withRoleNames(List.of("retained"))
+ .withRoleIds(missingIds ? null : List.of(10L, 11L))
+ .withAuditInfo(auditInfo)
+ .build();
+ store.put(observed, false);
+ RoleManager roles = Mockito.mock(RoleManager.class);
+ RoleEntity target = membershipRole(20L, "target");
+ Mockito.when(roles.getRole(METALAKE, target.name())).thenReturn(target);
+ PermissionManager manager = new PermissionManager(store, roles);
+ reset(authorizationPlugin);
+ Assertions.assertThrows(
+ IllegalRoleException.class,
+ () -> manager.grantRolesToUser(METALAKE, List.of(target.name()),
USER));
+ Assertions.assertSame(observed, membershipPrincipal(store, false));
+ Mockito.verifyNoInteractions(authorizationPlugin);
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(booleans = {false, true})
+ void testRevokeOldRoleDoesNotRemoveReplacementName(boolean group) throws
IOException {
+ try (EntityStore store = new TestMemoryEntityStore.InMemoryEntityStore()) {
+ store.initialize(config);
+ RoleEntity old = membershipRole(10L, "recreated");
+ RoleEntity replacement = membershipRole(11L, old.name());
+ RoleEntity removed = membershipRole(20L, "removed");
+ putMembershipPrincipal(store, group, List.of(replacement, removed));
+ RoleManager roles = Mockito.mock(RoleManager.class);
+ // The request resolves the old ID, then a concurrent operation grants
the replacement
+ // before the principal snapshot is read. Revoking the old ID must
preserve the new pair.
+ Mockito.when(roles.getRole(METALAKE, old.name())).thenReturn(old);
+ Mockito.when(roles.getRole(METALAKE,
removed.name())).thenReturn(removed);
+ PermissionManager manager = new PermissionManager(store, roles);
+ changeMembership(manager, group, false, List.of(old.name(),
removed.name(), removed.name()));
+ Entity updated = membershipPrincipal(store, group);
+ Assertions.assertEquals(
+ List.of(replacement.name()),
+ group ? ((GroupEntity) updated).roleNames() : ((UserEntity)
updated).roleNames());
+ Assertions.assertEquals(
+ List.of(replacement.id()),
+ group ? ((GroupEntity) updated).roleIds() : ((UserEntity)
updated).roleIds());
+ }
+ }
+
@Test
public void testGrantRoleToUser() {
reset(authorizationPlugin);
@@ -465,4 +578,67 @@ public class TestAccessControlManagerForPermissions {
accessControlManager.overridePrivilegesInRole(
METALAKE, notExist, Lists.newArrayList()));
}
+
+ private RoleEntity membershipRole(long id, String name) {
+ return RoleEntity.builder()
+ .withId(id)
+ .withName(name)
+ .withNamespace(AuthorizationUtils.ofRoleNamespace(METALAKE))
+ .withProperties(Maps.newHashMap())
+ .withSecurableObjects(roleEntity.securableObjects())
+ .withAuditInfo(auditInfo)
+ .build();
+ }
+
+ private void putMembershipPrincipal(EntityStore store, boolean group,
List<RoleEntity> roles)
+ throws IOException {
+ List<String> names = roles.stream().map(RoleEntity::name).toList();
+ List<Long> ids = roles.stream().map(RoleEntity::id).toList();
+ if (group) {
+ store.put(
+ GroupEntity.builder()
+ .withId(100L)
+ .withName(GROUP)
+ .withNamespace(AuthorizationUtils.ofGroupNamespace(METALAKE))
+ .withRoleNames(names)
+ .withRoleIds(ids)
+ .withAuditInfo(auditInfo)
+ .build(),
+ false);
+ } else {
+ store.put(
+ UserEntity.builder()
+ .withId(100L)
+ .withName(USER)
+ .withNamespace(AuthorizationUtils.ofUserNamespace(METALAKE))
+ .withRoleNames(names)
+ .withRoleIds(ids)
+ .withAuditInfo(auditInfo)
+ .build(),
+ false);
+ }
+ }
+
+ private Entity membershipPrincipal(EntityStore store, boolean group) throws
IOException {
+ return group
+ ? store.get(
+ AuthorizationUtils.ofGroup(METALAKE, GROUP),
Entity.EntityType.GROUP, GroupEntity.class)
+ : store.get(
+ AuthorizationUtils.ofUser(METALAKE, USER), Entity.EntityType.USER,
UserEntity.class);
+ }
+
+ private void changeMembership(
+ PermissionManager manager, boolean group, boolean grant, List<String>
roles) {
+ if (group) {
+ if (grant) {
+ manager.grantRolesToGroup(METALAKE, roles, GROUP);
+ } else {
+ manager.revokeRolesFromGroup(METALAKE, roles, GROUP);
+ }
+ } else if (grant) {
+ manager.grantRolesToUser(METALAKE, roles, USER);
+ } else {
+ manager.revokeRolesFromUser(METALAKE, roles, USER);
+ }
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestRoleMembershipWrites.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestRoleMembershipWrites.java
new file mode 100644
index 0000000000..eae6773ddf
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestRoleMembershipWrites.java
@@ -0,0 +1,761 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.gravitino.storage.relational.service;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.IOException;
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import org.apache.gravitino.Entity;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.authorization.AuthorizationUtils;
+import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.exceptions.NoSuchRoleException;
+import org.apache.gravitino.meta.GroupEntity;
+import org.apache.gravitino.meta.RoleEntity;
+import org.apache.gravitino.meta.UserEntity;
+import org.apache.gravitino.storage.RandomIdGenerator;
+import org.apache.gravitino.storage.relational.TestJDBCBackend;
+import org.apache.gravitino.storage.relational.mapper.GroupRoleRelMapper;
+import org.apache.gravitino.storage.relational.mapper.MetalakeMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.RoleMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.UserRoleRelMapper;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.gravitino.storage.relational.session.SqlSessions;
+import org.apache.gravitino.storage.relational.utils.SessionUtils;
+import org.apache.ibatis.session.SqlSession;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.TestTemplate;
+import org.junit.jupiter.api.function.Executable;
+
+class TestRoleMembershipWrites extends TestJDBCBackend {
+ private static final String METALAKE = "membership_metalake";
+ private static final String CATALOG = "catalog";
+
+ @TestTemplate
+ void testGrantRejectsRoleDeletedAfterObservation() throws Exception {
+ initialize();
+ RoleEntity retained = role(METALAKE, "retained", true);
+ for (boolean group : List.of(false, true)) {
+ RoleEntity deleted = role(METALAKE, "deleted_" + group, true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(retained), false);
+ long version = version(group, id);
+ assertThrows(
+ NoSuchRoleException.class,
+ () ->
+ updatePrincipal(
+ group,
+ List.of(retained, deleted),
+ () ->
RoleMetaService.getInstance().deleteRole(deleted.nameIdentifier())));
+ assertEquals(version, version(group, id));
+ assertEquals(1, memberships(group, id));
+ }
+ }
+
+ @TestTemplate
+ void testInvalidRoleBatchRollsBackInsertOverwriteAndUpdate() throws
Exception {
+ initialize();
+ String otherMetalake = "other_membership_metalake";
+ createAndInsertMakeLake(otherMetalake);
+ createAndInsertCatalog(otherMetalake, CATALOG);
+ RoleEntity retained = role(METALAKE, "retained", true);
+ RoleEntity valid = role(METALAKE, "valid", true);
+ RoleEntity missing = role(METALAKE, "missing", false);
+ RoleEntity foreign = role(otherMetalake, "foreign", true);
+ RoleEntity deleted = role(METALAKE, "deleted", true);
+ RoleMetaService.getInstance().deleteRole(deleted.nameIdentifier());
+ RoleEntity replaced = role(METALAKE, "recreated", true);
+ RoleMetaService.getInstance().deleteRole(replaced.nameIdentifier());
+ RoleEntity replacement = role(METALAKE, "recreated", true);
+ for (boolean group : List.of(false, true)) {
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ for (RoleEntity invalid : List.of(missing, foreign, deleted, replaced)) {
+ assertThrows(
+ NoSuchRoleException.class,
+ () -> insertPrincipal(group, id, List.of(valid, invalid), false));
+ assertFalse(backend.exists(identifier(group), type(group)));
+ assertEquals(0, memberships(group, id));
+ }
+ insertPrincipal(group, id, List.of(retained), false);
+ long version = version(group, id);
+ for (RoleEntity invalid : List.of(missing, foreign, deleted, replaced)) {
+ assertThrows(
+ NoSuchRoleException.class,
+ () -> insertPrincipal(group, id, List.of(valid, invalid), true));
+ assertEquals(version, version(group, id));
+ assertMemberships(group, id, retained);
+ assertFalse(SessionUtils.isInTransaction());
+ assertThrows(
+ NoSuchRoleException.class,
+ () -> updatePrincipal(group, List.of(valid, invalid), () -> {}));
+ assertEquals(version, version(group, id));
+ assertMemberships(group, id, retained);
+ assertFalse(SessionUtils.isInTransaction());
+ }
+ // The old ID must fail, but a fresh reference to the replacement
remains usable.
+ updatePrincipal(group, List.of(replacement), () -> {});
+ assertMemberships(group, id, replacement);
+ assertFalse(SessionUtils.isInTransaction());
+ }
+ }
+
+ @TestTemplate
+ void testValidOverwriteGrantAndRevokeRemainIdempotent() throws Exception {
+ initialize();
+ RoleEntity first = role(METALAKE, "first", true);
+ RoleEntity second = role(METALAKE, "second", true);
+ for (boolean group : List.of(false, true)) {
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(first), false);
+ insertPrincipal(group, id, List.of(second), true);
+ assertMemberships(group, id, second);
+ updatePrincipal(group, List.of(second, first), () -> {});
+ updatePrincipal(group, List.of(second, first), () -> {});
+ assertMemberships(group, id, first, second);
+ updatePrincipal(group, List.of(), () -> {});
+ assertMemberships(group, id);
+ updatePrincipal(group, List.of(), () -> {});
+ assertMemberships(group, id);
+ }
+ }
+
+ @TestTemplate
+ void testGrantWaitsForUncommittedRoleDelete() throws Exception {
+ initialize();
+ for (boolean group : List.of(false, true)) {
+ RoleEntity role = role(METALAKE, "delete_first_" + group, true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(), false);
+ long version = version(group, id);
+ Throwable failure =
+ whileTransactionHeld(
+ () ->
RoleMetaService.getInstance().deleteRole(role.nameIdentifier()),
+ () -> updatePrincipal(group, List.of(role), () -> {}));
+ Assertions.assertInstanceOf(NoSuchRoleException.class, failure);
+ assertEquals(version, version(group, id));
+ assertEquals(0, memberships(group, id));
+ }
+ }
+
+ @TestTemplate
+ void testInsertAndOverwriteWithoutRoles() throws Exception {
+ initialize();
+ RoleEntity role = role(METALAKE, "cleared_by_overwrite", true);
+ for (boolean group : List.of(false, true)) {
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ for (boolean overwrite : List.of(false, true)) {
+ if (overwrite) {
+ updatePrincipal(group, List.of(role), () -> {});
+ assertMemberships(group, id, role);
+ }
+ if (group) {
+ GroupMetaService.getInstance()
+ .insertGroup(
+ createGroupEntity(
+ id,
+ AuthorizationUtils.ofGroupNamespace(METALAKE),
+ "group",
+ AUDIT_INFO,
+ null,
+ null),
+ overwrite);
+ } else {
+ UserMetaService.getInstance()
+ .insertUser(
+ createUserEntity(
+ id,
+ AuthorizationUtils.ofUserNamespace(METALAKE),
+ "user",
+ AUDIT_INFO,
+ null,
+ null),
+ overwrite);
+ }
+ assertMemberships(group, id);
+ assertFalse(SessionUtils.isInTransaction());
+ }
+ }
+ }
+
+ @TestTemplate
+ void testMembershipLocksRequireTransaction() throws Exception {
+ initialize();
+ RoleEntity role = role(METALAKE, "transaction_required", true);
+ long metalakeId =
MetalakeMetaService.getInstance().getMetalakeIdByName(METALAKE);
+ assertFalse(SessionUtils.isInTransaction());
+ IllegalStateException failure =
+ assertThrows(
+ IllegalStateException.class,
+ () ->
+ RoleMetaService.getInstance()
+ .lockRolesForMembership(metalakeId, List.of(role.id())));
+ assertEquals("Role membership locks require an active transaction",
failure.getMessage());
+ assertFalse(SessionUtils.isInTransaction());
+
+ // A failed standalone call must not leave a session or lock behind.
+ SessionUtils.doMultipleWithCommit(
+ () -> RoleMetaService.getInstance().lockRolesForMembership(metalakeId,
List.of(role.id())));
+ assertFalse(SessionUtils.isInTransaction());
+
assertTrue(RoleMetaService.getInstance().deleteRole(role.nameIdentifier()));
+ }
+
+ @TestTemplate
+ void testGrantSucceedsAfterRoleDeleteRollsBack() throws Exception {
+ initialize();
+ for (boolean group : List.of(false, true)) {
+ RoleEntity role = role(METALAKE, "delete_rollback_" + group, true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(), false);
+ long oldVersion = version(group, id);
+ assertNull(
+ whileTransactionHeld(
+ () ->
RoleMetaService.getInstance().deleteRole(role.nameIdentifier()),
+ () -> updatePrincipal(group, List.of(role), () -> {}),
+ () -> {},
+ false));
+ assertEquals(oldVersion + 1, version(group, id));
+ assertMemberships(group, id, role);
+ assertTrue(backend.exists(role.nameIdentifier(),
Entity.EntityType.ROLE));
+ }
+ }
+
+ @TestTemplate
+ void testRoleDeleteSucceedsAfterGrantRollsBack() throws Exception {
+ initialize();
+ for (boolean group : List.of(false, true)) {
+ RoleEntity role = role(METALAKE, "grant_rollback_" + group, true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(), false);
+ long oldVersion = version(group, id);
+ assertNull(
+ whileTransactionHeld(
+ () -> updatePrincipal(group, List.of(role), () -> {}),
+ () ->
RoleMetaService.getInstance().deleteRole(role.nameIdentifier()),
+ () -> {},
+ false));
+ assertEquals(oldVersion, version(group, id));
+ assertMemberships(group, id);
+ assertFalse(backend.exists(role.nameIdentifier(),
Entity.EntityType.ROLE));
+ }
+ }
+
+ @TestTemplate
+ void testRoleDeleteWaitsForGrantAndCleansMembership() throws Exception {
+ initialize();
+ for (boolean group : List.of(false, true)) {
+ RoleEntity role = role(METALAKE, "grant_first_" + group, true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(), false);
+ assertNull(
+ whileTransactionHeld(
+ () -> updatePrincipal(group, List.of(role), () -> {}),
+ () ->
RoleMetaService.getInstance().deleteRole(role.nameIdentifier())));
+ assertEquals(0, memberships(group, id));
+ assertEquals(
+ 1,
+ queryLong(
+ "SELECT COUNT(*) FROM "
+ + relationTable(group)
+ + " WHERE "
+ + principalColumn(group)
+ + " = "
+ + id));
+ }
+ }
+
+ @TestTemplate
+ void testIndependentGrantsShareRoleLock() throws Exception {
+ initialize();
+ RoleEntity role = role(METALAKE, "shared", true);
+ long userId = RandomIdGenerator.INSTANCE.nextId();
+ long groupId = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(false, userId, List.of(), false);
+ insertPrincipal(true, groupId, List.of(), false);
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ CompletableFuture<Long> started = new CompletableFuture<>();
+ SessionUtils.beginTransaction();
+ try {
+ long holderId = prepareTransaction();
+ updatePrincipal(false, List.of(role), () -> {});
+ Future<Throwable> grant =
+ submitTransaction(
+ executor, started, () -> updatePrincipal(true, List.of(role), ()
-> {}));
+ long contenderId = started.get(10, TimeUnit.SECONDS);
+ if ("h2".equalsIgnoreCase(backendType)) {
+ awaitBlockedBy(grant, contenderId, holderId);
+ } else {
+ // Independent principals can commit grants while the first
transaction still holds its
+ // shared locks on the metalake and role.
+ assertNull(grant.get(10, TimeUnit.SECONDS));
+ }
+ SessionUtils.commitTransaction();
+ assertNull(grant.get(10, TimeUnit.SECONDS));
+ assertEquals(1, memberships(false, userId));
+ assertEquals(1, memberships(true, groupId));
+ } finally {
+ SessionUtils.rollbackTransaction();
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+ }
+ }
+
+ @TestTemplate
+ void testMetalakeCascadeWaitsForGrant() throws Exception {
+ for (boolean group : List.of(false, true)) {
+ initialize();
+ RoleEntity role = role(METALAKE, "cascade", true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(), false);
+ assertNull(
+ whileTransactionHeld(
+ () -> updatePrincipal(group, List.of(role), () -> {}),
+ () -> backend.delete(NameIdentifier.of(METALAKE),
Entity.EntityType.METALAKE, true)));
+ assertEquals(0, memberships(group, id));
+ assertFalse(backend.exists(identifier(group), type(group)));
+ }
+ }
+
+ @TestTemplate
+ void testMetalakeCascadeWaitsForMetadataOnlyUpdate() throws Exception {
+ for (boolean group : List.of(false, true)) {
+ initialize();
+ RoleEntity role = role(METALAKE, "metadata_cascade", true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(role), false);
+ assertNull(
+ whileTransactionHeld(
+ () -> updatePrincipal(group, List.of(role), () -> {}),
+ () -> backend.delete(NameIdentifier.of(METALAKE),
Entity.EntityType.METALAKE, true)));
+ assertEquals(0, memberships(group, id));
+ assertFalse(backend.exists(identifier(group), type(group)));
+ }
+ }
+
+ @TestTemplate
+ void testRevokeWaitsForMetalakeCascade() throws Exception {
+ for (boolean group : List.of(false, true)) {
+ initialize();
+ RoleEntity role = role(METALAKE, "revoke_cascade", true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(role), false);
+ long metalakeId =
MetalakeMetaService.getInstance().getMetalakeIdByName(METALAKE);
+ Throwable failure =
+ whileTransactionHeld(
+ () -> {
+ lockMetalake(metalakeId);
+ // Pause a cascade after membership cleanup but before its
principal write.
+ if (group) {
+ SessionUtils.doWithoutCommit(
+ GroupRoleRelMapper.class,
+ mapper ->
mapper.softDeleteGroupRoleRelByMetalakeId(metalakeId));
+ } else {
+ SessionUtils.doWithoutCommit(
+ UserRoleRelMapper.class,
+ mapper ->
mapper.softDeleteUserRoleRelByMetalakeId(metalakeId));
+ }
+ },
+ () -> updatePrincipal(group, List.of(), () -> {}),
+ () -> backend.delete(NameIdentifier.of(METALAKE),
Entity.EntityType.METALAKE, true));
+ Assertions.assertInstanceOf(NoSuchEntityException.class, failure);
+ assertEquals(0, memberships(group, id));
+ assertFalse(backend.exists(identifier(group), type(group)));
+ }
+ }
+
+ @TestTemplate
+ void testUnchangedMembershipsAndRevokesAvoidRoleLocks() throws Exception {
+ initialize();
+ RoleEntity role = role(METALAKE, "unchanged", true);
+ for (boolean group : List.of(false, true)) {
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(role), false);
+ for (boolean revoke : List.of(false, true)) {
+ long oldVersion = version(group, id);
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ SessionUtils.beginTransaction();
+ try {
+ // Unchanged memberships need no metalake lock. A revoke needs the
metalake lock, but
+ // neither operation needs to lock a retained or removed role.
+ if (!revoke) {
+
lockMetalake(MetalakeMetaService.getInstance().getMetalakeIdByName(METALAKE));
+ }
+ SessionUtils.getWithoutCommit(
+ RoleMetaMapper.class, mapper ->
mapper.selectRoleMetaByIdForUpdate(role.id()));
+ Future<?> update =
+ executor.submit(
+ () -> {
+ updatePrincipal(group, revoke ? List.of() : List.of(role),
() -> {});
+ return null;
+ });
+ update.get(10, TimeUnit.SECONDS);
+ } finally {
+ SessionUtils.rollbackTransaction();
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+ }
+ assertEquals(oldVersion + 1, version(group, id));
+ assertEquals(revoke ? 0 : 1, memberships(group, id));
+ }
+ }
+ }
+
+ @TestTemplate
+ void testRoleDeleteWaitsForInsertAndOverwrite() throws Exception {
+ initialize();
+ for (boolean group : List.of(false, true)) {
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ for (boolean overwrite : List.of(false, true)) {
+ RoleEntity role = role(METALAKE, "insert_" + group + "_" + overwrite,
true);
+ assertNull(
+ whileTransactionHeld(
+ () -> insertPrincipal(group, id, List.of(role), overwrite),
+ () ->
RoleMetaService.getInstance().deleteRole(role.nameIdentifier())));
+ assertEquals(0, memberships(group, id));
+ }
+ }
+ }
+
+ @TestTemplate
+ void testGrantRejectsMetalakeDeletedAfterObservation() throws Exception {
+ for (boolean group : List.of(false, true)) {
+ initialize();
+ RoleEntity role = role(METALAKE, "deleted_metalake", true);
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ insertPrincipal(group, id, List.of(), false);
+ assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ updatePrincipal(
+ group,
+ List.of(role),
+ () ->
+ Assertions.assertDoesNotThrow(
+ () ->
+ backend.delete(
+ NameIdentifier.of(METALAKE),
Entity.EntityType.METALAKE, true))));
+ assertEquals(0, memberships(group, id));
+ assertFalse(backend.exists(identifier(group), type(group)));
+ }
+ }
+
+ private void lockMetalake(long metalakeId) {
+ SessionUtils.getWithoutCommit(
+ MetalakeMetaMapper.class, mapper ->
mapper.selectMetalakeMetaByIdForUpdate(metalakeId));
+ }
+
+ private void initialize() throws IOException {
+ createAndInsertMakeLake(METALAKE);
+ createAndInsertCatalog(METALAKE, CATALOG);
+ }
+
+ private RoleEntity role(String metalake, String name, boolean insert) throws
IOException {
+ RoleEntity role =
+ createRoleEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ AuthorizationUtils.ofRoleNamespace(metalake),
+ name,
+ AUDIT_INFO,
+ CATALOG);
+ if (insert) {
+ RoleMetaService.getInstance().insertRole(role, false);
+ }
+ return role;
+ }
+
+ private UserEntity user(long id, List<RoleEntity> roles) {
+ return createUserEntity(
+ id,
+ AuthorizationUtils.ofUserNamespace(METALAKE),
+ "user",
+ AUDIT_INFO,
+ roles.stream().map(RoleEntity::name).collect(Collectors.toList()),
+ roles.stream().map(RoleEntity::id).collect(Collectors.toList()));
+ }
+
+ private GroupEntity group(long id, List<RoleEntity> roles) {
+ return createGroupEntity(
+ id,
+ AuthorizationUtils.ofGroupNamespace(METALAKE),
+ "group",
+ AUDIT_INFO,
+ roles.stream().map(RoleEntity::name).collect(Collectors.toList()),
+ roles.stream().map(RoleEntity::id).collect(Collectors.toList()));
+ }
+
+ private void insertPrincipal(boolean group, long id, List<RoleEntity> roles,
boolean overwrite)
+ throws IOException {
+ if (group) {
+ GroupMetaService.getInstance().insertGroup(group(id, roles), overwrite);
+ } else {
+ UserMetaService.getInstance().insertUser(user(id, roles), overwrite);
+ }
+ }
+
+ private void updatePrincipal(boolean group, List<RoleEntity> roles, Runnable
beforeWrite)
+ throws IOException {
+ if (group) {
+ GroupMetaService.getInstance()
+ .updateGroup(
+ identifier(true),
+ (GroupEntity old) -> {
+ beforeWrite.run();
+ return group(old.id(), roles);
+ });
+ } else {
+ UserMetaService.getInstance()
+ .updateUser(
+ identifier(false),
+ (UserEntity old) -> {
+ beforeWrite.run();
+ return user(old.id(), roles);
+ });
+ }
+ }
+
+ private NameIdentifier identifier(boolean group) {
+ return group
+ ? AuthorizationUtils.ofGroup(METALAKE, "group")
+ : AuthorizationUtils.ofUser(METALAKE, "user");
+ }
+
+ private Entity.EntityType type(boolean group) {
+ return group ? Entity.EntityType.GROUP : Entity.EntityType.USER;
+ }
+
+ private String relationTable(boolean group) {
+ return group ? "group_role_rel" : "user_role_rel";
+ }
+
+ private String principalColumn(boolean group) {
+ return group ? "group_id" : "user_id";
+ }
+
+ private void assertMemberships(boolean group, long id, RoleEntity...
expected) throws Exception {
+ Set<Long> actual = new HashSet<>();
+ String sql =
+ "SELECT role_id FROM "
+ + relationTable(group)
+ + " WHERE "
+ + principalColumn(group)
+ + " = "
+ + id
+ + " AND deleted_at = 0";
+ try (SqlSession session =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Statement statement = session.getConnection().createStatement();
+ ResultSet rows = statement.executeQuery(sql)) {
+ while (rows.next()) {
+ assertTrue(actual.add(rows.getLong(1)), "Duplicate active role
membership");
+ }
+ }
+ Set<Long> expectedIds = new HashSet<>();
+ for (RoleEntity role : expected) {
+ expectedIds.add(role.id());
+ }
+ assertEquals(expectedIds, actual);
+ }
+
+ private long memberships(boolean group, long id) throws Exception {
+ return queryLong(
+ "SELECT COUNT(*) FROM "
+ + relationTable(group)
+ + " WHERE "
+ + principalColumn(group)
+ + " = "
+ + id
+ + " AND deleted_at = 0");
+ }
+
+ private long version(boolean group, long id) throws Exception {
+ return queryLong(
+ "SELECT current_version FROM "
+ + (group ? "group_meta" : "user_meta")
+ + " WHERE "
+ + principalColumn(group)
+ + " = "
+ + id
+ + " AND deleted_at = 0");
+ }
+
+ private long queryLong(String sql) throws Exception {
+ try (SqlSession session =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Statement statement = session.getConnection().createStatement();
+ ResultSet rows = statement.executeQuery(sql)) {
+ assertTrue(rows.next());
+ return rows.getLong(1);
+ }
+ }
+
+ private Throwable whileTransactionHeld(Executable holder, Executable
contender) throws Exception {
+ return whileTransactionHeld(holder, contender, () -> {});
+ }
+
+ private Throwable whileTransactionHeld(
+ Executable holder, Executable contender, Executable beforeCommit) throws
Exception {
+ return whileTransactionHeld(holder, contender, beforeCommit, true);
+ }
+
+ private Throwable whileTransactionHeld(
+ Executable holder, Executable contender, Executable beforeCompletion,
boolean commitHolder)
+ throws Exception {
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ CompletableFuture<Long> started = new CompletableFuture<>();
+ SessionUtils.beginTransaction();
+ try {
+ long holderId = prepareTransaction();
+ Assertions.assertDoesNotThrow(holder);
+ Future<Throwable> result = submitTransaction(executor, started,
contender);
+ awaitBlockedBy(result, started.get(10, TimeUnit.SECONDS), holderId);
+ Assertions.assertDoesNotThrow(beforeCompletion);
+ if (commitHolder) {
+ SessionUtils.commitTransaction();
+ } else {
+ SessionUtils.rollbackTransaction();
+ }
+ return result.get(10, TimeUnit.SECONDS);
+ } finally {
+ SessionUtils.rollbackTransaction();
+ executor.shutdownNow();
+ assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+ }
+ }
+
+ private Future<Throwable> submitTransaction(
+ ExecutorService executor, CompletableFuture<Long> started, Executable
operation) {
+ return executor.submit(
+ () -> {
+ SessionUtils.beginTransaction();
+ try {
+ started.complete(prepareTransaction());
+ operation.execute();
+ SessionUtils.commitTransaction();
+ return null;
+ } catch (Throwable failure) {
+ started.completeExceptionally(failure);
+ return failure;
+ } finally {
+ SessionUtils.rollbackTransaction();
+ }
+ });
+ }
+
+ private long prepareTransaction() throws SQLException {
+ SqlSession session = SqlSessions.getSqlSession();
+ try (Statement statement = session.getConnection().createStatement()) {
+ String sessionIdQuery;
+ switch (backendType) {
+ case "h2":
+ // Keep the engine timeout above the test's lock-observation
deadline.
+ statement.execute("SET LOCK_TIMEOUT 30000");
+ sessionIdQuery = "SELECT SESSION_ID()";
+ break;
+ case "mysql":
+ statement.execute("SET SESSION innodb_lock_wait_timeout = 30");
+ sessionIdQuery = "SELECT CONNECTION_ID()";
+ break;
+ case "postgresql":
+ statement.execute("SET LOCAL lock_timeout = '30s'");
+ sessionIdQuery = "SELECT pg_backend_pid()";
+ break;
+ default:
+ throw new IllegalStateException("Unsupported backend: " +
backendType);
+ }
+ try (ResultSet rows = statement.executeQuery(sessionIdQuery)) {
+ assertTrue(rows.next());
+ return rows.getLong(1);
+ }
+ } finally {
+ SqlSessions.closeSqlSession();
+ }
+ }
+
+ private void awaitBlockedBy(Future<Throwable> result, long contenderId, long
holderId)
+ throws Exception {
+ String query;
+ switch (backendType) {
+ case "h2":
+ query =
+ "SELECT COUNT(*) FROM INFORMATION_SCHEMA.SESSIONS WHERE SESSION_ID
= "
+ + contenderId
+ + " AND BLOCKER_ID = "
+ + holderId;
+ break;
+ case "mysql":
+ query =
+ "SELECT COUNT(*) FROM performance_schema.data_lock_waits w"
+ + " JOIN performance_schema.threads r ON r.THREAD_ID =
w.REQUESTING_THREAD_ID"
+ + " JOIN performance_schema.threads b ON b.THREAD_ID =
w.BLOCKING_THREAD_ID"
+ + " WHERE r.PROCESSLIST_ID = "
+ + contenderId
+ + " AND b.PROCESSLIST_ID = "
+ + holderId;
+ break;
+ case "postgresql":
+ query =
+ "SELECT COUNT(*) FROM unnest(pg_blocking_pids("
+ + contenderId
+ + ")) AS blocker(pid) WHERE pid = "
+ + holderId;
+ break;
+ default:
+ throw new IllegalStateException("Unsupported backend: " + backendType);
+ }
+ // Observe the actual waiter/blocker pair. A slow thread or connection
checkout alone cannot
+ // satisfy this assertion, and an unexpectedly completed operation fails
immediately.
+ try (SqlSession observer =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true))
{
+ Connection connection = observer.getConnection();
+ await()
+ .pollInSameThread()
+ .atMost(10, TimeUnit.SECONDS)
+ .pollInterval(10, TimeUnit.MILLISECONDS)
+ .until(
+ () -> {
+ if (result.isDone()) {
+ throw new AssertionError(
+ "Operation completed without waiting for the holder",
result.get());
+ }
+ try (Statement statement = connection.createStatement();
+ ResultSet rows = statement.executeQuery(query)) {
+ return rows.next() && rows.getLong(1) > 0;
+ }
+ });
+ }
+ }
+}