From a3785bfcd6fd4deaeb30369f2cc7e8a4b886c3c9 Mon Sep 17 00:00:00 2001 From: iting0321 Date: Sun, 19 Jul 2026 21:13:21 +0800 Subject: [PATCH 1/4] make entity drop and cleanup task creation atomic --- .../NonFunctionalBasePersistence.java | 8 ++++ .../jdbc/JdbcBasePersistenceImpl.java | 34 ++++++++++++++ ...anagerWithJdbcBasePersistenceImplTest.java | 5 +++ .../AtomicOperationMetaStoreManager.java | 33 +++++++++----- .../core/persistence/BasePersistence.java | 17 +++++++ .../AbstractTransactionalPersistence.java | 25 +++++++++++ ...apAtomicOperationMetaStoreManagerTest.java | 5 +++ .../BasePolarisMetaStoreManagerTest.java | 45 +++++++++++++++++++ 8 files changed, 160 insertions(+), 12 deletions(-) diff --git a/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java b/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java index f88a5db9d9f..6fd650fdbe4 100644 --- a/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java +++ b/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java @@ -82,6 +82,14 @@ public void deleteEntity(@NonNull PolarisCallContext callCtx, @NonNull PolarisBa throw unimplemented(); } + @Override + public void deleteEntityAndCreateEntities( + @Nonnull PolarisCallContext callCtx, + @Nonnull PolarisBaseEntity entityToDelete, + @Nonnull List entitiesToCreate) { + throw useMetaStoreManager("create/update/rename/delete"); + } + @Override public void deleteFromGrantRecords( @NonNull PolarisCallContext callCtx, @NonNull PolarisGrantRecord grantRec) { diff --git a/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java b/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java index 365cb8262c1..45b1c0b9a38 100644 --- a/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java +++ b/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java @@ -352,6 +352,40 @@ public void deleteEntity(@NonNull PolarisCallContext callCtx, @NonNull PolarisBa } } + @Override + public void deleteEntityAndCreateEntities( + @Nonnull PolarisCallContext callCtx, + @Nonnull PolarisBaseEntity entityToDelete, + @Nonnull List entitiesToCreate) { + ModelEntity modelEntity = ModelEntity.fromEntity(entityToDelete, schemaVersion); + Map params = + Map.of( + "id", + modelEntity.getId(), + "catalog_id", + modelEntity.getCatalogId(), + "realm_id", + realmId); + try { + datasourceOperations.runWithinTransaction( + connection -> { + datasourceOperations.execute( + connection, + QueryGenerator.generateDeleteQuery( + ModelEntity.getAllColumnNames(schemaVersion), ModelEntity.TABLE_NAME, params)); + for (PolarisBaseEntity entityToCreate : entitiesToCreate) { + persistEntity( + callCtx, entityToCreate, null, connection, datasourceOperations::execute); + } + return true; + }); + } catch (SQLException e) { + throw new RuntimeException( + String.format("Failed to delete entity and create entities due to %s", e.getMessage()), + e); + } + } + @Override public void deleteFromGrantRecords( @NonNull PolarisCallContext callCtx, @NonNull PolarisGrantRecord grantRec) { diff --git a/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java b/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java index 4ccee4483c9..689c53017d7 100644 --- a/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java +++ b/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java @@ -48,6 +48,11 @@ public abstract class AtomicMetastoreManagerWithJdbcBasePersistenceImplTest extends BasePolarisMetaStoreManagerTest { + @Test + void testCleanupTaskCreationFailureRollsBackEntityDrop() { + assertCleanupTaskCreationFailureRollsBackEntityDrop(); + } + protected DatabaseType databaseType() { return DatabaseType.H2; } diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java index 0570758f864..13d0a43878e 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java @@ -172,11 +172,19 @@ private EntityResult persistNewEntity( * @param callCtx call context * @param ms meta store * @param entity the entity being dropped - */ + */ private void dropEntity( @NonNull PolarisCallContext callCtx, @NonNull BasePersistence ms, @NonNull PolarisBaseEntity entity) { + dropEntity(callCtx, ms, entity, List.of()); + } + + private void dropEntity( + @NonNull PolarisCallContext callCtx, + @NonNull BasePersistence ms, + @NonNull PolarisBaseEntity entity, + @NonNull List entitiesToCreateAtomically) { // validate the entity type and subtype getDiagnostics().checkNotNull(entity, "unexpected_null_dpo"); @@ -187,7 +195,11 @@ private void dropEntity( // Remove the main entity itself first-thing; once its id no longer resolves successfully // it will be pruned out of any grant-record lookups anyways. - ms.deleteEntity(callCtx, entity); + if (entitiesToCreateAtomically.isEmpty()) { + ms.deleteEntity(callCtx, entity); + } else { + ms.deleteEntityAndCreateEntities(callCtx, entity, entitiesToCreateAtomically); + } // Best-effort cleanup - drop grant records, update grantRecordVersions for affected // other entities. @@ -1182,10 +1194,6 @@ public void deletePrincipalSecrets( } } - // simply delete that entity. Will be removed from entities_active, added to the - // entities_dropped and its version will be changed. - this.dropEntity(callCtx, ms, refreshEntityToDrop); - // if cleanup, schedule a cleanup task for the entity. do this here, so that drop and scheduling // the cleanup task is transactional. Otherwise, we'll be unable to schedule the cleanup task // later @@ -1208,15 +1216,16 @@ public void deletePrincipalSecrets( if (cleanupProperties != null) { taskEntityBuilder.internalPropertiesAsMap(cleanupProperties); } - // TODO: Add a way to create the task entities atomically with dropping the entity; - // in the meantime, if the server fails partway through a dropEntity, it's possible that - // the entity is dropped but we don't have any persisted task records that will carry - // out the cleanup. - PolarisBaseEntity taskEntity = taskEntityBuilder.build(); - createEntityIfNotExists(callCtx, null, taskEntity); + PolarisBaseEntity taskEntity = + prepareToPersistNewEntity(callCtx, ms, taskEntityBuilder.build()); + this.dropEntity(callCtx, ms, refreshEntityToDrop, List.of(taskEntity)); return new DropEntityResult(taskEntity.getId()); } + // simply delete that entity. Will be removed from entities_active, added to the + // entities_dropped and its version will be changed. + this.dropEntity(callCtx, ms, refreshEntityToDrop); + // done, return success return new DropEntityResult(); } diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java index 59081179f98..422eba25985 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java @@ -158,6 +158,23 @@ void writeToGrantRecords( */ void deleteEntity(@NonNull PolarisCallContext callCtx, @NonNull PolarisBaseEntity entity); + /** + * Delete one entity and create the supplied entities in one atomic persistence operation. If the + * operation succeeds, the deleted entity must be durably removed and every created entity must be + * durably visible; if it fails, none of those changes may be applied. + * + *

The created entities use the same create semantics as {@link + * #writeEntities(PolarisCallContext, List, List)} with {@code originalEntities == null}. + * + * @param callCtx call context + * @param entityToDelete entity to delete + * @param entitiesToCreate entities to create atomically with the delete + */ + void deleteEntityAndCreateEntities( + @Nonnull PolarisCallContext callCtx, + @Nonnull PolarisBaseEntity entityToDelete, + @Nonnull List entitiesToCreate); + /** * Delete the specified grantRecord to the grant_records table. * diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java index e79cc997746..1e56c526e6f 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java @@ -285,6 +285,31 @@ public void deleteEntity(@NonNull PolarisCallContext callCtx, @NonNull PolarisBa runActionInTransaction(callCtx, () -> this.deleteEntityInCurrentTxn(callCtx, entity)); } + /** {@inheritDoc} */ + @Override + public void deleteEntityAndCreateEntities( + @Nonnull PolarisCallContext callCtx, + @Nonnull PolarisBaseEntity entityToDelete, + @Nonnull List entitiesToCreate) { + runActionInTransaction( + callCtx, + () -> { + this.deleteEntityInCurrentTxn(callCtx, entityToDelete); + for (PolarisBaseEntity entityToCreate : entitiesToCreate) { + try { + this.checkConditionsForWriteEntityInCurrentTxn(callCtx, entityToCreate, null); + } catch (EntityAlreadyExistsException e) { + // Matching ids indicate a retried create whose id was already reserved for the same + // entity. Treat that as idempotent, matching writeEntities. + if (e.getExistingEntity().getId() != entityToCreate.getId()) { + throw e; + } + } + this.writeEntityInCurrentTxn(callCtx, entityToCreate, true, null); + } + }); + } + /** {@inheritDoc} */ @Override public void deleteFromGrantRecords( diff --git a/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java b/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java index d5d526527a1..47ee020d285 100644 --- a/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java +++ b/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java @@ -31,6 +31,11 @@ public class PolarisTreeMapAtomicOperationMetaStoreManagerTest extends BasePolarisMetaStoreManagerTest { + @Test + void testCleanupTaskCreationFailureRollsBackEntityDrop() { + assertCleanupTaskCreationFailureRollsBackEntityDrop(); + } + @Override public PolarisTestMetaStoreManager createPolarisTestMetaStoreManager() { PolarisDiagnostics diagServices = new PolarisDefaultDiagServiceImpl(); diff --git a/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java b/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java index d497714ba71..fafa0f7dbc3 100644 --- a/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java +++ b/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java @@ -272,6 +272,51 @@ protected void testDropEntities() { polarisTestMetaStoreManager.testDropEntities(); } + protected void assertCleanupTaskCreationFailureRollsBackEntityDrop() { + PolarisMetaStoreManager metaStoreManager = polarisTestMetaStoreManager.polarisMetaStoreManager; + PolarisCallContext callCtx = polarisTestMetaStoreManager.polarisCallContext; + PolarisBaseEntity entityToDrop = + polarisTestMetaStoreManager.createEntity( + null, + PolarisEntityType.PRINCIPAL_ROLE, + PolarisEntitySubType.NULL_SUBTYPE, + "principal_role_to_drop"); + TaskEntity conflictingTask = + createTask( + "entityCleanup_" + entityToDrop.getId(), + metaStoreManager.generateNewEntityId(callCtx).getId()); + metaStoreManager.createEntitiesIfNotExist(callCtx, null, List.of(conflictingTask)); + + List taskIdsBeforeDrop = + metaStoreManager + .listFullEntitiesAll( + callCtx, null, PolarisEntityType.TASK, PolarisEntitySubType.NULL_SUBTYPE) + .stream() + .map(PolarisBaseEntity::getId) + .toList(); + + Assertions.assertThatThrownBy( + () -> metaStoreManager.dropEntityIfExists(callCtx, null, entityToDrop, Map.of(), true)) + .isInstanceOf(EntityAlreadyExistsException.class); + + Assertions.assertThat( + metaStoreManager + .loadEntity( + callCtx, + entityToDrop.getCatalogId(), + entityToDrop.getId(), + entityToDrop.getType()) + .getEntity()) + .isNotNull() + .extracting(PolarisBaseEntity::getId) + .isEqualTo(entityToDrop.getId()); + Assertions.assertThat( + metaStoreManager.listFullEntitiesAll( + callCtx, null, PolarisEntityType.TASK, PolarisEntitySubType.NULL_SUBTYPE)) + .extracting(PolarisBaseEntity::getId) + .containsExactlyInAnyOrderElementsOf(taskIdsBeforeDrop); + } + /** Test that granting/revoking privileges works well */ @Test protected void testPrivileges() { From 489547aeaaa9ffb8dd392802687ad94cf4018d02 Mon Sep 17 00:00:00 2001 From: iting0321 Date: Sun, 19 Jul 2026 21:35:47 +0800 Subject: [PATCH 2/4] Preserve JDBC idempotency for atomic entity drops --- .../NonFunctionalBasePersistence.java | 6 +- .../jdbc/JdbcBasePersistenceImpl.java | 62 ++++++++++++++----- ...anagerWithJdbcBasePersistenceImplTest.java | 5 ++ .../AtomicOperationMetaStoreManager.java | 2 +- .../core/persistence/BasePersistence.java | 6 +- .../AbstractTransactionalPersistence.java | 6 +- .../BasePolarisMetaStoreManagerTest.java | 31 ++++++++++ 7 files changed, 91 insertions(+), 27 deletions(-) diff --git a/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java b/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java index 6fd650fdbe4..dd2c2fe8b54 100644 --- a/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java +++ b/persistence/nosql/persistence/metastore/src/main/java/org/apache/polaris/persistence/nosql/metastore/NonFunctionalBasePersistence.java @@ -84,9 +84,9 @@ public void deleteEntity(@NonNull PolarisCallContext callCtx, @NonNull PolarisBa @Override public void deleteEntityAndCreateEntities( - @Nonnull PolarisCallContext callCtx, - @Nonnull PolarisBaseEntity entityToDelete, - @Nonnull List entitiesToCreate) { + @NonNull PolarisCallContext callCtx, + @NonNull PolarisBaseEntity entityToDelete, + @NonNull List entitiesToCreate) { throw useMetaStoreManager("create/update/rename/delete"); } diff --git a/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java b/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java index 45b1c0b9a38..72b421a156d 100644 --- a/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java +++ b/persistence/relational-jdbc/src/main/java/org/apache/polaris/persistence/relational/jdbc/JdbcBasePersistenceImpl.java @@ -354,9 +354,9 @@ public void deleteEntity(@NonNull PolarisCallContext callCtx, @NonNull PolarisBa @Override public void deleteEntityAndCreateEntities( - @Nonnull PolarisCallContext callCtx, - @Nonnull PolarisBaseEntity entityToDelete, - @Nonnull List entitiesToCreate) { + @NonNull PolarisCallContext callCtx, + @NonNull PolarisBaseEntity entityToDelete, + @NonNull List entitiesToCreate) { ModelEntity modelEntity = ModelEntity.fromEntity(entityToDelete, schemaVersion); Map params = Map.of( @@ -374,6 +374,15 @@ public void deleteEntityAndCreateEntities( QueryGenerator.generateDeleteQuery( ModelEntity.getAllColumnNames(schemaVersion), ModelEntity.TABLE_NAME, params)); for (PolarisBaseEntity entityToCreate : entitiesToCreate) { + PolarisBaseEntity existingEntity = + lookupEntity( + connection, + entityToCreate.getCatalogId(), + entityToCreate.getId(), + entityToCreate.getTypeCode()); + if (existingEntity != null) { + continue; + } persistEntity( callCtx, entityToCreate, null, connection, datasourceOperations::execute); } @@ -455,11 +464,25 @@ public void deleteAll(@NonNull PolarisCallContext callCtx) { @Override public PolarisBaseEntity lookupEntity( @NonNull PolarisCallContext callCtx, long catalogId, long entityId, int typeCode) { + return getPolarisBaseEntity(entityLookupQuery(catalogId, entityId, typeCode)); + } + + private PolarisBaseEntity lookupEntity( + @NonNull Connection connection, long catalogId, long entityId, int typeCode) + throws SQLException { + return getPolarisBaseEntity( + datasourceOperations.executeSelect( + connection, + entityLookupQuery(catalogId, entityId, typeCode), + new ModelEntity(schemaVersion))); + } + + private QueryGenerator.PreparedQuery entityLookupQuery( + long catalogId, long entityId, int typeCode) { Map params = Map.of("catalog_id", catalogId, "id", entityId, "type_code", typeCode, "realm_id", realmId); - return getPolarisBaseEntity( - QueryGenerator.generateSelectQuery( - ModelEntity.getAllColumnNames(schemaVersion), ModelEntity.TABLE_NAME, params)); + return QueryGenerator.generateSelectQuery( + ModelEntity.getAllColumnNames(schemaVersion), ModelEntity.TABLE_NAME, params); } @Override @@ -489,23 +512,28 @@ public PolarisBaseEntity lookupEntityByName( @Nullable private PolarisBaseEntity getPolarisBaseEntity(QueryGenerator.PreparedQuery query) { try { - var results = datasourceOperations.executeSelect(query, new ModelEntity(schemaVersion)); - if (results.isEmpty()) { - return null; - } else if (results.size() > 1) { - throw new IllegalStateException( - String.format( - "More than one(%s) entities were found for a given type code : %s", - results.size(), results.getFirst().getTypeCode())); - } else { - return results.getFirst(); - } + return getPolarisBaseEntity( + datasourceOperations.executeSelect(query, new ModelEntity(schemaVersion))); } catch (SQLException e) { throw new RuntimeException( String.format("Failed to retrieve polaris entity due to %s", e.getMessage()), e); } } + @Nullable + private PolarisBaseEntity getPolarisBaseEntity(List results) { + if (results.isEmpty()) { + return null; + } else if (results.size() > 1) { + throw new IllegalStateException( + String.format( + "More than one(%s) entities were found for a given type code : %s", + results.size(), results.getFirst().getTypeCode())); + } else { + return results.getFirst(); + } + } + @NonNull @Override public List lookupEntities( diff --git a/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java b/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java index 689c53017d7..e2ade4e9e05 100644 --- a/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java +++ b/persistence/relational-jdbc/src/test/java/org/apache/polaris/persistence/relational/jdbc/AtomicMetastoreManagerWithJdbcBasePersistenceImplTest.java @@ -53,6 +53,11 @@ void testCleanupTaskCreationFailureRollsBackEntityDrop() { assertCleanupTaskCreationFailureRollsBackEntityDrop(); } + @Test + void testCleanupTaskCreationRetryIsIdempotent() { + assertCleanupTaskCreationRetryIsIdempotent(); + } + protected DatabaseType databaseType() { return DatabaseType.H2; } diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java index 13d0a43878e..bc2124c86a2 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java @@ -172,7 +172,7 @@ private EntityResult persistNewEntity( * @param callCtx call context * @param ms meta store * @param entity the entity being dropped - */ + */ private void dropEntity( @NonNull PolarisCallContext callCtx, @NonNull BasePersistence ms, diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java index 422eba25985..b47da95028e 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/BasePersistence.java @@ -171,9 +171,9 @@ void writeToGrantRecords( * @param entitiesToCreate entities to create atomically with the delete */ void deleteEntityAndCreateEntities( - @Nonnull PolarisCallContext callCtx, - @Nonnull PolarisBaseEntity entityToDelete, - @Nonnull List entitiesToCreate); + @NonNull PolarisCallContext callCtx, + @NonNull PolarisBaseEntity entityToDelete, + @NonNull List entitiesToCreate); /** * Delete the specified grantRecord to the grant_records table. diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java index 1e56c526e6f..4304aa33405 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java @@ -288,9 +288,9 @@ public void deleteEntity(@NonNull PolarisCallContext callCtx, @NonNull PolarisBa /** {@inheritDoc} */ @Override public void deleteEntityAndCreateEntities( - @Nonnull PolarisCallContext callCtx, - @Nonnull PolarisBaseEntity entityToDelete, - @Nonnull List entitiesToCreate) { + @NonNull PolarisCallContext callCtx, + @NonNull PolarisBaseEntity entityToDelete, + @NonNull List entitiesToCreate) { runActionInTransaction( callCtx, () -> { diff --git a/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java b/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java index fafa0f7dbc3..34bb459c2ad 100644 --- a/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java +++ b/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java @@ -317,6 +317,37 @@ protected void assertCleanupTaskCreationFailureRollsBackEntityDrop() { .containsExactlyInAnyOrderElementsOf(taskIdsBeforeDrop); } + protected void assertCleanupTaskCreationRetryIsIdempotent() { + PolarisMetaStoreManager metaStoreManager = polarisTestMetaStoreManager.polarisMetaStoreManager; + PolarisCallContext callCtx = polarisTestMetaStoreManager.polarisCallContext; + PolarisBaseEntity entityToDrop = + polarisTestMetaStoreManager.createEntity( + null, + PolarisEntityType.PRINCIPAL_ROLE, + PolarisEntitySubType.NULL_SUBTYPE, + "principal_role_to_drop_on_retry"); + + var dropResult = + metaStoreManager.dropEntityIfExists(callCtx, null, entityToDrop, Map.of(), true); + PolarisBaseEntity cleanupTask = + metaStoreManager + .loadEntity(callCtx, 0L, dropResult.getCleanupTaskId(), PolarisEntityType.TASK) + .getEntity(); + + Assertions.assertThat(cleanupTask).isNotNull(); + Assertions.assertThatCode( + () -> + callCtx + .getMetaStore() + .deleteEntityAndCreateEntities(callCtx, entityToDrop, List.of(cleanupTask))) + .doesNotThrowAnyException(); + Assertions.assertThat( + metaStoreManager + .loadEntity(callCtx, 0L, cleanupTask.getId(), PolarisEntityType.TASK) + .getEntity()) + .isEqualTo(cleanupTask); + } + /** Test that granting/revoking privileges works well */ @Test protected void testPrivileges() { From f74982fb3bb1e9bee76d8fb73769a4e88969df3e Mon Sep 17 00:00:00 2001 From: iting0321 Date: Sun, 19 Jul 2026 21:39:54 +0800 Subject: [PATCH 3/4] preserve idempotency for atomic entity creation retries --- .../transactional/AbstractTransactionalPersistence.java | 1 + .../PolarisTreeMapAtomicOperationMetaStoreManagerTest.java | 5 +++++ 2 files changed, 6 insertions(+) diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java index 4304aa33405..965a2cc4392 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/AbstractTransactionalPersistence.java @@ -304,6 +304,7 @@ public void deleteEntityAndCreateEntities( if (e.getExistingEntity().getId() != entityToCreate.getId()) { throw e; } + continue; } this.writeEntityInCurrentTxn(callCtx, entityToCreate, true, null); } diff --git a/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java b/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java index 47ee020d285..d052b2aadbf 100644 --- a/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java +++ b/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java @@ -36,6 +36,11 @@ void testCleanupTaskCreationFailureRollsBackEntityDrop() { assertCleanupTaskCreationFailureRollsBackEntityDrop(); } + @Test + void testCleanupTaskCreationRetryIsIdempotent() { + assertCleanupTaskCreationRetryIsIdempotent(); + } + @Override public PolarisTestMetaStoreManager createPolarisTestMetaStoreManager() { PolarisDiagnostics diagServices = new PolarisDefaultDiagServiceImpl(); From 92afc9eb31f313d5712a47567831f89f5bcb34c3 Mon Sep 17 00:00:00 2001 From: iting0321 Date: Tue, 21 Jul 2026 19:58:55 +0800 Subject: [PATCH 4/4] fix idempotent entity drop cleanup retries and conflict handling --- .../AtomicOperationMetaStoreManager.java | 26 ++++++- .../TransactionalMetaStoreManagerImpl.java | 53 +++++++++++--- .../PolarisTreeMapMetaStoreManagerTest.java | 11 +++ .../BasePolarisMetaStoreManagerTest.java | 17 +++-- .../service/admin/PolarisAdminService.java | 15 +++- .../catalog/iceberg/LocalIcebergCatalog.java | 12 ++++ .../admin/PolarisAdminServiceTest.java | 23 ++++++ .../AbstractLocalIcebergCatalogTest.java | 70 +++++++++++++++++++ 8 files changed, 202 insertions(+), 25 deletions(-) diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java index bc2124c86a2..2005105f13d 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/AtomicOperationMetaStoreManager.java @@ -1104,6 +1104,18 @@ public void deletePrincipalSecrets( // if this entity was not found, return failure if (refreshEntityToDrop == null) { + if (cleanup && entityToDrop.getType() != PolarisEntityType.POLICY) { + PolarisBaseEntity cleanupTask = + ms.lookupEntityByName( + callCtx, + PolarisEntityConstants.getNullId(), + PolarisEntityConstants.getNullId(), + PolarisEntityType.TASK.getCode(), + "entityCleanup_" + entityToDrop.getId()); + if (cleanupTask != null) { + return new DropEntityResult(cleanupTask.getId()); + } + } return new DropEntityResult(BaseResult.ReturnStatus.ENTITY_NOT_FOUND, null); } @@ -1209,7 +1221,7 @@ public void deletePrincipalSecrets( .propertiesAsMap(properties) .id(ms.generateNewId(callCtx)) .catalogId(0L) - .name("entityCleanup_" + entityToDrop.getId()) + .name("entityCleanup_" + refreshEntityToDrop.getId()) .typeCode(PolarisEntityType.TASK.getCode()) .subTypeCode(PolarisEntitySubType.NULL_SUBTYPE.getCode()) .createTimestamp(clock.millis()); @@ -1218,7 +1230,17 @@ public void deletePrincipalSecrets( } PolarisBaseEntity taskEntity = prepareToPersistNewEntity(callCtx, ms, taskEntityBuilder.build()); - this.dropEntity(callCtx, ms, refreshEntityToDrop, List.of(taskEntity)); + try { + this.dropEntity(callCtx, ms, refreshEntityToDrop, List.of(taskEntity)); + } catch (EntityAlreadyExistsException e) { + return new DropEntityResult( + BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS, + String.format( + "Existing entity id: '%s', type %s subtype %s", + e.getExistingEntity().getId(), + e.getExistingEntity().getTypeCode(), + e.getExistingEntity().getSubTypeCode())); + } return new DropEntityResult(taskEntity.getId()); } diff --git a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/TransactionalMetaStoreManagerImpl.java b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/TransactionalMetaStoreManagerImpl.java index 7402c6c0d54..eaa72850051 100644 --- a/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/TransactionalMetaStoreManagerImpl.java +++ b/polaris-core/src/main/java/org/apache/polaris/core/persistence/transactional/TransactionalMetaStoreManagerImpl.java @@ -51,6 +51,7 @@ import org.apache.polaris.core.entity.PolarisTaskConstants; import org.apache.polaris.core.entity.PrincipalEntity; import org.apache.polaris.core.persistence.BaseMetaStoreManager; +import org.apache.polaris.core.persistence.EntityAlreadyExistsException; import org.apache.polaris.core.persistence.PolarisMetaStoreManager; import org.apache.polaris.core.persistence.PolarisObjectMapperUtil; import org.apache.polaris.core.persistence.PolicyMappingAlreadyExistsException; @@ -1316,6 +1317,26 @@ public void deletePrincipalSecrets( // entity cannot be null getDiagnostics().checkNotNull(entityToDrop, "unexpected_null_entity"); + // Resolve the entity by id before validating its path so a retry after a successful commit can + // recover the cleanup task even though the dropped entity no longer resolves. + PolarisBaseEntity refreshEntityToDrop = + ms.lookupEntityInCurrentTxn( + callCtx, entityToDrop.getCatalogId(), entityToDrop.getId(), entityToDrop.getTypeCode()); + if (refreshEntityToDrop == null + && cleanup + && entityToDrop.getType() != PolarisEntityType.POLICY) { + PolarisBaseEntity cleanupTask = + ms.lookupEntityByNameInCurrentTxn( + callCtx, + PolarisEntityConstants.getNullId(), + PolarisEntityConstants.getNullId(), + PolarisEntityType.TASK.getCode(), + "entityCleanup_" + entityToDrop.getId()); + if (cleanupTask != null) { + return new DropEntityResult(cleanupTask.getId()); + } + } + // re-resolve everything including that entity PolarisEntityResolver resolver = new PolarisEntityResolver(getDiagnostics(), callCtx, ms, catalogPath, entityToDrop); @@ -1325,11 +1346,6 @@ public void deletePrincipalSecrets( return new DropEntityResult(BaseResult.ReturnStatus.CATALOG_PATH_CANNOT_BE_RESOLVED, null); } - // first find the entity to drop - PolarisBaseEntity refreshEntityToDrop = - ms.lookupEntityInCurrentTxn( - callCtx, entityToDrop.getCatalogId(), entityToDrop.getId(), entityToDrop.getTypeCode()); - // if this entity was not found, return failure if (refreshEntityToDrop == null) { return new DropEntityResult(BaseResult.ReturnStatus.ENTITY_NOT_FOUND, null); @@ -1443,7 +1459,12 @@ public void deletePrincipalSecrets( taskEntityBuilder.internalPropertiesAsMap(cleanupProperties); } PolarisBaseEntity taskEntity = taskEntityBuilder.build(); - createEntityIfNotExists(callCtx, ms, null, taskEntity); + EntityResult createTaskResult = createEntityIfNotExists(callCtx, ms, null, taskEntity); + if (!createTaskResult.isSuccess()) { + ms.rollback(); + return new DropEntityResult( + createTaskResult.getReturnStatus(), createTaskResult.getExtraInformation()); + } return new DropEntityResult(taskEntity.getId()); } @@ -1463,11 +1484,21 @@ public void deletePrincipalSecrets( TransactionalPersistence ms = ((TransactionalPersistence) callCtx.getMetaStore()); // need to run inside a read/write transaction - return ms.runInTransaction( - callCtx, - () -> - this.dropEntityIfExists( - callCtx, ms, catalogPath, entityToDrop, cleanupProperties, cleanup)); + try { + return ms.runInTransaction( + callCtx, + () -> + this.dropEntityIfExists( + callCtx, ms, catalogPath, entityToDrop, cleanupProperties, cleanup)); + } catch (EntityAlreadyExistsException e) { + return new DropEntityResult( + BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS, + String.format( + "Existing entity id: '%s', type %s subtype %s", + e.getExistingEntity().getId(), + e.getExistingEntity().getTypeCode(), + e.getExistingEntity().getSubTypeCode())); + } } /** diff --git a/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapMetaStoreManagerTest.java b/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapMetaStoreManagerTest.java index 4ecfa027f53..4dc48009bc6 100644 --- a/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapMetaStoreManagerTest.java +++ b/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapMetaStoreManagerTest.java @@ -26,9 +26,20 @@ import org.apache.polaris.core.persistence.transactional.TransactionalMetaStoreManagerImpl; import org.apache.polaris.core.persistence.transactional.TreeMapMetaStore; import org.apache.polaris.core.persistence.transactional.TreeMapTransactionalPersistenceImpl; +import org.junit.jupiter.api.Test; import org.mockito.Mockito; public class PolarisTreeMapMetaStoreManagerTest extends BasePolarisMetaStoreManagerTest { + @Test + void testCleanupTaskCreationFailureRollsBackEntityDrop() { + assertCleanupTaskCreationFailureRollsBackEntityDrop(); + } + + @Test + void testCleanupTaskCreationRetryIsIdempotent() { + assertCleanupTaskCreationRetryIsIdempotent(); + } + @Override public PolarisTestMetaStoreManager createPolarisTestMetaStoreManager() { PolarisDiagnostics diagServices = new PolarisDefaultDiagServiceImpl(); diff --git a/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java b/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java index 34bb459c2ad..9585099bf91 100644 --- a/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java +++ b/polaris-core/src/testFixtures/java/org/apache/polaris/core/persistence/BasePolarisMetaStoreManagerTest.java @@ -295,9 +295,10 @@ protected void assertCleanupTaskCreationFailureRollsBackEntityDrop() { .map(PolarisBaseEntity::getId) .toList(); - Assertions.assertThatThrownBy( - () -> metaStoreManager.dropEntityIfExists(callCtx, null, entityToDrop, Map.of(), true)) - .isInstanceOf(EntityAlreadyExistsException.class); + Assertions.assertThat( + metaStoreManager.dropEntityIfExists(callCtx, null, entityToDrop, Map.of(), true)) + .extracting(BaseResult::getReturnStatus) + .isEqualTo(BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS); Assertions.assertThat( metaStoreManager @@ -335,12 +336,10 @@ protected void assertCleanupTaskCreationRetryIsIdempotent() { .getEntity(); Assertions.assertThat(cleanupTask).isNotNull(); - Assertions.assertThatCode( - () -> - callCtx - .getMetaStore() - .deleteEntityAndCreateEntities(callCtx, entityToDrop, List.of(cleanupTask))) - .doesNotThrowAnyException(); + var retryResult = + metaStoreManager.dropEntityIfExists(callCtx, null, entityToDrop, Map.of(), true); + Assertions.assertThat(retryResult.isSuccess()).isTrue(); + Assertions.assertThat(retryResult.getCleanupTaskId()).isEqualTo(dropResult.getCleanupTaskId()); Assertions.assertThat( metaStoreManager .loadEntity(callCtx, 0L, cleanupTask.getId(), PolarisEntityType.TASK) diff --git a/runtime/service/src/main/java/org/apache/polaris/service/admin/PolarisAdminService.java b/runtime/service/src/main/java/org/apache/polaris/service/admin/PolarisAdminService.java index faf6fd91d34..f2e0e39d6db 100644 --- a/runtime/service/src/main/java/org/apache/polaris/service/admin/PolarisAdminService.java +++ b/runtime/service/src/main/java/org/apache/polaris/service/admin/PolarisAdminService.java @@ -932,7 +932,10 @@ public void deleteCatalog(String name) { // at least some handling of error if (!dropEntityResult.isSuccess()) { - if (dropEntityResult.failedBecauseNotEmpty()) { + if (dropEntityResult.alreadyExists()) { + throw new CommitConflictException( + "Concurrent cleanup task creation while dropping catalog '%s'", entity.getName()); + } else if (dropEntityResult.failedBecauseNotEmpty()) { throw new BadRequestException( "Catalog '%s' cannot be dropped, it is not empty", entity.getName()); } else { @@ -1380,7 +1383,10 @@ public void deletePrincipalRole(String name) { // at least some handling of error if (!dropEntityResult.isSuccess()) { - if (dropEntityResult.isEntityUnDroppable()) { + if (dropEntityResult.alreadyExists()) { + throw new CommitConflictException( + "Concurrent cleanup task creation while dropping principal role '%s'", name); + } else if (dropEntityResult.isEntityUnDroppable()) { throw new BadRequestException("Polaris service admin principal role cannot be dropped"); } else { throw new BadRequestException( @@ -1500,7 +1506,10 @@ public void deleteCatalogRole(String catalogName, String name) { // at least some handling of error if (!dropEntityResult.isSuccess()) { - if (dropEntityResult.isEntityUnDroppable()) { + if (dropEntityResult.alreadyExists()) { + throw new CommitConflictException( + "Concurrent cleanup task creation while dropping catalog role '%s'", name); + } else if (dropEntityResult.isEntityUnDroppable()) { throw new BadRequestException("Catalog admin role cannot be dropped"); } else { throw new BadRequestException( diff --git a/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/LocalIcebergCatalog.java b/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/LocalIcebergCatalog.java index 5b77cdeb196..a095a513e06 100644 --- a/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/LocalIcebergCatalog.java +++ b/runtime/service/src/main/java/org/apache/polaris/service/catalog/iceberg/LocalIcebergCatalog.java @@ -605,6 +605,10 @@ public boolean dropTable(TableIdentifier tableIdentifier, boolean purge) { "Table %s cannot be dropped: %s", tableIdentifier, dropEntityResult.getExtraInformation()); + case BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS: + throw new CommitConflictException( + "Concurrent cleanup task creation while dropping table %s", tableIdentifier); + default: throw new ServiceFailureException( "Failed to drop table %s, status=%s, extraInfo=%s", @@ -856,6 +860,10 @@ public boolean dropNamespace(Namespace namespace) throws NamespaceNotEmptyExcept dropEntityResult.getExtraInformation()); return false; + case BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS: + throw new CommitConflictException( + "Concurrent cleanup task creation while dropping namespace %s", namespace); + default: throw new ServiceFailureException( "Failed to drop namespace %s, status=%s, extraInfo=%s", @@ -1150,6 +1158,10 @@ public boolean dropView(TableIdentifier identifier) { throw new ForbiddenException( "View %s cannot be dropped: %s", identifier, dropEntityResult.getExtraInformation()); + case BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS: + throw new CommitConflictException( + "Concurrent cleanup task creation while dropping view %s", identifier); + default: throw new ServiceFailureException( "Failed to drop view %s, status=%s, extraInfo=%s", diff --git a/runtime/service/src/test/java/org/apache/polaris/service/admin/PolarisAdminServiceTest.java b/runtime/service/src/test/java/org/apache/polaris/service/admin/PolarisAdminServiceTest.java index 7d47b27c7a4..dd347937fff 100644 --- a/runtime/service/src/test/java/org/apache/polaris/service/admin/PolarisAdminServiceTest.java +++ b/runtime/service/src/test/java/org/apache/polaris/service/admin/PolarisAdminServiceTest.java @@ -22,6 +22,7 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doThrow; @@ -67,11 +68,14 @@ import org.apache.polaris.core.entity.PolarisEntityType; import org.apache.polaris.core.entity.PolarisPrivilege; import org.apache.polaris.core.entity.table.IcebergTableLikeEntity; +import org.apache.polaris.core.exceptions.CommitConflictException; import org.apache.polaris.core.identity.provider.ServiceIdentityProvider; import org.apache.polaris.core.persistence.PolarisMetaStoreManager; import org.apache.polaris.core.persistence.PolarisResolvedPathWrapper; +import org.apache.polaris.core.persistence.ResolvedPolarisEntity; import org.apache.polaris.core.persistence.dao.entity.BaseResult; import org.apache.polaris.core.persistence.dao.entity.CreateCatalogResult; +import org.apache.polaris.core.persistence.dao.entity.DropEntityResult; import org.apache.polaris.core.persistence.dao.entity.EntityResult; import org.apache.polaris.core.persistence.dao.entity.GenerateEntityIdResult; import org.apache.polaris.core.persistence.dao.entity.PrivilegeResult; @@ -183,6 +187,25 @@ void testCreateCatalogCleansUpInlineBearerSecretWhenCatalogAlreadyExists() { verify(userSecretsManager).deleteSecret(secretReference); } + @Test + void testDeleteCatalogWithCleanupTaskConflict() { + String catalogName = "test-catalog"; + PolarisEntity catalogEntity = createEntity(catalogName, PolarisEntityType.CATALOG); + when(resolutionManifest.getResolvedCatalogEntity()).thenReturn(CatalogEntity.of(catalogEntity)); + when(resolutionManifest.getResolvedTopLevelEntity(catalogName, PolarisEntityType.CATALOG)) + .thenReturn(resolvedPathWrapper); + when(resolvedPathWrapper.getResolvedLeafEntity()) + .thenReturn(new ResolvedPolarisEntity(catalogEntity, List.of(), List.of())); + when(realmConfig.getConfig(FeatureConfiguration.CLEANUP_ON_CATALOG_DROP)).thenReturn(true); + when(metaStoreManager.dropEntityIfExists(any(), any(), any(), any(), anyBoolean())) + .thenReturn(new DropEntityResult(BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS, null)); + + assertThatThrownBy(() -> adminService.deleteCatalog(catalogName)) + .isInstanceOf(CommitConflictException.class) + .hasMessageContaining("Concurrent cleanup task creation") + .hasMessageContaining(catalogName); + } + @Test void testCreateCatalogCleanupFailureDoesNotHideOriginalFailure() { SecretReference secretReference = diff --git a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java index 90b8310c5b6..49ce790a743 100644 --- a/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java +++ b/runtime/service/src/test/java/org/apache/polaris/service/catalog/iceberg/AbstractLocalIcebergCatalogTest.java @@ -2072,6 +2072,27 @@ public void testDropNamespaceWithUnexpectedError() { .hasMessageContaining("UNEXPECTED_ERROR_SIGNALED"); } + @Test + public void testDropNamespaceWithCleanupTaskConflict() { + catalog.createNamespace(NS); + + PolarisMetaStoreManager spiedManager = spy(metaStoreManager); + doReturn(new DropEntityResult(BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS, null)) + .when(spiedManager) + .dropEntityIfExists(any(), anyList(), any(), anyMap(), anyBoolean()); + + LocalIcebergCatalog spiedCatalog = newIcebergCatalog(CATALOG_NAME, spiedManager, fileIOFactory); + spiedCatalog.initialize( + CATALOG_NAME, + ImmutableMap.of( + CatalogProperties.FILE_IO_IMPL, "org.apache.iceberg.inmemory.InMemoryFileIO")); + + Assertions.assertThatThrownBy(() -> spiedCatalog.dropNamespace(NS)) + .isInstanceOf(CommitConflictException.class) + .hasMessageContaining("Concurrent cleanup task creation") + .hasMessageContaining("namespace"); + } + @Test public void testDropTableWithUnexpectedError() { catalog.createNamespace(NS); @@ -2097,6 +2118,28 @@ public void testDropTableWithUnexpectedError() { .hasMessageContaining("UNEXPECTED_ERROR_SIGNALED"); } + @Test + public void testDropTableWithCleanupTaskConflict() { + catalog.createNamespace(NS); + catalog.buildTable(TABLE, SCHEMA).create(); + + PolarisMetaStoreManager spiedManager = spy(metaStoreManager); + doReturn(new DropEntityResult(BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS, null)) + .when(spiedManager) + .dropEntityIfExists(any(), anyList(), any(), anyMap(), anyBoolean()); + + LocalIcebergCatalog spiedCatalog = newIcebergCatalog(CATALOG_NAME, spiedManager, fileIOFactory); + spiedCatalog.initialize( + CATALOG_NAME, + ImmutableMap.of( + CatalogProperties.FILE_IO_IMPL, "org.apache.iceberg.inmemory.InMemoryFileIO")); + + Assertions.assertThatThrownBy(() -> spiedCatalog.dropTable(TABLE, false)) + .isInstanceOf(CommitConflictException.class) + .hasMessageContaining("Concurrent cleanup task creation") + .hasMessageContaining("table"); + } + @Test public void testDropTableWithUndroppableEntity() { catalog.createNamespace(NS); @@ -2150,6 +2193,33 @@ public void testDropViewWithUnexpectedError() { .hasMessageContaining("UNEXPECTED_ERROR_SIGNALED"); } + @Test + public void testDropViewWithCleanupTaskConflict() { + catalog.createNamespace(NS); + catalog + .buildView(TABLE) + .withSchema(SCHEMA) + .withDefaultNamespace(NS) + .withQuery("spark", "SELECT * FROM ns.tbl") + .create(); + + PolarisMetaStoreManager spiedManager = spy(metaStoreManager); + doReturn(new DropEntityResult(BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS, null)) + .when(spiedManager) + .dropEntityIfExists(any(), anyList(), any(), anyMap(), anyBoolean()); + + LocalIcebergCatalog spiedCatalog = newIcebergCatalog(CATALOG_NAME, spiedManager, fileIOFactory); + spiedCatalog.initialize( + CATALOG_NAME, + ImmutableMap.of( + CatalogProperties.FILE_IO_IMPL, "org.apache.iceberg.inmemory.InMemoryFileIO")); + + Assertions.assertThatThrownBy(() -> spiedCatalog.dropView(TABLE)) + .isInstanceOf(CommitConflictException.class) + .hasMessageContaining("Concurrent cleanup task creation") + .hasMessageContaining("view"); + } + @Test public void testDropViewWithUndroppableEntity() { catalog.createNamespace(NS);