diff --git a/runtime/service/src/main/java/org/apache/polaris/service/catalog/policy/PolicyCatalog.java b/runtime/service/src/main/java/org/apache/polaris/service/catalog/policy/PolicyCatalog.java index 9bd94716d77..557a1a68fe8 100644 --- a/runtime/service/src/main/java/org/apache/polaris/service/catalog/policy/PolicyCatalog.java +++ b/runtime/service/src/main/java/org/apache/polaris/service/catalog/policy/PolicyCatalog.java @@ -45,6 +45,7 @@ import org.apache.polaris.core.entity.PolarisEntityCore; import org.apache.polaris.core.entity.PolarisEntitySubType; import org.apache.polaris.core.entity.PolarisEntityType; +import org.apache.polaris.core.exceptions.CommitConflictException; import org.apache.polaris.core.persistence.PolarisMetaStoreManager; import org.apache.polaris.core.persistence.PolarisResolvedPathWrapper; import org.apache.polaris.core.persistence.PolicyMappingAlreadyExistsException; @@ -249,12 +250,10 @@ public Policy updatePolicy( newPolicyEntity) .getEntity()) .map(PolicyEntity::of) - .orElse(null); - - if (newPolicyEntity == null) { - throw new IllegalStateException( - String.format("Failed to update policy %s", policyIdentifier)); - } + .orElseThrow( + () -> + new CommitConflictException( + "Concurrent modification on policy '%s'; retry later", policyIdentifier)); return constructPolicy(newPolicyEntity); } diff --git a/runtime/service/src/test/java/org/apache/polaris/service/catalog/policy/AbstractPolicyCatalogTest.java b/runtime/service/src/test/java/org/apache/polaris/service/catalog/policy/AbstractPolicyCatalogTest.java index 3a1152bd5e2..91013d2a826 100644 --- a/runtime/service/src/test/java/org/apache/polaris/service/catalog/policy/AbstractPolicyCatalogTest.java +++ b/runtime/service/src/test/java/org/apache/polaris/service/catalog/policy/AbstractPolicyCatalogTest.java @@ -58,9 +58,12 @@ import org.apache.polaris.core.entity.CatalogEntity; import org.apache.polaris.core.entity.PolarisEntity; import org.apache.polaris.core.entity.PrincipalEntity; +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.PolicyMappingAlreadyExistsException; +import org.apache.polaris.core.persistence.dao.entity.BaseResult; +import org.apache.polaris.core.persistence.dao.entity.EntityResult; import org.apache.polaris.core.persistence.resolver.ResolutionManifestFactory; import org.apache.polaris.core.persistence.resolver.ResolverFactory; import org.apache.polaris.core.policy.PredefinedPolicyTypes; @@ -420,6 +423,31 @@ public void testUpdatePolicyWithWrongVersion() { .isInstanceOf(PolicyVersionMismatchException.class); } + @Test + public void testUpdatePolicyLosingConcurrentUpdateIsRetryableConflict() { + icebergCatalog.createNamespace(NS); + policyCatalog.createPolicy( + POLICY1, PredefinedPolicyTypes.DATA_COMPACTION.getName(), "test", "{\"enable\": false}"); + + // Simulate another writer winning the compare-and-swap on the policy entity. + PolarisMetaStoreManager concurrentlyModified = Mockito.spy(metaStoreManager); + Mockito.doReturn( + new EntityResult( + BaseResult.ReturnStatus.TARGET_ENTITY_CONCURRENTLY_MODIFIED, "simulated")) + .when(concurrentlyModified) + .updateEntityPropertiesIfNotChanged(Mockito.any(), Mockito.any(), Mockito.any()); + + PolicyCatalog catalog = + new PolicyCatalog( + concurrentlyModified, + polarisContext, + new PolarisPassthroughResolutionView( + resolutionManifestFactory, authenticatedRoot, CATALOG_NAME)); + + assertThatThrownBy(() -> catalog.updatePolicy(POLICY1, "updated", "{\"enable\": true}", 0)) + .isInstanceOf(CommitConflictException.class); + } + @Test public void testUpdatePolicyWithInvalidContent() { icebergCatalog.createNamespace(NS); diff --git a/spec/generated/bundled-polaris-catalog-service.yaml b/spec/generated/bundled-polaris-catalog-service.yaml index 27302100176..719dbb911c8 100644 --- a/spec/generated/bundled-polaris-catalog-service.yaml +++ b/spec/generated/bundled-polaris-catalog-service.yaml @@ -1766,7 +1766,7 @@ paths: PolicyToUpdateDoesNotExist: $ref: '#/components/examples/NoSuchPolicyError' '409': - description: The policy version doesn't match the current-policy-version; retry after fetching latest version + description: The policy version doesn't match the current-policy-version, or the policy was concurrently modified; retry after fetching latest version content: application/json: schema: diff --git a/spec/polaris-catalog-apis/policy-apis.yaml b/spec/polaris-catalog-apis/policy-apis.yaml index 9aa15620cf2..23a5ea87d85 100644 --- a/spec/polaris-catalog-apis/policy-apis.yaml +++ b/spec/polaris-catalog-apis/policy-apis.yaml @@ -195,7 +195,7 @@ paths: PolicyToUpdateDoesNotExist: $ref: '#/components/examples/NoSuchPolicyError' 409: - description: "The policy version doesn't match the current-policy-version; retry after fetching latest version" + description: "The policy version doesn't match the current-policy-version, or the policy was concurrently modified; retry after fetching latest version" content: application/json: schema: