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 7f303cadad [#13002] fix(core): Fence owner assignment on the principal
and keep one live owner per object (#13293)
7f303cadad is described below
commit 7f303cadad7623dfe8438ae8c2b9269f8e21e919
Author: Qi Yu <[email protected]>
AuthorDate: Mon Sep 21 21:53:29 2026 +0800
[#13002] fix(core): Fence owner assignment on the principal and keep one
live owner per object (#13293)
### What changes were proposed in this pull request?
- Fence the metalake, owned object, and owner principal (User/Group) in
one `setOwner` or `batchSetOwners` transaction. Assignments lock the
owned object's row with `FOR UPDATE`, including when it has no owner
yet. Batches lock objects in ID order. A metalake's own owner assignment
takes an exclusive metalake lock.
- Reject a principal that was deleted, replaced under the same name, or
moved to another metalake while the assignment was waiting.
- Keep the existing `owner_meta` unique key. The 2.0.0 upgrade scripts
merge duplicate live owner rows left by older concurrent writes,
retaining the row with the largest ID.
### Why are the changes needed?
An assignment previously resolved the principal before its transaction
and could insert a live owner row after another server deleted that
principal. Two concurrent first assignments could also insert different
live owners for the same object because no owner row existed to lock and
the existing unique key includes `owner_id`. Locking the stable
owned-object row serializes those assignments without adding a column or
replacing the unique key.
Fix: #13002
### Does this PR introduce _any_ user-facing change?
No API change. An assignment waiting on a deleted object or principal
now fails with the existing not-found error. Concurrent assignments
through this service leave one live owner, with the later assignment
taking effect.
### How was this patch tested?
- `TestOwnerAssignmentWrites`: concurrent first assignments,
reassignments, principal deletion and replacement, object deletion,
metalake ownership, and batch rollback and retry.
- `TestSQLScripts`: the 1.3.0 to 2.0.0 upgrade merges three live owners
of one object while preserving historical rows that share a deletion
timestamp.
- `TestOwnerMetaService`, `TestOwnerAssignmentWrites`, and
`TestSQLScripts` passed against H2, MySQL, and PostgreSQL.
`:core:spotlessApply` passed.
---
.../storage/relational/mapper/GroupMetaMapper.java | 12 +
.../mapper/GroupMetaSQLProviderFactory.java | 5 +
.../storage/relational/mapper/OwnerMetaMapper.java | 16 +
.../mapper/OwnerMetaSQLProviderFactory.java | 82 ++++
.../storage/relational/mapper/UserMetaMapper.java | 12 +
.../mapper/UserMetaSQLProviderFactory.java | 5 +
.../provider/base/GroupMetaBaseSQLProvider.java | 12 +-
.../provider/base/UserMetaBaseSQLProvider.java | 12 +-
.../mapper/provider/h2/GroupMetaH2Provider.java | 6 +
.../mapper/provider/h2/UserMetaH2Provider.java | 6 +
.../postgresql/GroupMetaPostgreSQLProvider.java | 5 +
.../postgresql/UserMetaPostgreSQLProvider.java | 5 +
.../relational/service/OwnerMetaService.java | 147 +++++-
.../apache/gravitino/storage/TestSQLScripts.java | 75 +++
.../service/TestOwnerAssignmentWrites.java | 529 +++++++++++++++++++++
scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql | 13 +
scripts/mysql/schema-2.0.0-mysql.sql | 2 +-
scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql | 15 +
.../upgrade-1.3.0-to-2.0.0-postgresql.sql | 13 +
19 files changed, 953 insertions(+), 19 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.java
index a5bd65d5c3..333be53597 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaMapper.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.ExtendedGroupPO;
import org.apache.gravitino.storage.relational.po.GroupPO;
import org.apache.gravitino.storage.relational.po.auth.GroupUpdatedAt;
@@ -57,6 +58,17 @@ public interface GroupMetaMapper {
@SelectProvider(type = GroupMetaSQLProviderFactory.class, method =
"selectGroupMetaByIdForUpdate")
GroupPO selectGroupMetaByIdForUpdate(@Param("groupId") Long groupId);
+ /**
+ * Returns an active group 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 group, or null if it does not exist
+ */
+ @Nullable
+ @SelectProvider(type = GroupMetaSQLProviderFactory.class, method =
"selectGroupMetaByIdForShare")
+ GroupPO selectGroupMetaByIdForShare(@Param("groupId") Long groupId);
+
@SelectProvider(
type = GroupMetaSQLProviderFactory.class,
method = "listExtendedGroupPOsByMetalakeIdAndNames")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
index 14c6ab732e..215b4083c6 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/GroupMetaSQLProviderFactory.java
@@ -64,6 +64,11 @@ public class GroupMetaSQLProviderFactory {
return getProvider().selectGroupMetaByIdForUpdate(groupId);
}
+ /** Returns SQL that selects an active group by ID and locks it for shared
access. */
+ public static String selectGroupMetaByIdForShare(@Param("groupId") Long
groupId) {
+ return getProvider().selectGroupMetaByIdForShare(groupId);
+ }
+
public static String listExtendedGroupPOsByMetalakeIdAndNames(
@Param("metalakeId") Long metalakeId, @Param("groupNames") List<String>
groupNames) {
return getProvider().listExtendedGroupPOsByMetalakeIdAndNames(metalakeId,
groupNames);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
index e4ef4b6096..fc456f10cb 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaMapper.java
@@ -19,6 +19,8 @@
package org.apache.gravitino.storage.relational.mapper;
import java.util.List;
+import javax.annotation.Nullable;
+import org.apache.gravitino.Entity;
import org.apache.gravitino.storage.relational.po.GroupOwnerRelPO;
import org.apache.gravitino.storage.relational.po.GroupPO;
import org.apache.gravitino.storage.relational.po.OwnerRelForDeletion;
@@ -44,6 +46,20 @@ public interface OwnerMetaMapper {
String OWNER_TABLE_NAME = "owner_meta";
+ /**
+ * Locks an active metadata object before owner assignment.
+ *
+ * @return the object ID, or null if the object is no longer active
+ */
+ @Nullable
+ @SelectProvider(
+ type = OwnerMetaSQLProviderFactory.class,
+ method = "selectMetadataObjectIdForUpdate")
+ Long selectMetadataObjectIdForUpdate(
+ @Param("entityId") Long entityId,
+ @Param("metalakeId") Long metalakeId,
+ @Param("entityType") Entity.EntityType entityType);
+
@SelectProvider(
type = OwnerMetaSQLProviderFactory.class,
method = "selectUserOwnerMetaByMetadataObjectIdAndType")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
index 5957c6d91b..777904dabb 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/OwnerMetaSQLProviderFactory.java
@@ -21,6 +21,7 @@ package org.apache.gravitino.storage.relational.mapper;
import com.google.common.collect.ImmutableMap;
import java.util.List;
import java.util.Map;
+import org.apache.gravitino.Entity;
import org.apache.gravitino.storage.relational.JDBCBackend.JDBCBackendType;
import
org.apache.gravitino.storage.relational.mapper.provider.base.OwnerMetaBaseSQLProvider;
import
org.apache.gravitino.storage.relational.mapper.provider.postgresql.OwnerMetaPostgreSQLProvider;
@@ -52,6 +53,87 @@ public class OwnerMetaSQLProviderFactory {
static class OwnerMetaH2Provider extends OwnerMetaBaseSQLProvider {}
+ /** Returns SQL that locks an active metadata object before assigning its
owner. */
+ public static String selectMetadataObjectIdForUpdate(
+ @Param("entityId") Long entityId,
+ @Param("metalakeId") Long metalakeId,
+ @Param("entityType") Entity.EntityType entityType) {
+ String table;
+ String idColumn;
+ switch (entityType) {
+ case CATALOG:
+ table = CatalogMetaMapper.TABLE_NAME;
+ idColumn = "catalog_id";
+ break;
+ case SCHEMA:
+ table = SchemaMetaMapper.TABLE_NAME;
+ idColumn = "schema_id";
+ break;
+ case TABLE:
+ table = TableMetaMapper.TABLE_NAME;
+ idColumn = "table_id";
+ break;
+ case COLUMN:
+ table = TableColumnMapper.COLUMN_TABLE_NAME;
+ idColumn = "column_id";
+ break;
+ case FILESET:
+ table = FilesetMetaMapper.META_TABLE_NAME;
+ idColumn = "fileset_id";
+ break;
+ case TOPIC:
+ table = TopicMetaMapper.TABLE_NAME;
+ idColumn = "topic_id";
+ break;
+ case MODEL:
+ table = ModelMetaMapper.TABLE_NAME;
+ idColumn = "model_id";
+ break;
+ case VIEW:
+ table = ViewMetaMapper.TABLE_NAME;
+ idColumn = "view_id";
+ break;
+ case FUNCTION:
+ table = FunctionMetaMapper.TABLE_NAME;
+ idColumn = "function_id";
+ break;
+ case ROLE:
+ table = RoleMetaMapper.ROLE_TABLE_NAME;
+ idColumn = "role_id";
+ break;
+ case TAG:
+ table = TagMetaMapper.TAG_TABLE_NAME;
+ idColumn = "tag_id";
+ break;
+ case POLICY:
+ table = PolicyMetaMapper.POLICY_META_TABLE_NAME;
+ idColumn = "policy_id";
+ break;
+ case JOB_TEMPLATE:
+ table = JobTemplateMetaMapper.TABLE_NAME;
+ idColumn = "job_template_id";
+ break;
+ case JOB:
+ table = JobMetaMapper.TABLE_NAME;
+ idColumn = "job_run_id";
+ break;
+ default:
+ throw new IllegalArgumentException("Unsupported owned object type: " +
entityType);
+ }
+ // Column versions can share a column ID; always lock the same oldest live
row.
+ String orderColumn = entityType == Entity.EntityType.COLUMN ? "id" :
idColumn;
+ return "SELECT "
+ + idColumn
+ + " FROM "
+ + table
+ + " WHERE "
+ + idColumn
+ + " = #{entityId} AND metalake_id = #{metalakeId} AND deleted_at = 0"
+ + " ORDER BY "
+ + orderColumn
+ + " LIMIT 1 FOR UPDATE";
+ }
+
public static String selectUserOwnerMetaByMetadataObjectIdAndType(
@Param("metadataObjectId") Long metadataObjectId,
@Param("metadataObjectType") String metadataObjectType) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.java
index ba79e5dca0..4a10ac8dca 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaMapper.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.ExtendedUserPO;
import org.apache.gravitino.storage.relational.po.UserPO;
import org.apache.gravitino.storage.relational.po.auth.AuthPrefetchRow;
@@ -58,6 +59,17 @@ public interface UserMetaMapper {
@SelectProvider(type = UserMetaSQLProviderFactory.class, method =
"selectUserMetaByIdForUpdate")
UserPO selectUserMetaByIdForUpdate(@Param("userId") Long userId);
+ /**
+ * Returns an active user 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 user, or null if it does not exist
+ */
+ @Nullable
+ @SelectProvider(type = UserMetaSQLProviderFactory.class, method =
"selectUserMetaByIdForShare")
+ UserPO selectUserMetaByIdForShare(@Param("userId") Long userId);
+
@InsertProvider(type = UserMetaSQLProviderFactory.class, method =
"insertUserMeta")
void insertUserMeta(@Param("userMeta") UserPO userPO);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
index 0779d9920f..bc52e64d52 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/UserMetaSQLProviderFactory.java
@@ -66,6 +66,11 @@ public class UserMetaSQLProviderFactory {
return getProvider().selectUserMetaByIdForUpdate(userId);
}
+ /** Returns SQL that selects an active user by ID and locks it for shared
access. */
+ public static String selectUserMetaByIdForShare(@Param("userId") Long
userId) {
+ return getProvider().selectUserMetaByIdForShare(userId);
+ }
+
public static String insertUserMeta(@Param("userMeta") UserPO userPO) {
return getProvider().insertUserMeta(userPO);
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
index 002d68d2b5..4afa7ad3d2 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/GroupMetaBaseSQLProvider.java
@@ -141,13 +141,23 @@ public class GroupMetaBaseSQLProvider {
/** Returns SQL that selects and locks an active group by ID. */
public String selectGroupMetaByIdForUpdate(@Param("groupId") Long groupId) {
+ return selectGroupMetaById(groupId) + " FOR UPDATE";
+ }
+
+ /** Returns SQL that selects an active group by ID and locks it for shared
access. */
+ public String selectGroupMetaByIdForShare(@Param("groupId") Long groupId) {
+ return selectGroupMetaById(groupId) + " LOCK IN SHARE MODE";
+ }
+
+ /** Returns SQL that selects an active group by ID. */
+ protected String selectGroupMetaById(Long groupId) {
return "SELECT group_id as groupId, group_name as groupName,"
+ " metalake_id as metalakeId, audit_info as auditInfo,"
+ " current_version as currentVersion, last_version as lastVersion,"
+ " deleted_at as deletedAt"
+ " FROM "
+ GROUP_TABLE_NAME
- + " WHERE group_id = #{groupId} AND deleted_at = 0 FOR UPDATE";
+ + " WHERE group_id = #{groupId} AND deleted_at = 0";
}
public String listExtendedGroupPOsByMetalakeIdAndNames(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
index 98e8931f1d..dbe362b7db 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/UserMetaBaseSQLProvider.java
@@ -55,13 +55,23 @@ public class UserMetaBaseSQLProvider {
/** Returns SQL that selects and locks an active user by ID. */
public String selectUserMetaByIdForUpdate(@Param("userId") Long userId) {
+ return selectUserMetaById(userId) + " FOR UPDATE";
+ }
+
+ /** Returns SQL that selects an active user by ID and locks it for shared
access. */
+ public String selectUserMetaByIdForShare(@Param("userId") Long userId) {
+ return selectUserMetaById(userId) + " LOCK IN SHARE MODE";
+ }
+
+ /** Returns SQL that selects an active user by ID. */
+ protected String selectUserMetaById(Long userId) {
return "SELECT user_id as userId, user_name as userName,"
+ " metalake_id as metalakeId,"
+ " audit_info as auditInfo, current_version as currentVersion,"
+ " last_version as lastVersion, deleted_at as deletedAt"
+ " FROM "
+ USER_TABLE_NAME
- + " WHERE user_id = #{userId} AND deleted_at = 0 FOR UPDATE";
+ + " WHERE user_id = #{userId} AND deleted_at = 0";
}
public String insertUserMeta(@Param("userMeta") UserPO userPO) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
index 2d542428cd..e25aeb4592 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/GroupMetaH2Provider.java
@@ -28,6 +28,12 @@ import
org.apache.gravitino.storage.relational.mapper.provider.base.GroupMetaBas
import org.apache.ibatis.annotations.Param;
public class GroupMetaH2Provider extends GroupMetaBaseSQLProvider {
+ @Override
+ public String selectGroupMetaByIdForShare(Long groupId) {
+ // H2 has no shared row-lock syntax, matching the other parent-fencing
providers.
+ return selectGroupMetaByIdForUpdate(groupId);
+ }
+
@Override
public String listExtendedGroupPOsByMetalakeId(@Param("metalakeId") Long
metalakeId) {
return "SELECT gt.group_id as groupId, gt.group_name as groupName,"
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
index 83fe6774b5..9e28a2ad9f 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/h2/UserMetaH2Provider.java
@@ -27,6 +27,12 @@ import
org.apache.gravitino.storage.relational.mapper.provider.base.UserMetaBase
import org.apache.ibatis.annotations.Param;
public class UserMetaH2Provider extends UserMetaBaseSQLProvider {
+ @Override
+ public String selectUserMetaByIdForShare(Long userId) {
+ // H2 has no shared row-lock syntax, matching the other parent-fencing
providers.
+ return selectUserMetaByIdForUpdate(userId);
+ }
+
@Override
public String listExtendedUserPOsByMetalakeId(@Param("metalakeId") Long
metalakeId) {
return "SELECT ut.user_id as userId, ut.user_name as userName,"
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
index 1a24acec0f..bac95d6f99 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/GroupMetaPostgreSQLProvider.java
@@ -30,6 +30,11 @@ import org.apache.gravitino.storage.relational.po.GroupPO;
import org.apache.ibatis.annotations.Param;
public class GroupMetaPostgreSQLProvider extends GroupMetaBaseSQLProvider {
+ @Override
+ public String selectGroupMetaByIdForShare(Long groupId) {
+ return selectGroupMetaById(groupId) + " FOR SHARE";
+ }
+
@Override
public String softDeleteGroupMetaByGroupId(Long groupId, Long
currentVersion) {
return "UPDATE "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
index 4562901e8b..5ec9ef7030 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/UserMetaPostgreSQLProvider.java
@@ -29,6 +29,11 @@ import org.apache.gravitino.storage.relational.po.UserPO;
import org.apache.ibatis.annotations.Param;
public class UserMetaPostgreSQLProvider extends UserMetaBaseSQLProvider {
+ @Override
+ public String selectUserMetaByIdForShare(Long userId) {
+ return selectUserMetaById(userId) + " FOR SHARE";
+ }
+
@Override
public String softDeleteUserMetaByUserId(Long userId, Long currentVersion) {
return "UPDATE "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
index 2b261e35b3..4221254694 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/OwnerMetaService.java
@@ -23,6 +23,7 @@ import static
org.apache.gravitino.metrics.source.MetricsSource.GRAVITINO_RELATI
import com.google.common.base.Preconditions;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -38,7 +39,10 @@ import org.apache.gravitino.authorization.AuthorizationUtils;
import org.apache.gravitino.meta.GroupEntity;
import org.apache.gravitino.meta.UserEntity;
import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.mapper.GroupMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.MetalakeMetaMapper;
import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.UserMetaMapper;
import org.apache.gravitino.storage.relational.po.GroupOwnerRelPO;
import org.apache.gravitino.storage.relational.po.GroupPO;
import org.apache.gravitino.storage.relational.po.OwnerRelForDeletion;
@@ -52,10 +56,10 @@ import org.apache.gravitino.utils.NameIdentifierUtil;
/** This class is an utilization class to retrieve owner relation. */
public class OwnerMetaService {
- private OwnerMetaService() {}
-
private static final OwnerMetaService INSTANCE = new OwnerMetaService();
+ private OwnerMetaService() {}
+
public static OwnerMetaService getInstance() {
return INSTANCE;
}
@@ -173,9 +177,8 @@ public class OwnerMetaService {
Entity.EntityType entityType,
NameIdentifier owner,
Entity.EntityType ownerType) {
- long metalakeId =
- MetalakeMetaService.getInstance()
- .getMetalakeIdByName(NameIdentifierUtil.getMetalake(entity));
+ String metalake = NameIdentifierUtil.getMetalake(entity);
+ long metalakeId =
MetalakeMetaService.getInstance().getMetalakeIdByName(metalake);
Long entityId = EntityIdService.getEntityId(entity, entityType);
Long ownerId = EntityIdService.getEntityId(owner, ownerType);
@@ -183,14 +186,22 @@ public class OwnerMetaService {
OwnerRelPO ownerRelPO =
POConverters.initializeOwnerRelPOsWithVersion(
metalakeId, ownerType.name(), ownerId, entityType.name(),
entityId);
- SessionUtils.doMultipleWithCommit(
+ String metadataObjectType =
+ NameIdentifierUtil.toMetadataObject(entity, entityType).type().name();
+ assignOwner(
+ metalake,
+ metalakeId,
+ owner,
+ ownerType,
+ ownerId,
+ List.of(entityId),
+ entityType,
() ->
SessionUtils.doWithoutCommit(
OwnerMetaMapper.class,
mapper ->
mapper.softDeleteOwnerRelByMetadataObjectIdAndType(
- entityId,
- NameIdentifierUtil.toMetadataObject(entity,
entityType).type().name())),
+ entityId, metadataObjectType)),
() ->
SessionUtils.doWithoutCommit(
OwnerMetaMapper.class, mapper ->
mapper.insertOwnerRel(ownerRelPO)));
@@ -218,20 +229,33 @@ public class OwnerMetaService {
long metalakeId =
MetalakeMetaService.getInstance().getMetalakeIdByName(metalake);
Long ownerId = EntityIdService.getEntityId(ownerIdent, ownerType);
- List<OwnerRelForDeletion> deletions = new ArrayList<>(ownedObjects.size());
- List<OwnerRelPO> ownerRelPOs = new ArrayList<>(ownedObjects.size());
+ // Resolve every object first and write in stable id order, so two batches
that overlap take
+ // the same row locks in the same order and cannot deadlock each other.
+ List<Long> entityIds = new ArrayList<>(ownedObjects.size());
for (NameIdentifier entity : ownedObjects) {
- Long entityId = EntityIdService.getEntityId(entity, ownedObjectType);
- deletions.add(
- new OwnerRelForDeletion(
- entityId,
- NameIdentifierUtil.toMetadataObject(entity,
ownedObjectType).type().name()));
+ entityIds.add(EntityIdService.getEntityId(entity, ownedObjectType));
+ }
+ entityIds.sort(Comparator.naturalOrder());
+ String metadataObjectType =
+ NameIdentifierUtil.toMetadataObject(ownedObjects.get(0),
ownedObjectType).type().name();
+
+ List<OwnerRelForDeletion> deletions = new ArrayList<>(entityIds.size());
+ List<OwnerRelPO> ownerRelPOs = new ArrayList<>(entityIds.size());
+ for (Long entityId : entityIds) {
+ deletions.add(new OwnerRelForDeletion(entityId, metadataObjectType));
ownerRelPOs.add(
POConverters.initializeOwnerRelPOsWithVersion(
metalakeId, ownerType.name(), ownerId, ownedObjectType.name(),
entityId));
}
- SessionUtils.doMultipleWithCommit(
+ assignOwner(
+ metalake,
+ metalakeId,
+ ownerIdent,
+ ownerType,
+ ownerId,
+ entityIds,
+ ownedObjectType,
() ->
SessionUtils.doWithoutCommit(
OwnerMetaMapper.class,
@@ -240,4 +264,95 @@ public class OwnerMetaService {
SessionUtils.doWithoutCommit(
OwnerMetaMapper.class, mapper ->
mapper.batchInsertOwnerRels(ownerRelPOs)));
}
+
+ /**
+ * Serializes assignments on the owned object's row, including the first
assignment when no owner
+ * relation exists. The metalake is locked first to fence cascade deletion;
object rows are locked
+ * in stable ID order; the principal is then fenced before owner relations
change.
+ */
+ private void assignOwner(
+ String metalake,
+ long metalakeId,
+ NameIdentifier owner,
+ Entity.EntityType ownerType,
+ long ownerId,
+ List<Long> entityIds,
+ Entity.EntityType ownedObjectType,
+ Runnable retirePreviousOwners,
+ Runnable insertOwners) {
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ lockMetalakeForOwnerWrite(
+ metalake, metalakeId, ownedObjectType ==
Entity.EntityType.METALAKE),
+ () -> lockOwnedObjectsForOwnerWrite(entityIds, ownedObjectType,
metalakeId),
+ () -> lockPrincipalForOwnerWrite(owner, ownerType, ownerId,
metalakeId),
+ retirePreviousOwners,
+ insertOwners);
+ }
+
+ private void lockMetalakeForOwnerWrite(String metalake, long metalakeId,
boolean exclusive) {
+ OccWriteSupport.lockParentForChildWrite(
+ metalake,
+ Entity.EntityType.METALAKE,
+ () ->
+ SessionUtils.getWithoutCommit(
+ MetalakeMetaMapper.class,
+ mapper ->
+ exclusive
+ ? mapper.selectMetalakeMetaByIdForUpdate(metalakeId)
+ : mapper.selectMetalakeMetaByIdForShare(metalakeId)),
+ null,
+ current -> Objects.equals(current.getMetalakeName(), metalake));
+ }
+
+ private void lockOwnedObjectsForOwnerWrite(
+ List<Long> entityIds, Entity.EntityType entityType, long metalakeId) {
+ if (entityType == Entity.EntityType.METALAKE) {
+ return;
+ }
+ for (Long entityId : entityIds) {
+ OccWriteSupport.lockParentForChildWrite(
+ String.valueOf(entityId),
+ entityType,
+ () ->
+ SessionUtils.getWithoutCommit(
+ OwnerMetaMapper.class,
+ mapper ->
+ mapper.selectMetadataObjectIdForUpdate(entityId,
metalakeId, entityType)),
+ null,
+ current -> Objects.equals(current, entityId));
+ }
+ }
+
+ private void lockPrincipalForOwnerWrite(
+ NameIdentifier owner, Entity.EntityType ownerType, long ownerId, long
metalakeId) {
+ switch (ownerType) {
+ case USER:
+ OccWriteSupport.lockParentForChildWrite(
+ owner.name(),
+ ownerType,
+ () ->
+ SessionUtils.getWithoutCommit(
+ UserMetaMapper.class, mapper ->
mapper.selectUserMetaByIdForShare(ownerId)),
+ null,
+ current ->
+ Objects.equals(current.getMetalakeId(), metalakeId)
+ && Objects.equals(current.getUserName(), owner.name()));
+ return;
+ case GROUP:
+ OccWriteSupport.lockParentForChildWrite(
+ owner.name(),
+ ownerType,
+ () ->
+ SessionUtils.getWithoutCommit(
+ GroupMetaMapper.class, mapper ->
mapper.selectGroupMetaByIdForShare(ownerId)),
+ null,
+ current ->
+ Objects.equals(current.getMetalakeId(), metalakeId)
+ && Objects.equals(current.getGroupName(), owner.name()));
+ return;
+ default:
+ throw new IllegalArgumentException("Unsupported owner type: " +
ownerType);
+ }
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
b/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
index ba06be927c..c758614778 100644
--- a/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
+++ b/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
@@ -126,6 +126,81 @@ public class TestSQLScripts extends TestJDBCBackend {
}
}
+ /**
+ * The owner unique key allows one live row per (owner, object). Rows left
by concurrent
+ * assignments are merged during the upgrade: the newest live row (largest
id) stays, while older
+ * ones are soft-deleted.
+ */
+ @TestTemplate
+ public void testUpgradeToTwoZeroMergesDuplicateLiveOwners() throws
SQLException, IOException {
+ String gravitinoHome = System.getenv("GRAVITINO_HOME");
+ Assertions.assertNotNull(gravitinoHome, "GRAVITINO_HOME environment
variable is not set");
+ Path scriptDir = Path.of(gravitinoHome, "scripts",
backendType.toLowerCase());
+ String suffix = "-" + backendType.toLowerCase() + ".sql";
+ dropAllTables();
+ executeScript(scriptDir.resolve("schema-1.3.0" + suffix).toFile());
+
+ String insert =
+ "INSERT INTO owner_meta (id, metalake_id, owner_id, owner_type,
metadata_object_id,"
+ + " metadata_object_type, audit_info, current_version,
last_version, deleted_at,"
+ + " updated_at) VALUES (%d, 1, %d, 'USER', %d, 'CATALOG', '{}', 1,
1, %d, 0)";
+ List<String> rows =
+ List.of(
+ // Three owners of object 10 leave two rows to retire in the same
statement.
+ String.format(insert, 1, 100, 10, 0),
+ String.format(insert, 2, 200, 10, 0),
+ String.format(insert, 3, 300, 10, 0),
+ // Historical rows may already share a deletion timestamp.
+ String.format(insert, 4, 100, 20, 0),
+ String.format(insert, 5, 200, 20, 5),
+ String.format(insert, 6, 300, 20, 5),
+ // Object 10 as a SCHEMA is a different object.
+ String.format(insert, 7, 400, 10, 0).replace("'CATALOG'",
"'SCHEMA'"));
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement()) {
+ for (String row : rows) {
+ statement.execute(row);
+ }
+ }
+
+ executeScript(scriptDir.resolve("upgrade-1.3.0-to-2.0.0" +
suffix).toFile());
+
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet live =
+ statement.executeQuery("SELECT id FROM owner_meta WHERE deleted_at
= 0 ORDER BY id")) {
+ List<Long> liveIds = new ArrayList<>();
+ while (live.next()) {
+ liveIds.add(live.getLong(1));
+ }
+ Assertions.assertEquals(List.of(3L, 4L, 7L), liveIds);
+ }
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet retired =
+ statement.executeQuery("SELECT deleted_at, updated_at FROM
owner_meta WHERE id = 1")) {
+ Assertions.assertTrue(retired.next());
+ Assertions.assertTrue(retired.getLong(1) > 0, "older duplicate must be
soft-deleted");
+ Assertions.assertEquals(retired.getLong(1), retired.getLong(2));
+ }
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement()) {
+ // Historical rows can share a deletion timestamp after the upgrade.
+ statement.execute(String.format(insert, 8, 500, 20, 5));
+ // The existing key still rejects a second live row for the same owner
and object.
+ Assertions.assertThrows(
+ SQLException.class, () -> statement.execute(String.format(insert, 9,
100, 20, 0)));
+ }
+ }
+
private void executeScript(File scriptFile) throws IOException, SQLException
{
List<String> ddls = extractStatements(scriptFile.toPath());
try (SqlSession sqlSession =
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOwnerAssignmentWrites.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOwnerAssignmentWrites.java
new file mode 100644
index 0000000000..9c0e92dc75
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestOwnerAssignmentWrites.java
@@ -0,0 +1,529 @@
+/*
+ * 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.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+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.ArrayList;
+import java.util.List;
+import java.util.Optional;
+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 javax.annotation.Nullable;
+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.meta.CatalogEntity;
+import org.apache.gravitino.meta.GroupEntity;
+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.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;
+
+/**
+ * Races between owner assignment and deletion of the owned object or owner
principal, and between
+ * two assignments on the same object. Every scenario is driven by a real
transaction held open on
+ * one thread while the contender runs on another, so the assertions describe
what two Gravitino
+ * servers sharing one database would observe.
+ */
+class TestOwnerAssignmentWrites extends TestJDBCBackend {
+ private static final String METALAKE = "owner_write_metalake";
+
+ @TestTemplate
+ void testAssignmentWaitsForUncommittedPrincipalDeleteAndFails() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ for (boolean group : List.of(false, true)) {
+ CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned_" + group);
+ NameIdentifier principal = insertPrincipal(group, "deleted_" + group);
+ long id = principalId(group, principal);
+ Throwable failure =
+ whileTransactionHeld(
+ () -> deletePrincipal(group, principal),
+ () -> setOwner(owned, principal, type(group)));
+ assertInstanceOf(NoSuchEntityException.class, failure);
+ assertFalse(owner(owned).isPresent());
+ assertEquals(0, liveOwnerRows(id));
+ }
+ }
+
+ @TestTemplate
+ void testAssignmentWaitsForUncommittedOwnedObjectDeleteAndFails() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned");
+ NameIdentifier principal = insertPrincipal(false, "owner");
+ Throwable failure =
+ whileTransactionHeld(
+ () ->
+ assertTrue(
+
CatalogMetaService.getInstance().deleteCatalog(owned.nameIdentifier(), false)),
+ () -> setOwner(owned, principal, Entity.EntityType.USER));
+ assertInstanceOf(NoSuchEntityException.class, failure);
+ assertEquals(0, liveOwnerRowsForObject(owned));
+ }
+
+ @TestTemplate
+ void testPrincipalDeleteWaitsForUncommittedAssignmentAndCleansIt() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ for (boolean group : List.of(false, true)) {
+ CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned_" + group);
+ NameIdentifier principal = insertPrincipal(group, "deleted_" + group);
+ long id = principalId(group, principal);
+ assertNull(
+ whileTransactionHeld(
+ () -> setOwner(owned, principal, type(group)),
+ () -> deletePrincipal(group, principal)));
+ assertFalse(owner(owned).isPresent());
+ assertEquals(0, liveOwnerRows(id));
+ }
+ }
+
+ @TestTemplate
+ void testAssignmentRejectsSameNameReplacementObservedByStaleId() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ for (boolean group : List.of(false, true)) {
+ CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned_" + group);
+ NameIdentifier principal = insertPrincipal(group, "recreated_" + group);
+ long staleId = principalId(group, principal);
+ Throwable failure =
+ whileTransactionHeld(
+ () -> {
+ deletePrincipal(group, principal);
+ insertPrincipal(group, principal.name());
+ },
+ () -> setOwner(owned, principal, type(group)));
+ assertInstanceOf(NoSuchEntityException.class, failure);
+ assertEquals(0, liveOwnerRows(staleId));
+ assertFalse(owner(owned).isPresent());
+
+ // A fresh lookup resolves the replacement and succeeds.
+ setOwner(owned, principal, type(group));
+ assertEquals(principal.name(), ownerName(owned));
+ assertEquals(1, liveOwnerRows(principalId(group, principal)));
+ }
+ }
+
+ @TestTemplate
+ void testConcurrentInitialAssignmentsLeaveExactlyOneLiveOwner() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned");
+ NameIdentifier first = insertPrincipal(false, "first");
+ NameIdentifier second = insertPrincipal(true, "second");
+ assertNull(
+ whileTransactionHeld(
+ () -> setOwner(owned, first, Entity.EntityType.USER),
+ () -> setOwner(owned, second, Entity.EntityType.GROUP),
+ true));
+ assertEquals(1, liveOwnerRowsForObject(owned));
+ assertEquals("second", ownerName(owned));
+ }
+
+ @TestTemplate
+ void testConcurrentMetalakeAssignmentsSerializeOnTheMetalakeRow() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ NameIdentifier owned = NameIdentifier.of(METALAKE);
+ NameIdentifier first = insertPrincipal(false, "first");
+ NameIdentifier second = insertPrincipal(true, "second");
+ assertNull(
+ whileTransactionHeld(
+ () ->
+ OwnerMetaService.getInstance()
+ .setOwner(owned, Entity.EntityType.METALAKE, first,
Entity.EntityType.USER),
+ () ->
+ OwnerMetaService.getInstance()
+ .setOwner(owned, Entity.EntityType.METALAKE, second,
Entity.EntityType.GROUP)));
+ long metalakeId =
MetalakeMetaService.getInstance().getMetalakeIdByName(METALAKE);
+ assertEquals(
+ 1,
+ queryLong(
+ "SELECT COUNT(*) FROM owner_meta WHERE metadata_object_id = "
+ + metalakeId
+ + " AND metadata_object_type = 'METALAKE' AND deleted_at =
0"));
+ assertEquals(
+ "second",
+ assertInstanceOf(
+ GroupEntity.class,
+ OwnerMetaService.getInstance()
+ .getOwner(owned, Entity.EntityType.METALAKE)
+ .orElseThrow())
+ .name());
+ }
+
+ @TestTemplate
+ void testConcurrentReassignmentsSerializeOnTheExistingOwnerRow() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ CatalogEntity owned = createAndInsertCatalog(METALAKE, "owned");
+ NameIdentifier initial = insertPrincipal(false, "initial");
+ NameIdentifier first = insertPrincipal(false, "first");
+ NameIdentifier second = insertPrincipal(true, "second");
+ setOwner(owned, initial, Entity.EntityType.USER);
+ assertNull(
+ whileTransactionHeld(
+ () -> setOwner(owned, first, Entity.EntityType.USER),
+ () -> setOwner(owned, second, Entity.EntityType.GROUP),
+ true));
+ assertEquals(1, liveOwnerRowsForObject(owned));
+ assertEquals("second", ownerName(owned));
+ }
+
+ @TestTemplate
+ void testBatchAssignmentRollsBackWhenPrincipalIsDeleted() throws Exception {
+ createAndInsertMakeLake(METALAKE);
+ for (boolean group : List.of(false, true)) {
+ NameIdentifier previous = insertPrincipal(false, "previous_" + group);
+ NameIdentifier principal = insertPrincipal(group, "deleted_" + group);
+ long id = principalId(group, principal);
+ List<CatalogEntity> owned = new ArrayList<>();
+ for (int i = 0; i < 3; i++) {
+ CatalogEntity catalog = createAndInsertCatalog(METALAKE, "owned_" +
group + "_" + i);
+ owned.add(catalog);
+ if (i > 0) {
+ setOwner(catalog, previous, Entity.EntityType.USER);
+ }
+ }
+ Throwable failure =
+ whileTransactionHeld(
+ () -> deletePrincipal(group, principal),
+ () -> batchSetOwners(owned, principal, type(group)));
+ assertInstanceOf(NoSuchEntityException.class, failure);
+ assertFalse(SessionUtils.isInTransaction());
+ assertFalse(owner(owned.get(0)).isPresent());
+ assertEquals(previous.name(), ownerName(owned.get(1)));
+ assertEquals(previous.name(), ownerName(owned.get(2)));
+ assertEquals(0, liveOwnerRows(id));
+ }
+ }
+
+ @TestTemplate
+ void testBatchAssignmentSucceedsAfterPrincipalDeleteRollsBack() throws
Exception {
+ createAndInsertMakeLake(METALAKE);
+ for (boolean group : List.of(false, true)) {
+ NameIdentifier principal = insertPrincipal(group, "kept_" + group);
+ List<CatalogEntity> owned =
+ List.of(
+ createAndInsertCatalog(METALAKE, "owned_" + group + "_0"),
+ createAndInsertCatalog(METALAKE, "owned_" + group + "_1"));
+ assertNull(
+ whileTransactionHeld(
+ () -> deletePrincipal(group, principal),
+ () -> batchSetOwners(owned, principal, type(group)),
+ () -> {},
+ false));
+ for (CatalogEntity catalog : owned) {
+ assertEquals(principal.name(), ownerName(catalog));
+ }
+ assertEquals(2, liveOwnerRows(principalId(group, principal)));
+ }
+ }
+
+ private void setOwner(CatalogEntity owned, NameIdentifier owner,
Entity.EntityType ownerType) {
+ OwnerMetaService.getInstance()
+ .setOwner(owned.nameIdentifier(), Entity.EntityType.CATALOG, owner,
ownerType);
+ }
+
+ private void batchSetOwners(
+ List<CatalogEntity> owned, NameIdentifier owner, Entity.EntityType
ownerType) {
+ List<NameIdentifier> identifiers = new ArrayList<>();
+ for (CatalogEntity catalog : owned) {
+ identifiers.add(catalog.nameIdentifier());
+ }
+ OwnerMetaService.getInstance()
+ .batchSetOwners(identifiers, Entity.EntityType.CATALOG, owner,
ownerType);
+ }
+
+ private Optional<Entity> owner(CatalogEntity owned) {
+ return OwnerMetaService.getInstance()
+ .getOwner(owned.nameIdentifier(), Entity.EntityType.CATALOG);
+ }
+
+ private String ownerName(CatalogEntity owned) {
+ Entity entity = owner(owned).orElseThrow(() -> new AssertionError("No
owner"));
+ return entity instanceof UserEntity
+ ? ((UserEntity) entity).name()
+ : ((GroupEntity) entity).name();
+ }
+
+ private NameIdentifier insertPrincipal(boolean group, String name) throws
IOException {
+ long id = RandomIdGenerator.INSTANCE.nextId();
+ if (group) {
+ GroupMetaService.getInstance()
+ .insertGroup(
+ createGroupEntity(
+ id, AuthorizationUtils.ofGroupNamespace(METALAKE), name,
AUDIT_INFO, null, null),
+ false);
+ return AuthorizationUtils.ofGroup(METALAKE, name);
+ }
+ UserMetaService.getInstance()
+ .insertUser(
+ createUserEntity(id, AuthorizationUtils.ofUserNamespace(METALAKE),
name, AUDIT_INFO),
+ false);
+ return AuthorizationUtils.ofUser(METALAKE, name);
+ }
+
+ private void deletePrincipal(boolean group, NameIdentifier principal) {
+ if (group) {
+ assertTrue(GroupMetaService.getInstance().deleteGroup(principal));
+ } else {
+ assertTrue(UserMetaService.getInstance().deleteUser(principal));
+ }
+ }
+
+ private long principalId(boolean group, NameIdentifier principal) {
+ return EntityIdService.getEntityId(principal, type(group));
+ }
+
+ private Entity.EntityType type(boolean group) {
+ return group ? Entity.EntityType.GROUP : Entity.EntityType.USER;
+ }
+
+ private long liveOwnerRows(long ownerId) throws Exception {
+ return queryLong(
+ "SELECT COUNT(*) FROM owner_meta WHERE owner_id = " + ownerId + " AND
deleted_at = 0");
+ }
+
+ private long liveOwnerRowsForObject(CatalogEntity owned) throws Exception {
+ return queryLong(
+ "SELECT COUNT(*) FROM owner_meta WHERE metadata_object_id = "
+ + owned.id()
+ + " AND metadata_object_type = 'CATALOG' 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, () -> {}, true, false);
+ }
+
+ private Throwable whileTransactionHeld(
+ Executable holder, Executable contender, boolean standaloneContender)
throws Exception {
+ return whileTransactionHeld(holder, contender, () -> {}, true,
standaloneContender);
+ }
+
+ private Throwable whileTransactionHeld(
+ Executable holder, Executable contender, Executable beforeCompletion,
boolean commitHolder)
+ throws Exception {
+ return whileTransactionHeld(holder, contender, beforeCompletion,
commitHolder, false);
+ }
+
+ /**
+ * Runs {@code holder} inside a transaction held open on this thread, starts
{@code contender} on
+ * another thread, waits until the database reports the contender blocked on
the holder, then
+ * commits or rolls back the holder and returns the contender's failure, or
null when it
+ * committed.
+ *
+ * <p>By default the contender is wrapped in a transaction of its own so
that its session can be
+ * identified. A {@code standaloneContender} runs exactly as production
does, owning its
+ * transactions, which is what the owner assignment needs to replay a lost
race: its session is
+ * then not known in advance and the wait is recognised by the holder side
alone.
+ */
+ private Throwable whileTransactionHeld(
+ Executable holder,
+ Executable contender,
+ Executable beforeCompletion,
+ boolean commitHolder,
+ boolean standaloneContender)
+ throws Exception {
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ CompletableFuture<Long> started = new CompletableFuture<>();
+ if (standaloneContender) {
+ raiseDefaultLockTimeout();
+ }
+ SessionUtils.beginTransaction();
+ try {
+ long holderId = prepareTransaction();
+ Assertions.assertDoesNotThrow(holder);
+ Future<Throwable> result =
+ standaloneContender
+ ? executor.submit(
+ () -> {
+ try {
+ contender.execute();
+ return null;
+ } catch (Throwable failure) {
+ return failure;
+ }
+ })
+ : submitTransaction(executor, started, contender);
+ awaitBlockedBy(
+ result, standaloneContender ? null : 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();
+ }
+ }
+
+ /**
+ * H2 gives new sessions a one-second lock timeout. A standalone contender
opens its own sessions,
+ * so the database-wide default is raised instead of a per-session setting.
+ */
+ private void raiseDefaultLockTimeout() throws SQLException {
+ if (!"h2".equals(backendType)) {
+ return;
+ }
+ try (SqlSession session =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Statement statement = session.getConnection().createStatement()) {
+ statement.execute("SET DEFAULT_LOCK_TIMEOUT 30000");
+ }
+ }
+
+ /**
+ * Waits until the database reports a session blocked by the holder: the
given contender session
+ * when known, otherwise any session. On MySQL a waiter on the implicit lock
of an uncommitted
+ * insert is reported as blocked by itself rather than by the inserting
transaction, so that shape
+ * is accepted too.
+ */
+ private void awaitBlockedBy(Future<Throwable> result, @Nullable Long
contenderId, long holderId)
+ throws Exception {
+ String query;
+ switch (backendType) {
+ case "h2":
+ query =
+ "SELECT COUNT(*) FROM INFORMATION_SCHEMA.SESSIONS WHERE BLOCKER_ID
= "
+ + holderId
+ + (contenderId == null ? "" : " AND SESSION_ID = " +
contenderId);
+ 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 (b.PROCESSLIST_ID = "
+ + holderId
+ + " OR b.THREAD_ID = r.THREAD_ID)"
+ + (contenderId == null ? "" : " AND r.PROCESSLIST_ID = " +
contenderId);
+ break;
+ case "postgresql":
+ query =
+ "SELECT COUNT(*) FROM pg_stat_activity a WHERE "
+ + holderId
+ + " = ANY(pg_blocking_pids(a.pid))"
+ + (contenderId == null ? "" : " AND a.pid = " + contenderId);
+ 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;
+ }
+ });
+ }
+ }
+}
diff --git a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
index 947bef162d..2b1e63bf06 100644
--- a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
+++ b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
@@ -101,3 +101,16 @@ CREATE TABLE IF NOT EXISTS `semantic_model_version_info` (
KEY `idx_smvi_cid` (`catalog_id`),
KEY `idx_smvi_sid` (`schema_id`)
) ENGINE=InnoDB COMMENT 'semantic model version information';
+
+-- Merge duplicate live owners left by concurrent assignments: the newest live
row
+-- (largest id) wins, and older ones are soft-deleted.
+UPDATE `owner_meta`
+ SET `deleted_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND FROM
CURRENT_TIMESTAMP(3)) / 1000),
+ `updated_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND FROM
CURRENT_TIMESTAMP(3)) / 1000)
+ WHERE `deleted_at` = 0
+ AND `id` < (
+ SELECT MAX(d.`id`) FROM `owner_meta` d
+ WHERE d.`deleted_at` = 0
+ AND d.`metadata_object_id` = `owner_meta`.`metadata_object_id`
+ AND d.`metadata_object_type` = `owner_meta`.`metadata_object_type`
+ );
diff --git a/scripts/mysql/schema-2.0.0-mysql.sql
b/scripts/mysql/schema-2.0.0-mysql.sql
index 70505d3a9e..68e9f9cf23 100644
--- a/scripts/mysql/schema-2.0.0-mysql.sql
+++ b/scripts/mysql/schema-2.0.0-mysql.sql
@@ -325,7 +325,7 @@ CREATE TABLE IF NOT EXISTS `owner_meta` (
`deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'owner
relation deleted at',
`updated_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'updated at',
PRIMARY KEY (`id`),
- UNIQUE KEY `uk_ow_me_del` (`owner_id`, `metadata_object_id`,
`metadata_object_type`,`deleted_at`),
+ UNIQUE KEY `uk_ow_me_del` (`owner_id`, `metadata_object_id`,
`metadata_object_type`, `deleted_at`),
KEY `idx_oid` (`owner_id`),
KEY `idx_meid` (`metadata_object_id`),
KEY `idx_owner_meta_del_upd_obj` (`deleted_at`, `updated_at`,
`metadata_object_id`)
diff --git a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
index 678767d88d..4bad49e12d 100644
--- a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
+++ b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
@@ -182,3 +182,18 @@ CREATE TABLE IF NOT EXISTS `semantic_model_version_info` (
KEY `idx_smvi_cid` (`catalog_id`),
KEY `idx_smvi_sid` (`schema_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_bin COMMENT 'semantic
model version information';
+
+-- Merge duplicate live owners left by concurrent assignments: the newest live
row
+-- (largest id) wins, and older ones are soft-deleted.
+UPDATE `owner_meta` o
+ JOIN (
+ SELECT `metadata_object_id`, `metadata_object_type`, MAX(`id`) AS
keep_id
+ FROM `owner_meta`
+ WHERE `deleted_at` = 0
+ GROUP BY `metadata_object_id`, `metadata_object_type`
+ HAVING COUNT(*) > 1
+ ) d ON o.`metadata_object_id` = d.`metadata_object_id`
+ AND o.`metadata_object_type` = d.`metadata_object_type`
+ SET o.`deleted_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND
FROM CURRENT_TIMESTAMP(3)) / 1000),
+ o.`updated_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND
FROM CURRENT_TIMESTAMP(3)) / 1000)
+ WHERE o.`deleted_at` = 0 AND o.`id` <> d.keep_id;
diff --git a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
index 4a64d43798..774cec6b8e 100644
--- a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
+++ b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
@@ -150,3 +150,16 @@ COMMENT ON COLUMN
semantic_model_version_info.semantic_model_definition IS 'stru
COMMENT ON COLUMN semantic_model_version_info.properties IS 'semantic model
properties snapshot (JSON)';
COMMENT ON COLUMN semantic_model_version_info.audit_info IS 'semantic model
version audit info';
COMMENT ON COLUMN semantic_model_version_info.deleted_at IS 'version deleted
at';
+
+-- Merge duplicate live owners left by concurrent assignments: the newest live
row
+-- (largest id) wins, and older ones are soft-deleted.
+UPDATE owner_meta
+ SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS
BIGINT),
+ updated_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS
BIGINT)
+ WHERE deleted_at = 0
+ AND id < (
+ SELECT MAX(d.id) FROM owner_meta d
+ WHERE d.deleted_at = 0
+ AND d.metadata_object_id = owner_meta.metadata_object_id
+ AND d.metadata_object_type = owner_meta.metadata_object_type
+ );