This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 415bb34683 [Cherry-pick to branch-1.3] [#12405] fix(core): prevent
silent authorization updates on closed catalogs (#12793) (#12963)
415bb34683 is described below
commit 415bb346837ac26d62a08e9c66f04bd17328c74b
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Tue Sep 8 08:49:51 2026 +0800
[Cherry-pick to branch-1.3] [#12405] fix(core): prevent silent
authorization updates on closed catalogs (#12793) (#12963)
**Cherry-pick Information:**
- Original commit: d460830e5fce9179b4ccf71f11aebc81d58bfb8b
- Target branch: `branch-1.3`
- Status: ✅ Clean cherry-pick (no conflicts)
---------
Co-authored-by: Qi Yu <[email protected]>
---
.../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 83d1434225..11c380d669 100644
---
a/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
+++
b/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
@@ -224,7 +224,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);
}
@@ -341,16 +342,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);
@@ -482,8 +486,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 4ee0b9c794..04fc8a58f6 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;
@@ -45,6 +46,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;
@@ -254,20 +256,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;
}
public void initAuthorizationPluginInstance(IsolatedClassLoader classLoader)
{
@@ -280,7 +288,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)) {
@@ -364,6 +371,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 4ee37b4dde..a63425bbe9 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
@@ -26,6 +26,7 @@ import org.apache.gravitino.Namespace;
import org.apache.gravitino.TestCatalog;
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;
@@ -99,6 +100,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);
+ 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);
+
+ Assertions.assertNull(catalog.getAuthorizationPlugin());
+ }
+ }
+
@Test
public void testRangerHDFSAuthorization() {
AuthorizationPlugin rangerHDFSAuthPlugin =
filesetCatalog.getAuthorizationPlugin();