This is an automated email from the ASF dual-hosted git repository.
jerryshao pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 69aa582d4d [#13311] fix(core): Clean up relations, owner and
privileges when deleting a tag or policy (#13321)
69aa582d4d is described below
commit 69aa582d4de8bdd7dbf034f0e671be5c054ade36
Author: Jerry Shao <[email protected]>
AuthorDate: Fri Sep 18 19:49:14 2026 +0800
[#13311] fix(core): Clean up relations, owner and privileges when deleting
a tag or policy (#13321)
### What changes were proposed in this pull request?
On `branch-1.3`, deleting a tag or a policy now removes everything that
references it, in the same transaction as the tag or policy row.
- `TagMetaService.deleteTag` resolves the tag id first, then
soft-deletes:
- the tag row;
- the tag's object relations (`tag_relation_meta`);
- the tag's owner;
- its role securable objects.
- `PolicyMetaService.deletePolicy` resolves the policy id first, then
soft-deletes:
- the policy row;
- the policy versions;
- the policy's object relations (`policy_relation_meta`);
- its owner;
- its securable objects.
- The versions and object relations are removed by id through three new
mapper methods: `softDeletePolicyVersionsByPolicyId`,
`softDeleteTagMetadataObjectRelsByTagId` and
`softDeletePolicyMetadataObjectRelsByPolicyId`. Each has a base
(MySQL/H2) and a PostgreSQL provider. This matches `main`.
- If the tag or policy row is already gone when it is deleted (renamed
or deleted concurrently after its id was read), the transaction rolls
back and the delete returns `false`. The cleanup therefore never touches
the rows of an entity that is still live under another name.
`main` already does this cleanup, but as part of the OCC work (#12781,
#12782), which can't be cherry-picked. This change ports it by hand.
### Why are the changes needed?
- **Deleting a tag:** the tag row was soft-deleted first. The relation
cleanup (`softDeleteTagMetadataObjectRelsByMetalakeAndTagName`) only
matches live tags (`tm.deleted_at = 0`), so it removed nothing. The
owner and securable objects were never removed.
- **Deleting a policy:**
- `softDeletePolicyMetadataObjectRelsByMetalakeAndPolicyName` was never
called.
- The version cleanup ran after the policy row was deleted and also
matches live policies only, so no version was ever retired.
- The owner and securable objects were never removed.
- `branch-1.3` has no orphan relation GC, so these rows stayed forever.
Fixed: #13311
Part of #13303
### Does this PR introduce _any_ user-facing change?
No API or configuration change. Deleting a tag or policy now also
removes its assignments, its owner and the privileges granted on it.
A delete that races with a concurrent delete-and-recreate of the same
name on another server can still leave the new entity's relations
behind, as before this change. `main` handles that case with OCC.
### How was this patch tested?
- Added unit tests (H2 locally). Both fail without the fix and check the
relation rows directly:
- `TestTagMetaService.testDeleteTagCleansEveryDependentRelation`
- `TestPolicyMetaService.testDeletePolicyCleansEveryDependentRelation`
(also checks that every policy version is retired)
-
`TestPolicyMetaService.testDeletePolicyKeepsSameNamePolicyInAnotherMetalake`
- `./gradlew :core:test -PskipITs` passes locally.
- The new SQL needs the MySQL and PostgreSQL backends in CI.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
---------
Co-authored-by: Claude Opus 5 <[email protected]>
---
.../mapper/PolicyMetadataObjectRelMapper.java | 5 +
.../PolicyMetadataObjectRelSQLProviderFactory.java | 5 +
.../relational/mapper/PolicyVersionMapper.java | 5 +
.../mapper/PolicyVersionSQLProviderFactory.java | 4 +
.../mapper/TagMetadataObjectRelMapper.java | 5 +
.../TagMetadataObjectRelSQLProviderFactory.java | 4 +
.../PolicyMetadataObjectRelBaseSQLProvider.java | 8 ++
.../base/PolicyVersionBaseSQLProvider.java | 8 ++
.../base/TagMetadataObjectRelBaseSQLProvider.java | 8 ++
.../PolicyMetadataObjectRelPostgreSQLProvider.java | 9 ++
.../PolicyVersionPostgreSQLProvider.java | 8 ++
.../TagMetadataObjectRelPostgreSQLProvider.java | 8 ++
.../relational/service/PolicyMetaService.java | 63 +++++++---
.../storage/relational/service/TagMetaService.java | 56 ++++++---
.../relational/service/TestPolicyMetaService.java | 135 +++++++++++++++++++++
.../relational/service/TestTagMetaService.java | 75 ++++++++++++
16 files changed, 377 insertions(+), 29 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelMapper.java
index 31226a0eeb..7c8a85e99f 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelMapper.java
@@ -103,6 +103,11 @@ public interface PolicyMetadataObjectRelMapper {
@Param("metadataObjectType") String metadataObjectType,
@Param("policyIds") List<Long> policyIds);
+ @UpdateProvider(
+ type = PolicyMetadataObjectRelSQLProviderFactory.class,
+ method = "softDeletePolicyMetadataObjectRelsByPolicyId")
+ Integer softDeletePolicyMetadataObjectRelsByPolicyId(@Param("policyId") Long
policyId);
+
@UpdateProvider(
type = PolicyMetadataObjectRelSQLProviderFactory.class,
method = "softDeletePolicyMetadataObjectRelsByMetalakeAndPolicyName")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelSQLProviderFactory.java
index 8b38bfb3e0..76cfbfcbe7 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetadataObjectRelSQLProviderFactory.java
@@ -89,6 +89,11 @@ public class PolicyMetadataObjectRelSQLProviderFactory {
metadataObjectId, metadataObjectType, policyIds);
}
+ public static String softDeletePolicyMetadataObjectRelsByPolicyId(
+ @Param("policyId") Long policyId) {
+ return
getProvider().softDeletePolicyMetadataObjectRelsByPolicyId(policyId);
+ }
+
public static String
softDeletePolicyMetadataObjectRelsByMetalakeAndPolicyName(
@Param("metalakeName") String metalakeName, @Param("policyName") String
policyName) {
return getProvider()
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionMapper.java
index 9bfbc62ea7..cd3f68401d 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionMapper.java
@@ -44,6 +44,11 @@ public interface PolicyVersionMapper {
Integer softDeletePolicyVersionByMetalakeAndPolicyName(
@Param("metalakeName") String metalakeName, @Param("policyName") String
policyName);
+ @UpdateProvider(
+ type = PolicyVersionSQLProviderFactory.class,
+ method = "softDeletePolicyVersionsByPolicyId")
+ Integer softDeletePolicyVersionsByPolicyId(@Param("policyId") Long policyId);
+
@UpdateProvider(
type = PolicyVersionSQLProviderFactory.class,
method = "deletePolicyVersionsByLegacyTimeline")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionSQLProviderFactory.java
index 9fbe4ddc79..ab6d4b82fb 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyVersionSQLProviderFactory.java
@@ -62,6 +62,10 @@ public class PolicyVersionSQLProviderFactory {
return
getProvider().softDeletePolicyVersionByMetalakeAndPolicyName(metalakeName,
policyName);
}
+ public static String softDeletePolicyVersionsByPolicyId(@Param("policyId")
Long policyId) {
+ return getProvider().softDeletePolicyVersionsByPolicyId(policyId);
+ }
+
public static String deletePolicyVersionsByLegacyTimeline(
@Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
return getProvider().deletePolicyVersionsByLegacyTimeline(legacyTimeline,
limit);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
index fda0ac5819..b2c5f5c04f 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
@@ -64,6 +64,11 @@ public interface TagMetadataObjectRelMapper {
@Param("metadataObjectType") String metadataObjectType,
@Param("tagIds") List<Long> tagIds);
+ @UpdateProvider(
+ type = TagMetadataObjectRelSQLProviderFactory.class,
+ method = "softDeleteTagMetadataObjectRelsByTagId")
+ Integer softDeleteTagMetadataObjectRelsByTagId(@Param("tagId") Long tagId);
+
@UpdateProvider(
type = TagMetadataObjectRelSQLProviderFactory.class,
method = "softDeleteTagMetadataObjectRelsByMetalakeAndTagName")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
index 46cfaa2f2d..e6b8c7dfe8 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
@@ -85,6 +85,10 @@ public class TagMetadataObjectRelSQLProviderFactory {
metadataObjectId, metadataObjectType, tagIds);
}
+ public static String softDeleteTagMetadataObjectRelsByTagId(@Param("tagId")
Long tagId) {
+ return getProvider().softDeleteTagMetadataObjectRelsByTagId(tagId);
+ }
+
public static String softDeleteTagMetadataObjectRelsByMetalakeAndTagName(
@Param("metalakeName") String metalakeName, @Param("tagName") String
tagName) {
return
getProvider().softDeleteTagMetadataObjectRelsByMetalakeAndTagName(metalakeName,
tagName);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetadataObjectRelBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetadataObjectRelBaseSQLProvider.java
index 3e4af8580b..9be6b28385 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetadataObjectRelBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetadataObjectRelBaseSQLProvider.java
@@ -133,6 +133,14 @@ public class PolicyMetadataObjectRelBaseSQLProvider {
+ "</script>";
}
+ public String
softDeletePolicyMetadataObjectRelsByPolicyId(@Param("policyId") Long policyId) {
+ return "UPDATE "
+ +
PolicyMetadataObjectRelMapper.POLICY_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE policy_id = #{policyId} AND deleted_at = 0";
+ }
+
public String softDeletePolicyMetadataObjectRelsByMetalakeAndPolicyName(
@Param("metalakeName") String metalakeName, @Param("policyName") String
policyName) {
return "UPDATE "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyVersionBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyVersionBaseSQLProvider.java
index b345890d49..f0dd9a548e 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyVersionBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyVersionBaseSQLProvider.java
@@ -74,6 +74,14 @@ public class PolicyVersionBaseSQLProvider {
+ " AND pv.deleted_at = 0";
}
+ public String softDeletePolicyVersionsByPolicyId(@Param("policyId") Long
policyId) {
+ return "UPDATE "
+ + POLICY_VERSION_TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE policy_id = #{policyId} AND deleted_at = 0";
+ }
+
public String deletePolicyVersionsByLegacyTimeline(
@Param("legacyTimeline") Long legacyTimeline, @Param("limit") int limit)
{
return "DELETE FROM "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
index b617ef9b68..31d627d184 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
@@ -130,6 +130,14 @@ public class TagMetadataObjectRelBaseSQLProvider {
+ "</script>";
}
+ public String softDeleteTagMetadataObjectRelsByTagId(@Param("tagId") Long
tagId) {
+ return "UPDATE "
+ + TagMetadataObjectRelMapper.TAG_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE tag_id = #{tagId} AND deleted_at = 0";
+ }
+
public String softDeleteTagMetadataObjectRelsByMetalakeAndTagName(
@Param("metalakeName") String metalakeName, @Param("tagName") String
tagName) {
return "UPDATE "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetadataObjectRelPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetadataObjectRelPostgreSQLProvider.java
index e365f359da..2a7a10e58e 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetadataObjectRelPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetadataObjectRelPostgreSQLProvider.java
@@ -40,6 +40,15 @@ public class PolicyMetadataObjectRelPostgreSQLProvider
private static final String DELETED_AT_NOW_EXPRESSION =
" CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT)";
+ @Override
+ public String softDeletePolicyMetadataObjectRelsByPolicyId(Long policyId) {
+ return "UPDATE "
+ + POLICY_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at ="
+ + DELETED_AT_NOW_EXPRESSION
+ + " WHERE policy_id = #{policyId} AND deleted_at = 0";
+ }
+
@Override
public String softDeletePolicyMetadataObjectRelsByMetalakeAndPolicyName(
String metalakeName, String policyName) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyVersionPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyVersionPostgreSQLProvider.java
index 539af52177..5960fd70b7 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyVersionPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyVersionPostgreSQLProvider.java
@@ -41,6 +41,14 @@ public class PolicyVersionPostgreSQLProvider extends
PolicyVersionBaseSQLProvide
+ " AND deleted_at = 0";
}
+ @Override
+ public String softDeletePolicyVersionsByPolicyId(Long policyId) {
+ return "UPDATE "
+ + POLICY_VERSION_TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE policy_id = #{policyId} AND deleted_at = 0";
+ }
+
@Override
public String deletePolicyVersionsByLegacyTimeline(Long legacyTimeline, int
limit) {
return "DELETE FROM "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
index 48097d7ad8..e2318ad492 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
@@ -37,6 +37,14 @@ import
org.apache.gravitino.storage.relational.mapper.provider.base.TagMetadataO
import org.apache.ibatis.annotations.Param;
public class TagMetadataObjectRelPostgreSQLProvider extends
TagMetadataObjectRelBaseSQLProvider {
+ @Override
+ public String softDeleteTagMetadataObjectRelsByTagId(Long tagId) {
+ return "UPDATE "
+ + TAG_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE tag_id = #{tagId} AND deleted_at = 0";
+ }
+
@Override
public String softDeleteTagMetadataObjectRelsByMetalakeAndTagName(
String metalakeName, String tagName) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
index 0c623f96c0..abe777740b 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
@@ -38,9 +38,11 @@ import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.meta.GenericEntity;
import org.apache.gravitino.meta.PolicyEntity;
import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
import org.apache.gravitino.storage.relational.mapper.PolicyMetaMapper;
import
org.apache.gravitino.storage.relational.mapper.PolicyMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.mapper.PolicyVersionMapper;
+import org.apache.gravitino.storage.relational.mapper.SecurableObjectMapper;
import org.apache.gravitino.storage.relational.po.PolicyMaxVersionPO;
import org.apache.gravitino.storage.relational.po.PolicyMetadataObjectRelPO;
import org.apache.gravitino.storage.relational.po.PolicyPO;
@@ -187,25 +189,56 @@ public class PolicyMetaService {
baseMetricName = "deletePolicy")
public boolean deletePolicy(NameIdentifier ident) {
String metalakeName = ident.namespace().level(0);
- int[] policyMetaDeletedCount = new int[] {0};
- int[] policyVersionDeletedCount = new int[] {0};
+ PolicyPO policyPO;
+ try {
+ policyPO = getPolicyPOByMetalakeAndName(metalakeName, ident.name());
+ } catch (NoSuchEntityException e) {
+ return false;
+ }
+ long policyId = policyPO.getPolicyId();
+ String policyType = MetadataObject.Type.POLICY.name();
- // We should delete meta and version info
- SessionUtils.doMultipleWithCommit(
- () ->
- policyMetaDeletedCount[0] =
+ // Everything that references the policy, including its versions, is
removed by the policy's id
+ // in the same transaction, so the cleanup does not depend on the policy
row still being live.
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ Integer deleted =
SessionUtils.getWithoutCommit(
PolicyMetaMapper.class,
mapper ->
-
mapper.softDeletePolicyByMetalakeAndPolicyName(metalakeName, ident.name())),
- () ->
- policyVersionDeletedCount[0] =
- SessionUtils.getWithoutCommit(
- PolicyVersionMapper.class,
- mapper ->
- mapper.softDeletePolicyVersionByMetalakeAndPolicyName(
- metalakeName, ident.name())));
- return policyMetaDeletedCount[0] + policyVersionDeletedCount[0] > 0;
+
mapper.softDeletePolicyByMetalakeAndPolicyName(metalakeName, ident.name()));
+ // The policy was renamed or deleted after its id was read. Roll
back, so the cleanup
+ // below never touches the rows of a policy that is still live
under another name.
+ if (deleted == null || deleted == 0) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.POLICY.name().toLowerCase(),
+ ident.name());
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ PolicyVersionMapper.class,
+ mapper ->
mapper.softDeletePolicyVersionsByPolicyId(policyId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ PolicyMetadataObjectRelMapper.class,
+ mapper ->
mapper.softDeletePolicyMetadataObjectRelsByPolicyId(policyId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ OwnerMetaMapper.class,
+ mapper ->
+
mapper.softDeleteOwnerRelByMetadataObjectIdAndType(policyId, policyType)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SecurableObjectMapper.class,
+ mapper ->
mapper.softDeleteObjectRelsByMetadataObject(policyId, policyType)));
+ } catch (NoSuchEntityException e) {
+ return false;
+ }
+
+ return true;
}
@Monitored(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
index b56f703d2c..0e7e356916 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TagMetaService.java
@@ -41,6 +41,8 @@ import org.apache.gravitino.exceptions.NoSuchTagException;
import org.apache.gravitino.meta.GenericEntity;
import org.apache.gravitino.meta.TagEntity;
import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.SecurableObjectMapper;
import org.apache.gravitino.storage.relational.mapper.TagMetaMapper;
import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.po.TagMetadataObjectRelPO;
@@ -146,26 +148,52 @@ public class TagMetaService {
@Monitored(metricsSource = GRAVITINO_RELATIONAL_STORE_METRIC_NAME,
baseMetricName = "deleteTag")
public boolean deleteTag(NameIdentifier identifier) {
String metalakeName = identifier.namespace().level(0);
- int[] tagDeletedCount = new int[] {0};
- int[] tagMetadataObjectRelDeletedCount = new int[] {0};
+ TagPO tagPO;
+ try {
+ tagPO = getTagPOByMetalakeAndName(metalakeName, identifier.name());
+ } catch (NoSuchEntityException e) {
+ return false;
+ }
+ long tagId = tagPO.getTagId();
+ String tagType = MetadataObject.Type.TAG.name();
- SessionUtils.doMultipleWithCommit(
- () ->
- tagDeletedCount[0] =
+ // Everything that references the tag is removed by the tag's id in the
same transaction, so
+ // the cleanup does not depend on the tag row still being live.
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ Integer deleted =
SessionUtils.getWithoutCommit(
TagMetaMapper.class,
mapper ->
mapper.softDeleteTagMetaByMetalakeAndTagName(
- metalakeName, identifier.name())),
- () ->
- tagMetadataObjectRelDeletedCount[0] =
- SessionUtils.getWithoutCommit(
- TagMetadataObjectRelMapper.class,
- mapper ->
-
mapper.softDeleteTagMetadataObjectRelsByMetalakeAndTagName(
- metalakeName, identifier.name())));
+ metalakeName, identifier.name()));
+ // The tag was renamed or deleted after its id was read. Roll
back, so the cleanup
+ // below never touches the rows of a tag that is still live under
another name.
+ if (deleted == null || deleted == 0) {
+ throw new NoSuchEntityException(
+ NoSuchEntityException.NO_SUCH_ENTITY_MESSAGE,
+ Entity.EntityType.TAG.name().toLowerCase(),
+ identifier.name());
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ TagMetadataObjectRelMapper.class,
+ mapper ->
mapper.softDeleteTagMetadataObjectRelsByTagId(tagId)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ OwnerMetaMapper.class,
+ mapper ->
mapper.softDeleteOwnerRelByMetadataObjectIdAndType(tagId, tagType)),
+ () ->
+ SessionUtils.doWithoutCommit(
+ SecurableObjectMapper.class,
+ mapper -> mapper.softDeleteObjectRelsByMetadataObject(tagId,
tagType)));
+ } catch (NoSuchEntityException e) {
+ return false;
+ }
- return tagDeletedCount[0] + tagMetadataObjectRelDeletedCount[0] > 0;
+ return true;
}
@Monitored(
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
index 5dd6056e41..a1bb055024 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
@@ -25,6 +25,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Lists;
import java.io.IOException;
import java.sql.Connection;
import java.sql.ResultSet;
@@ -40,6 +41,9 @@ import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.MetadataObject;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
+import org.apache.gravitino.authorization.AuthorizationUtils;
+import org.apache.gravitino.authorization.Privileges;
+import org.apache.gravitino.authorization.SecurableObjects;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.meta.BaseMetalake;
import org.apache.gravitino.meta.CatalogEntity;
@@ -47,9 +51,11 @@ import org.apache.gravitino.meta.FilesetEntity;
import org.apache.gravitino.meta.GenericEntity;
import org.apache.gravitino.meta.ModelEntity;
import org.apache.gravitino.meta.PolicyEntity;
+import org.apache.gravitino.meta.RoleEntity;
import org.apache.gravitino.meta.SchemaEntity;
import org.apache.gravitino.meta.TableEntity;
import org.apache.gravitino.meta.TopicEntity;
+import org.apache.gravitino.meta.UserEntity;
import org.apache.gravitino.policy.Policy;
import org.apache.gravitino.policy.PolicyContent;
import org.apache.gravitino.policy.PolicyContents;
@@ -1040,6 +1046,119 @@ public class TestPolicyMetaService extends
TestJDBCBackend {
return new EntitiesToTest(catalog, schema, table, topic, fileset, model);
}
+ @TestTemplate
+ public void testDeletePolicyCleansEveryDependentRelation() throws
IOException {
+ createAndInsertMakeLake(METALAKE_NAME);
+ CatalogEntity catalog = createAndInsertCatalog(METALAKE_NAME,
"catalog_policy_cascade");
+ PolicyMetaService policyMetaService = PolicyMetaService.getInstance();
+ PolicyEntity policy =
+ createPolicy(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofPolicy(METALAKE_NAME),
+ "policy_cascade",
+ AUDIT_INFO);
+ policyMetaService.insertPolicy(policy, false);
+ policyMetaService.associatePoliciesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new NameIdentifier[] {policy.nameIdentifier()},
+ new NameIdentifier[0]);
+
+ UserEntity user =
+ createUserEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ AuthorizationUtils.ofUserNamespace(METALAKE_NAME),
+ "user_policy_cascade",
+ AUDIT_INFO);
+ backend.insert(user, false);
+ OwnerMetaService.getInstance()
+ .setOwner(
+ policy.nameIdentifier(), Entity.EntityType.POLICY,
user.nameIdentifier(), user.type());
+
+ RoleEntity role =
+ createRoleEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ AuthorizationUtils.ofRoleNamespace(METALAKE_NAME),
+ "role_policy_cascade",
+ AUDIT_INFO,
+ Lists.newArrayList(
+ SecurableObjects.ofPolicy(
+ policy.name(),
Lists.newArrayList(Privileges.ApplyPolicy.allow()))),
+ null);
+ backend.insert(role, false);
+
+ String policyAsMetadataObject =
+ String.format("metadata_object_id = %d AND metadata_object_type =
'POLICY'", policy.id());
+ String policyAsSecurableObject =
+ String.format("metadata_object_id = %d AND type = 'POLICY'",
policy.id());
+ assertEquals(1, countActivePolicyRel(policy.id()));
+ assertEquals(1, countActiveRows("owner_meta", policyAsMetadataObject));
+ assertEquals(1, countActiveRows("role_meta_securable_object",
policyAsSecurableObject));
+ assertFalse(listPolicyVersions(policy.id()).isEmpty());
+ assertTrue(listPolicyVersions(policy.id()).values().stream().allMatch(d ->
d == 0L));
+
+ assertTrue(policyMetaService.deletePolicy(policy.nameIdentifier()));
+
+ // Every version of the policy is retired with it.
+ assertTrue(listPolicyVersions(policy.id()).values().stream().allMatch(d ->
d > 0L));
+
+ assertEquals(0, countActivePolicyRel(policy.id()));
+ assertEquals(0, countActiveRows("owner_meta", policyAsMetadataObject));
+ assertEquals(0, countActiveRows("role_meta_securable_object",
policyAsSecurableObject));
+ assertTrue(
+ policyMetaService
+ .listPoliciesForMetadataObject(catalog.nameIdentifier(),
catalog.type())
+ .isEmpty());
+ assertFalse(policyMetaService.deletePolicy(policy.nameIdentifier()));
+ }
+
+ @TestTemplate
+ public void testDeletePolicyKeepsSameNamePolicyInAnotherMetalake() throws
IOException {
+ String anotherMetalakeName = METALAKE_NAME + "_another";
+ createAndInsertMakeLake(METALAKE_NAME);
+ createAndInsertMakeLake(anotherMetalakeName);
+ CatalogEntity catalog = createAndInsertCatalog(METALAKE_NAME,
"catalog_same_name");
+ CatalogEntity anotherCatalog = createAndInsertCatalog(anotherMetalakeName,
"catalog_same_name");
+ PolicyMetaService policyMetaService = PolicyMetaService.getInstance();
+
+ PolicyEntity policy =
+ createPolicy(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofPolicy(METALAKE_NAME),
+ "policy_same_name",
+ AUDIT_INFO);
+ PolicyEntity anotherPolicy =
+ createPolicy(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofPolicy(anotherMetalakeName),
+ "policy_same_name",
+ AUDIT_INFO);
+ policyMetaService.insertPolicy(policy, false);
+ policyMetaService.insertPolicy(anotherPolicy, false);
+ policyMetaService.associatePoliciesWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new NameIdentifier[] {policy.nameIdentifier()},
+ new NameIdentifier[0]);
+ policyMetaService.associatePoliciesWithMetadataObject(
+ anotherCatalog.nameIdentifier(),
+ anotherCatalog.type(),
+ new NameIdentifier[] {anotherPolicy.nameIdentifier()},
+ new NameIdentifier[0]);
+
+ assertTrue(policyMetaService.deletePolicy(policy.nameIdentifier()));
+
+ assertTrue(listPolicyVersions(policy.id()).values().stream().allMatch(d ->
d > 0L));
+ assertEquals(0, countActivePolicyRel(policy.id()));
+
+ // The policy with the same name in the other metalake is left untouched.
+ assertFalse(listPolicyVersions(anotherPolicy.id()).isEmpty());
+
assertTrue(listPolicyVersions(anotherPolicy.id()).values().stream().allMatch(d
-> d == 0L));
+ assertEquals(1, countActivePolicyRel(anotherPolicy.id()));
+ assertEquals(
+ anotherPolicy,
policyMetaService.getPolicyByIdentifier(anotherPolicy.nameIdentifier()));
+ }
+
private Integer countActivePolicyRel(Long policyId) {
try (SqlSession sqlSession =
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
@@ -1078,4 +1197,20 @@ public class TestPolicyMetaService extends
TestJDBCBackend {
throw new RuntimeException("SQL execution failed", se);
}
}
+
+ private int countActiveRows(String table, String condition) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet rs =
+ statement.executeQuery(
+ String.format(
+ "SELECT count(*) FROM %s WHERE %s AND deleted_at = 0",
table, condition))) {
+ Assertions.assertTrue(rs.next());
+ return rs.getInt(1);
+ } catch (SQLException se) {
+ throw new RuntimeException("SQL execution failed", se);
+ }
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
index 35cd9d7d2b..ed97827aef 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
@@ -36,6 +36,9 @@ import org.apache.gravitino.Entity;
import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
+import org.apache.gravitino.authorization.AuthorizationUtils;
+import org.apache.gravitino.authorization.Privileges;
+import org.apache.gravitino.authorization.SecurableObjects;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.meta.BaseMetalake;
import org.apache.gravitino.meta.CatalogEntity;
@@ -43,10 +46,12 @@ import org.apache.gravitino.meta.ColumnEntity;
import org.apache.gravitino.meta.FilesetEntity;
import org.apache.gravitino.meta.GenericEntity;
import org.apache.gravitino.meta.ModelEntity;
+import org.apache.gravitino.meta.RoleEntity;
import org.apache.gravitino.meta.SchemaEntity;
import org.apache.gravitino.meta.TableEntity;
import org.apache.gravitino.meta.TagEntity;
import org.apache.gravitino.meta.TopicEntity;
+import org.apache.gravitino.meta.UserEntity;
import org.apache.gravitino.rel.types.Types;
import org.apache.gravitino.storage.RandomIdGenerator;
import org.apache.gravitino.storage.relational.TestJDBCBackend;
@@ -1282,6 +1287,60 @@ public class TestTagMetaService extends TestJDBCBackend {
}
}
+ @TestTemplate
+ public void testDeleteTagCleansEveryDependentRelation() throws IOException {
+ createAndInsertMakeLake(METALAKE_NAME);
+ CatalogEntity catalog = createAndInsertCatalog(METALAKE_NAME,
"catalog_tag_cascade");
+ TagMetaService tagMetaService = TagMetaService.getInstance();
+ TagEntity tag = createAndInsertTagEntity("tag_cascade", "comment",
METALAKE_NAME);
+ tagMetaService.associateTagsWithMetadataObject(
+ catalog.nameIdentifier(),
+ catalog.type(),
+ new NameIdentifier[] {tag.nameIdentifier()},
+ new NameIdentifier[0]);
+
+ UserEntity user =
+ createUserEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ AuthorizationUtils.ofUserNamespace(METALAKE_NAME),
+ "user_tag_cascade",
+ AUDIT_INFO);
+ backend.insert(user, false);
+ OwnerMetaService.getInstance()
+ .setOwner(tag.nameIdentifier(), Entity.EntityType.TAG,
user.nameIdentifier(), user.type());
+
+ RoleEntity role =
+ createRoleEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ AuthorizationUtils.ofRoleNamespace(METALAKE_NAME),
+ "role_tag_cascade",
+ AUDIT_INFO,
+ Lists.newArrayList(
+ SecurableObjects.ofTag(
+ tag.name(),
Lists.newArrayList(Privileges.ApplyTag.allow()))),
+ null);
+ backend.insert(role, false);
+
+ String tagAsMetadataObject =
+ String.format("metadata_object_id = %d AND metadata_object_type =
'TAG'", tag.id());
+ String tagAsSecurableObject =
+ String.format("metadata_object_id = %d AND type = 'TAG'", tag.id());
+ assertEquals(1, countActiveTagRel(tag.id()));
+ assertEquals(1, countActiveRows("owner_meta", tagAsMetadataObject));
+ assertEquals(1, countActiveRows("role_meta_securable_object",
tagAsSecurableObject));
+
+ assertTrue(tagMetaService.deleteTag(tag.nameIdentifier()));
+
+ assertEquals(0, countActiveTagRel(tag.id()));
+ assertEquals(0, countActiveRows("owner_meta", tagAsMetadataObject));
+ assertEquals(0, countActiveRows("role_meta_securable_object",
tagAsSecurableObject));
+ assertTrue(
+ tagMetaService
+ .listTagsForMetadataObject(catalog.nameIdentifier(),
catalog.type())
+ .isEmpty());
+ assertFalse(tagMetaService.deleteTag(tag.nameIdentifier()));
+ }
+
private Integer countActiveTagRel(Long tagId) {
try (SqlSession sqlSession =
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
@@ -1301,4 +1360,20 @@ public class TestTagMetaService extends TestJDBCBackend {
throw new RuntimeException("SQL execution failed", se);
}
}
+
+ private int countActiveRows(String table, String condition) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet rs =
+ statement.executeQuery(
+ String.format(
+ "SELECT count(*) FROM %s WHERE %s AND deleted_at = 0",
table, condition))) {
+ Assertions.assertTrue(rs.next());
+ return rs.getInt(1);
+ } catch (SQLException se) {
+ throw new RuntimeException("SQL execution failed", se);
+ }
+ }
}