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..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 @@ -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..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 @@ -352,6 +352,49 @@ 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) { + PolarisBaseEntity existingEntity = + lookupEntity( + connection, + entityToCreate.getCatalogId(), + entityToCreate.getId(), + entityToCreate.getTypeCode()); + if (existingEntity != null) { + continue; + } + 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) { @@ -421,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 @@ -455,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 4ccee4483c9..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 @@ -48,6 +48,16 @@ public abstract class AtomicMetastoreManagerWithJdbcBasePersistenceImplTest extends BasePolarisMetaStoreManagerTest { + @Test + 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 0570758f864..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 @@ -177,6 +177,14 @@ 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. @@ -1092,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); } @@ -1182,10 +1206,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 @@ -1201,22 +1221,33 @@ 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()); 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()); + 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()); } + // 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..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 @@ -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..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 @@ -285,6 +285,32 @@ 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; + } + continue; + } + this.writeEntityInCurrentTxn(callCtx, entityToCreate, true, null); + } + }); + } + /** {@inheritDoc} */ @Override public void deleteFromGrantRecords( 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/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java b/polaris-core/src/test/java/org/apache/polaris/core/persistence/PolarisTreeMapAtomicOperationMetaStoreManagerTest.java index d5d526527a1..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 @@ -31,6 +31,16 @@ public class PolarisTreeMapAtomicOperationMetaStoreManagerTest extends BasePolarisMetaStoreManagerTest { + @Test + void testCleanupTaskCreationFailureRollsBackEntityDrop() { + assertCleanupTaskCreationFailureRollsBackEntityDrop(); + } + + @Test + void testCleanupTaskCreationRetryIsIdempotent() { + assertCleanupTaskCreationRetryIsIdempotent(); + } + @Override public PolarisTestMetaStoreManager createPolarisTestMetaStoreManager() { PolarisDiagnostics diagServices = new PolarisDefaultDiagServiceImpl(); 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 d497714ba71..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 @@ -272,6 +272,81 @@ 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.assertThat( + metaStoreManager.dropEntityIfExists(callCtx, null, entityToDrop, Map.of(), true)) + .extracting(BaseResult::getReturnStatus) + .isEqualTo(BaseResult.ReturnStatus.ENTITY_ALREADY_EXISTS); + + 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); + } + + 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(); + 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) + .getEntity()) + .isEqualTo(cleanupTask); + } + /** Test that granting/revoking privileges works well */ @Test protected void testPrivileges() { 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);