This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new d460830e5f [#12405] fix(core): prevent silent authorization updates on
closed catalogs (#12793)
d460830e5f is described below
commit d460830e5fce9179b4ccf71f11aebc81d58bfb8b
Author: Qi Yu <[email protected]>
AuthorDate: Mon Sep 7 19:14:06 2026 +0800
[#12405] fix(core): prevent silent authorization updates on closed catalogs
(#12793)
### What changes were proposed in this pull request?
- Validate the authorization-plugin lifecycle in the public
BaseCatalog.getAuthorizationPlugin() accessor.
- Derive whether a provider is configured from the existing catalog
property metadata and configuration instead of maintaining duplicate
state.
- Throw AuthorizationPluginException when a configured plugin is
unexpectedly unavailable instead of silently skipping the authorization
update.
- Use the leased catalog's canonical name when removing catalog
privileges, and reuse the caller's CatalogManager for metadata updates.
- Cover unavailable configured plugins, catalogs without authorization
providers, canonical catalog names, and multi-catalog updates.
### Why are the changes needed?
A catalog cache eviction can close the authorization plugin while a
catalog reference is still in use. Silently treating the resulting null
plugin as “authorization is disabled” can leave stale grants in the
external authorization system.
#12404 makes the authorization call paths lease-scoped. This patch
enforces the lifecycle invariant at the public plugin accessor so future
call sites cannot silently bypass it.
Fix: #12405
### Does this PR introduce _any_ user-facing change?
Yes. If an authorization provider is configured but its plugin is
unexpectedly unavailable, the operation now fails with
AuthorizationPluginException instead of silently succeeding without
updating the external authorization system.
There are no REST API or configuration-key changes.
### How was this patch tested?
```bash
./gradlew :core:spotlessApply
./gradlew :core:test \
--tests org.apache.gravitino.authorization.TestAuthorizationUtils \
--tests org.apache.gravitino.authorization.TestFutureGrantManager \
--tests org.apache.gravitino.connector.authorization.TestAuthorization \
-PskipITs -PskipDockerTests=true
./gradlew :core:javadoc
```
---
.../authorization/AuthorizationUtils.java | 25 ++++----
.../apache/gravitino/connector/BaseCatalog.java | 43 +++++++++----
.../authorization/TestAuthorizationUtils.java | 72 ++++++++++++++++++++++
.../connector/authorization/TestAuthorization.java | 65 +++++++++++++++++++
4 files changed, 181 insertions(+), 24 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
b/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
index aa4875f3ff..43f99539ad 100644
---
a/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
+++
b/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
@@ -227,7 +227,8 @@ public class AuthorizationUtils {
public static void callAuthorizationPluginForMetadataObject(
String metalake, MetadataObject metadataObject,
Consumer<AuthorizationPlugin> consumer) {
CatalogManager catalogManager =
GravitinoEnv.getInstance().catalogManager();
- List<NameIdentifier> catalogIdents = getMetadataObjectCatalogs(metalake,
metadataObject);
+ List<NameIdentifier> catalogIdents =
+ getMetadataObjectCatalogs(catalogManager, metalake, metadataObject);
for (NameIdentifier catalogIdent : catalogIdents) {
callAuthorizationPluginImpl(consumer, catalogManager, catalogIdent);
}
@@ -344,16 +345,19 @@ public class AuthorizationUtils {
}
}
+ /**
+ * Removes catalog privileges using the live catalog's name while its
operation lease is held.
+ *
+ * @param catalogIdent the identifier used to load the catalog
+ * @param locations the catalog storage locations
+ */
public static void removeCatalogPrivileges(NameIdentifier catalogIdent,
List<String> locations) {
- // If we enable authorization, we should remove the privileges about the
entity in the
- // authorization plugin.
- MetadataObject metadataObject =
- MetadataObjects.of(null, catalogIdent.name(),
MetadataObject.Type.CATALOG);
- MetadataObjectChange removeObject =
MetadataObjectChange.remove(metadataObject, locations);
-
callAuthorizationPluginImpl(
- authorizationPlugin -> {
- authorizationPlugin.onMetadataUpdated(removeObject);
+ (authorizationPlugin, catalogName) -> {
+ MetadataObject metadataObject =
+ MetadataObjects.of(null, catalogName,
MetadataObject.Type.CATALOG);
+ authorizationPlugin.onMetadataUpdated(
+ MetadataObjectChange.remove(metadataObject, locations));
},
GravitinoEnv.getInstance().catalogManager(),
catalogIdent);
@@ -485,8 +489,7 @@ public class AuthorizationUtils {
}
private static List<NameIdentifier> getMetadataObjectCatalogs(
- String metalake, MetadataObject metadataObject) {
- CatalogManager catalogManager =
GravitinoEnv.getInstance().catalogManager();
+ CatalogManager catalogManager, String metalake, MetadataObject
metadataObject) {
if (needApplyAuthorizationPluginAllCatalogs(metadataObject.type())) {
return
Arrays.asList(catalogManager.listCatalogs(Namespace.of(metalake)));
}
diff --git a/core/src/main/java/org/apache/gravitino/connector/BaseCatalog.java
b/core/src/main/java/org/apache/gravitino/connector/BaseCatalog.java
index dc077f2e56..5b0b664045 100644
--- a/core/src/main/java/org/apache/gravitino/connector/BaseCatalog.java
+++ b/core/src/main/java/org/apache/gravitino/connector/BaseCatalog.java
@@ -28,6 +28,7 @@ import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import javax.annotation.Nullable;
import org.apache.commons.lang3.StringUtils;
import org.apache.gravitino.Audit;
import org.apache.gravitino.Catalog;
@@ -46,6 +47,7 @@ import org.apache.gravitino.credential.CredentialConstants;
import org.apache.gravitino.credential.GCSTokenCredential;
import org.apache.gravitino.credential.OSSSecretKeyCredential;
import org.apache.gravitino.credential.S3SecretKeyCredential;
+import org.apache.gravitino.exceptions.AuthorizationPluginException;
import org.apache.gravitino.exceptions.CatalogNotInUseException;
import org.apache.gravitino.exceptions.MetalakeNotInUseException;
import org.apache.gravitino.meta.CatalogEntity;
@@ -268,20 +270,26 @@ public abstract class BaseCatalog<T extends BaseCatalog>
return Boolean.parseBoolean(catalogInUseStr);
}
- private boolean isInvokedBy(String methodName) {
- return StackWalker.getInstance()
- .walk(frames -> frames.anyMatch(frame ->
frame.getMethodName().equals(methodName)));
- }
-
+ /**
+ * Returns the authorization plugin configured for this catalog.
+ *
+ * <p>A configured provider without a plugin means the catalog's
authorization lifecycle is
+ * incomplete. Silently treating that state as authorization being disabled
could leave stale
+ * grants in the external authorization system.
+ *
+ * @return the authorization plugin, or null if no authorization provider is
configured
+ * @throws AuthorizationPluginException if a configured authorization plugin
is unavailable
+ */
+ @Nullable
public AuthorizationPlugin getAuthorizationPlugin() {
- if (authorizationPlugin == null) {
- synchronized (this) {
- if (authorizationPlugin == null) {
- return null;
- }
- }
+ AuthorizationPlugin plugin = authorizationPlugin;
+ if (plugin == null && isAuthorizationProviderConfigured()) {
+ throw new AuthorizationPluginException(
+ "The authorization plugin of catalog %s is unavailable although an
authorization "
+ + "provider is configured",
+ name());
}
- return authorizationPlugin;
+ return plugin;
}
/**
@@ -300,7 +308,6 @@ public abstract class BaseCatalog<T extends BaseCatalog>
LOG.info("Authorization provider is not set!");
return;
}
-
// use try-with-resources to auto-close authorization object if exit
with exception
try (BaseAuthorization<?> authorization =
BaseAuthorization.createAuthorization(classLoader,
authorizationProvider)) {
@@ -388,6 +395,16 @@ public abstract class BaseCatalog<T extends BaseCatalog>
return catalogCredentialManager;
}
+ private boolean isInvokedBy(String methodName) {
+ return StackWalker.getInstance()
+ .walk(frames -> frames.anyMatch(frame ->
frame.getMethodName().equals(methodName)));
+ }
+
+ private boolean isAuthorizationProviderConfigured() {
+ return conf != null
+ && catalogPropertiesMetadata().getOrDefault(conf,
AUTHORIZATION_PROVIDER) != null;
+ }
+
private CatalogOperations createOps(Map<String, String> conf) {
String customCatalogOperationClass = conf.get(CATALOG_OPERATION_IMPL);
return Optional.ofNullable(customCatalogOperationClass)
diff --git
a/core/src/test/java/org/apache/gravitino/authorization/TestAuthorizationUtils.java
b/core/src/test/java/org/apache/gravitino/authorization/TestAuthorizationUtils.java
index 102d4854a7..bef17095d0 100644
---
a/core/src/test/java/org/apache/gravitino/authorization/TestAuthorizationUtils.java
+++
b/core/src/test/java/org/apache/gravitino/authorization/TestAuthorizationUtils.java
@@ -29,6 +29,7 @@ import org.apache.gravitino.Catalog;
import org.apache.gravitino.Entity;
import org.apache.gravitino.GravitinoEnv;
import org.apache.gravitino.MetadataObject;
+import org.apache.gravitino.MetadataObjects;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
import org.apache.gravitino.Schema;
@@ -44,6 +45,7 @@ import
org.apache.gravitino.exceptions.IllegalNamespaceException;
import org.apache.gravitino.meta.AuditInfo;
import org.apache.gravitino.meta.RoleEntity;
import org.apache.gravitino.rel.Table;
+import org.apache.gravitino.utils.ThrowableFunction;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
@@ -458,4 +460,74 @@ class TestAuthorizationUtils {
Assertions.assertEquals("catalog.schema.table",
removeChange.metadataObject().fullName());
Assertions.assertEquals(locations, removeChange.getLocations());
}
+
+ @Test
+ void testRemoveCatalogPrivilegesUsesLeasedCatalogName() {
+ NameIdentifier requestedIdent = NameIdentifier.of("metalake",
"requested_catalog");
+ List<String> locations = List.of("/warehouse/catalog");
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ BaseCatalog<?> catalog = Mockito.mock(BaseCatalog.class);
+ AuthorizationPlugin plugin = Mockito.mock(AuthorizationPlugin.class);
+ CatalogTestUtils.mockDoWithCatalog(catalogManager, catalog);
+ Mockito.when(catalog.name()).thenReturn("canonical_catalog");
+ Mockito.when(catalog.getAuthorizationPlugin()).thenReturn(plugin);
+
+ GravitinoEnv env = Mockito.mock(GravitinoEnv.class);
+ Mockito.when(env.catalogManager()).thenReturn(catalogManager);
+ try (MockedStatic<GravitinoEnv> envStatic =
Mockito.mockStatic(GravitinoEnv.class)) {
+ envStatic.when(GravitinoEnv::getInstance).thenReturn(env);
+ AuthorizationUtils.removeCatalogPrivileges(requestedIdent, locations);
+ }
+
+ Mockito.verify(catalogManager).doWithCatalog(Mockito.eq(requestedIdent),
Mockito.any());
+ ArgumentCaptor<MetadataObjectChange[]> changes =
+ ArgumentCaptor.forClass(MetadataObjectChange[].class);
+ Mockito.verify(plugin).onMetadataUpdated(changes.capture());
+ Assertions.assertEquals(1, changes.getValue().length);
+ MetadataObjectChange.RemoveMetadataObject removal =
+ Assertions.assertInstanceOf(
+ MetadataObjectChange.RemoveMetadataObject.class,
changes.getValue()[0]);
+ Assertions.assertEquals(MetadataObject.Type.CATALOG,
removal.metadataObject().type());
+ Assertions.assertEquals("canonical_catalog",
removal.metadataObject().fullName());
+ Assertions.assertEquals(locations, removal.getLocations());
+ }
+
+ @Test
+ void testMetalakeUpdateVisitsEachCatalogWithAuthorization() {
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ NameIdentifier first = NameIdentifier.of(metalake, "first");
+ NameIdentifier second = NameIdentifier.of(metalake, "second");
+ NameIdentifier withoutAuthorization = NameIdentifier.of(metalake,
"without_auth");
+ Mockito.when(catalogManager.listCatalogs(Namespace.of(metalake)))
+ .thenReturn(new NameIdentifier[] {first, second,
withoutAuthorization});
+ BaseCatalog<?> firstCatalog = Mockito.mock(BaseCatalog.class);
+ BaseCatalog<?> secondCatalog = Mockito.mock(BaseCatalog.class);
+ BaseCatalog<?> plainCatalog = Mockito.mock(BaseCatalog.class);
+ AuthorizationPlugin firstPlugin = Mockito.mock(AuthorizationPlugin.class);
+ AuthorizationPlugin secondPlugin = Mockito.mock(AuthorizationPlugin.class);
+
Mockito.when(firstCatalog.getAuthorizationPlugin()).thenReturn(firstPlugin);
+
Mockito.when(secondCatalog.getAuthorizationPlugin()).thenReturn(secondPlugin);
+ CatalogTestUtils.mockDoWithCatalog(catalogManager, plainCatalog);
+ Mockito.doAnswer(
+ invocation -> {
+ ThrowableFunction<BaseCatalog, Object> operation =
invocation.getArgument(1);
+ NameIdentifier ident = invocation.getArgument(0);
+ return operation.apply(first.equals(ident) ? firstCatalog :
secondCatalog);
+ })
+ .when(catalogManager)
+ .doWithCatalog(
+ Mockito.argThat(ident -> first.equals(ident) ||
second.equals(ident)), Mockito.any());
+
+ List<AuthorizationPlugin> visited = Lists.newArrayList();
+ GravitinoEnv env = Mockito.mock(GravitinoEnv.class);
+ Mockito.when(env.catalogManager()).thenReturn(catalogManager);
+ try (MockedStatic<GravitinoEnv> envStatic =
Mockito.mockStatic(GravitinoEnv.class)) {
+ envStatic.when(GravitinoEnv::getInstance).thenReturn(env);
+ AuthorizationUtils.callAuthorizationPluginForMetadataObject(
+ metalake, MetadataObjects.of(null, metalake,
MetadataObject.Type.METALAKE), visited::add);
+ }
+
+ Assertions.assertEquals(List.of(firstPlugin, secondPlugin), visited);
+
Mockito.verify(catalogManager).doWithCatalog(Mockito.eq(withoutAuthorization),
Mockito.any());
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/connector/authorization/TestAuthorization.java
b/core/src/test/java/org/apache/gravitino/connector/authorization/TestAuthorization.java
index 35b13f7ec7..3f0f62524e 100644
---
a/core/src/test/java/org/apache/gravitino/connector/authorization/TestAuthorization.java
+++
b/core/src/test/java/org/apache/gravitino/connector/authorization/TestAuthorization.java
@@ -27,6 +27,7 @@ import org.apache.gravitino.TestCatalog;
import
org.apache.gravitino.connector.authorization.ranger.TestRangerAuthorization;
import
org.apache.gravitino.connector.authorization.ranger.TestRangerAuthorizationHDFSPlugin;
import
org.apache.gravitino.connector.authorization.ranger.TestRangerAuthorizationHadoopSQLPlugin;
+import org.apache.gravitino.exceptions.AuthorizationPluginException;
import org.apache.gravitino.meta.AuditInfo;
import org.apache.gravitino.meta.CatalogEntity;
import org.apache.gravitino.utils.IsolatedClassLoader;
@@ -103,6 +104,70 @@ public class TestAuthorization {
Assertions.assertTrue(testRangerAuthHadoopSQLPlugin.callOnCreateRole1);
}
+ @Test
+ public void testConfiguredAuthorizationPluginUnavailableAfterClose() throws
Exception {
+ AuditInfo auditInfo =
+
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build();
+ CatalogEntity entity =
+ CatalogEntity.builder()
+ .withId(3L)
+ .withName("catalog-test3")
+ .withNamespace(Namespace.of("default"))
+ .withType(Catalog.Type.RELATIONAL)
+ .withProvider("test")
+ .withAuditInfo(auditInfo)
+ .build();
+
+ TestCatalog catalog =
+ new TestCatalog()
+ .withCatalogConf(
+ ImmutableMap.of(
+ Catalog.AUTHORIZATION_PROVIDER,
+ "test-ranger",
+ "authorization.ranger.service.type",
+ "HadoopSQL"))
+ .withCatalogEntity(entity);
+ try (IsolatedClassLoader isolatedClassLoader =
+ new IsolatedClassLoader(
+ Collections.emptyList(), Collections.emptyList(),
Collections.emptyList())) {
+ catalog.initAuthorizationPluginInstance(isolatedClassLoader,
METALAKE_ID);
+ Assertions.assertNotNull(catalog.getAuthorizationPlugin());
+
+ catalog.close();
+ AuthorizationPluginException exception =
+ Assertions.assertThrows(
+ AuthorizationPluginException.class,
catalog::getAuthorizationPlugin);
+ Assertions.assertTrue(
+ exception.getMessage().contains("catalog-test3"),
exception.getMessage());
+ }
+ }
+
+ @Test
+ public void testAuthorizationProviderNotConfigured() throws Exception {
+ AuditInfo auditInfo =
+
AuditInfo.builder().withCreator("test").withCreateTime(Instant.now()).build();
+ CatalogEntity entity =
+ CatalogEntity.builder()
+ .withId(4L)
+ .withName("catalog-test4")
+ .withNamespace(Namespace.of("default"))
+ .withType(Catalog.Type.RELATIONAL)
+ .withProvider("test")
+ .withAuditInfo(auditInfo)
+ .build();
+
+ TestCatalog catalog =
+ new
TestCatalog().withCatalogConf(ImmutableMap.of()).withCatalogEntity(entity);
+ try (IsolatedClassLoader isolatedClassLoader =
+ new IsolatedClassLoader(
+ Collections.emptyList(), Collections.emptyList(),
Collections.emptyList());
+ catalog) {
+ catalog.initAuthorizationPluginInstance(isolatedClassLoader,
METALAKE_ID);
+
+ Assertions.assertNull(catalog.getAuthorizationPlugin());
+ }
+ }
+
@Test
public void testRangerHDFSAuthorization() {
AuthorizationPlugin rangerHDFSAuthPlugin =
filesetCatalog.getAuthorizationPlugin();