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 f6114c7bea [Cherry-pick to branch-1.3] [#12871] fix(authz): Invalidate
function grants after drop (#12873) (#12897)
f6114c7bea is described below
commit f6114c7bea59fec51929b8441c96a46aaafa6dee
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Fri Sep 4 20:58:27 2026 +0800
[Cherry-pick to branch-1.3] [#12871] fix(authz): Invalidate function grants
after drop (#12873) (#12897)
**Cherry-pick Information:**
- Original commit: 37f9d23c85dafa2f71c0766a4deb59957bf6645f
- Target branch: `branch-1.3`
- Status: ✅ Conflicts resolved manually
The resolution preserves the branch-1.3 service-level entity change-log
transaction pattern.
---------
Co-authored-by: Qi Yu <[email protected]>
---
.../java/org/apache/gravitino/GravitinoEnv.java | 3 +-
.../authorization/AuthorizationUtils.java | 13 +-
.../gravitino/catalog/CapabilityHelpers.java | 14 ++
.../catalog/FunctionNormalizeDispatcher.java | 3 +-
.../gravitino/hook/FunctionHookDispatcher.java | 38 +++-
.../relational/service/FunctionMetaService.java | 16 ++
.../authorization/TestAuthorizationUtils.java | 25 +++
.../gravitino/hook/TestFunctionHookDispatcher.java | 226 +++++++++++----------
.../service/TestEntityChangeLogService.java | 58 ++++++
.../jcasbin/TestJcasbinChangePoller.java | 14 ++
10 files changed, 287 insertions(+), 123 deletions(-)
diff --git a/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
b/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
index 4969be4a37..b60914afcf 100644
--- a/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
+++ b/core/src/main/java/org/apache/gravitino/GravitinoEnv.java
@@ -846,7 +846,8 @@ public class GravitinoEnv {
this.internalFunctionDispatcher = functionNormalizeDispatcher;
FunctionEventDispatcher functionEventDispatcher =
new FunctionEventDispatcher(eventBus, functionNormalizeDispatcher);
- this.functionDispatcher = new
FunctionHookDispatcher(functionEventDispatcher);
+ this.functionDispatcher =
+ new FunctionHookDispatcher(functionEventDispatcher,
this::ownerDispatcher, catalogManager);
// View operation chain: ViewHookDispatcher -> ViewEventDispatcher ->
ViewNormalizeDispatcher
// -> ViewOperationDispatcher.
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 280f4e7ffe..e6e8b95c8e 100644
---
a/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
+++
b/core/src/main/java/org/apache/gravitino/authorization/AuthorizationUtils.java
@@ -328,6 +328,7 @@ public class AuthorizationUtils {
// If we enable authorization, we should remove the privileges about the
entity in the
// authorization plugin.
if (GravitinoEnv.getInstance().internalAccessControlDispatcher() != null) {
+ notifyEntityNameIdMappingChange(ident, type);
MetadataObject metadataObject =
NameIdentifierUtil.toMetadataObject(ident, type);
String metalake =
type == Entity.EntityType.METALAKE ? ident.name() :
ident.namespace().level(0);
@@ -387,8 +388,16 @@ public class AuthorizationUtils {
}
}
- private static void notifyEntityNameIdMappingChange(
- NameIdentifier ident, Entity.EntityType type) {
+ /**
+ * Notifies the built-in authorizer that an entity name may now resolve to a
different ID.
+ *
+ * <p>This does not push a metadata change to the catalog authorization
plugin. Use it when only
+ * Gravitino's local authorization caches support the entity type.
+ *
+ * @param ident the entity identifier whose mapping changed
+ * @param type the entity type
+ */
+ public static void notifyEntityNameIdMappingChange(NameIdentifier ident,
Entity.EntityType type) {
GravitinoAuthorizer gravitinoAuthorizer =
GravitinoEnv.getInstance().gravitinoAuthorizer();
if (gravitinoAuthorizer == null) {
return;
diff --git
a/core/src/main/java/org/apache/gravitino/catalog/CapabilityHelpers.java
b/core/src/main/java/org/apache/gravitino/catalog/CapabilityHelpers.java
index 5825a97142..11a162274f 100644
--- a/core/src/main/java/org/apache/gravitino/catalog/CapabilityHelpers.java
+++ b/core/src/main/java/org/apache/gravitino/catalog/CapabilityHelpers.java
@@ -148,6 +148,20 @@ public class CapabilityHelpers {
return NameIdentifier.of(namespace, name);
}
+ /**
+ * Loads the catalog capability for {@code ident} and applies its
case-sensitivity rules.
+ *
+ * @param ident the identifier to normalize
+ * @param scope the identifier's capability scope
+ * @param catalogManager the catalog manager used to load the capability
+ * @return the case-normalized identifier
+ */
+ public static NameIdentifier applyCaseSensitive(
+ NameIdentifier ident, Capability.Scope scope, CatalogManager
catalogManager) {
+ Capability capability = getCapability(ident, catalogManager);
+ return applyCaseSensitive(ident, scope, capability);
+ }
+
public static Namespace applyCaseSensitive(
Namespace namespace, Capability.Scope identScope, Capability
capabilities) {
String metalake = namespace.level(0);
diff --git
a/core/src/main/java/org/apache/gravitino/catalog/FunctionNormalizeDispatcher.java
b/core/src/main/java/org/apache/gravitino/catalog/FunctionNormalizeDispatcher.java
index 810ca4ed5e..afd716218a 100644
---
a/core/src/main/java/org/apache/gravitino/catalog/FunctionNormalizeDispatcher.java
+++
b/core/src/main/java/org/apache/gravitino/catalog/FunctionNormalizeDispatcher.java
@@ -98,8 +98,7 @@ public class FunctionNormalizeDispatcher implements
FunctionDispatcher {
}
private NameIdentifier normalizeCaseSensitive(NameIdentifier functionIdent) {
- Capability capabilities = getCapability(functionIdent, catalogManager);
- return applyCaseSensitive(functionIdent, Capability.Scope.FUNCTION,
capabilities);
+ return applyCaseSensitive(functionIdent, Capability.Scope.FUNCTION,
catalogManager);
}
private NameIdentifier[] normalizeCaseSensitive(NameIdentifier[]
functionIdents) {
diff --git
a/core/src/main/java/org/apache/gravitino/hook/FunctionHookDispatcher.java
b/core/src/main/java/org/apache/gravitino/hook/FunctionHookDispatcher.java
index 20cfe2ce9c..33c70bda57 100644
--- a/core/src/main/java/org/apache/gravitino/hook/FunctionHookDispatcher.java
+++ b/core/src/main/java/org/apache/gravitino/hook/FunctionHookDispatcher.java
@@ -19,13 +19,15 @@
package org.apache.gravitino.hook;
+import java.util.function.Supplier;
import org.apache.gravitino.Entity;
-import org.apache.gravitino.GravitinoEnv;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
+import org.apache.gravitino.authorization.AuthorizationUtils;
import org.apache.gravitino.authorization.Owner;
import org.apache.gravitino.authorization.OwnerDispatcher;
import org.apache.gravitino.catalog.CapabilityHelpers;
+import org.apache.gravitino.catalog.CatalogManager;
import org.apache.gravitino.catalog.FunctionDispatcher;
import org.apache.gravitino.connector.capability.Capability;
import org.apache.gravitino.exceptions.FunctionAlreadyExistsException;
@@ -45,9 +47,24 @@ import org.apache.gravitino.utils.PrincipalUtils;
*/
public class FunctionHookDispatcher implements FunctionDispatcher {
private final FunctionDispatcher dispatcher;
+ private final Supplier<OwnerDispatcher> ownerDispatcher;
+ private final CatalogManager catalogManager;
- public FunctionHookDispatcher(FunctionDispatcher dispatcher) {
+ /**
+ * Creates a function hook dispatcher.
+ *
+ * @param dispatcher the underlying function dispatcher
+ * @param ownerDispatcher supplies the owner dispatcher, or {@code null}
when authorization is
+ * disabled
+ * @param catalogManager the catalog manager used to apply catalog
capabilities
+ */
+ public FunctionHookDispatcher(
+ FunctionDispatcher dispatcher,
+ Supplier<OwnerDispatcher> ownerDispatcher,
+ CatalogManager catalogManager) {
this.dispatcher = dispatcher;
+ this.ownerDispatcher = ownerDispatcher;
+ this.catalogManager = catalogManager;
}
@Override
@@ -82,15 +99,14 @@ public class FunctionHookDispatcher implements
FunctionDispatcher {
dispatcher.registerFunction(ident, comment, functionType,
deterministic, definitions);
// Set the creator as the owner of the function.
- OwnerDispatcher ownerManager =
GravitinoEnv.getInstance().internalOwnerDispatcher();
+ OwnerDispatcher ownerManager = ownerDispatcher.get();
if (ownerManager != null) {
// The inner NormalizeDispatcher case-folds the function name (and its
schema namespace)
// based on catalog capabilities, so the entity is stored under the
normalized identifier.
// Apply the same normalization here so the owner is attached to the
same identifier the
// manager sees.
NameIdentifier normalizedIdent =
- CapabilityHelpers.applyCapabilities(
- ident, Capability.Scope.FUNCTION,
GravitinoEnv.getInstance().catalogManager());
+ CapabilityHelpers.applyCapabilities(ident,
Capability.Scope.FUNCTION, catalogManager);
ownerManager.setOwner(
normalizedIdent.namespace().level(0),
NameIdentifierUtil.toMetadataObject(normalizedIdent,
Entity.EntityType.FUNCTION),
@@ -108,6 +124,16 @@ public class FunctionHookDispatcher implements
FunctionDispatcher {
@Override
public boolean dropFunction(NameIdentifier ident) {
- return dispatcher.dropFunction(ident);
+ boolean dropped = dispatcher.dropFunction(ident);
+ if (dropped) {
+ NameIdentifier normalizedIdent =
+ CapabilityHelpers.applyCaseSensitive(ident,
Capability.Scope.FUNCTION, catalogManager);
+ // Function privileges are managed by Gravitino. Catalog authorization
plugins such as
+ // Ranger HadoopSQL do not support FUNCTION metadata objects, so only
invalidate the built-in
+ // authorizer's name-to-ID mapping here.
+ AuthorizationUtils.notifyEntityNameIdMappingChange(
+ normalizedIdent, Entity.EntityType.FUNCTION);
+ }
+ return dropped;
}
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
index 03e6b3b8d9..32445cde21 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
@@ -41,6 +41,7 @@ import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.meta.FunctionEntity;
import org.apache.gravitino.meta.NamespacedEntityId;
import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
import org.apache.gravitino.storage.relational.mapper.FunctionMetaMapper;
import
org.apache.gravitino.storage.relational.mapper.FunctionVersionMetaMapper;
import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
@@ -49,6 +50,7 @@ import
org.apache.gravitino.storage.relational.mapper.SecurableObjectMapper;
import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.po.FunctionMaxVersionPO;
import org.apache.gravitino.storage.relational.po.FunctionPO;
+import org.apache.gravitino.storage.relational.po.cache.OperateType;
import org.apache.gravitino.storage.relational.utils.ExceptionUtils;
import org.apache.gravitino.storage.relational.utils.SessionUtils;
import org.apache.gravitino.utils.NameIdentifierUtil;
@@ -142,6 +144,8 @@ public class FunctionMetaService {
public boolean deleteFunction(NameIdentifier ident) {
FunctionPO functionPO = getFunctionPOByIdentifier(ident);
Long functionId = functionPO.functionId();
+ String metalakeName = NameIdentifierUtil.getMetalake(ident);
+ String functionFullName = ident.toString();
AtomicInteger functionDeletedCount = new AtomicInteger();
SessionUtils.doMultipleWithCommit(
@@ -179,6 +183,18 @@ public class FunctionMetaService {
mapper.softDeletePolicyMetadataObjectRelsByMetadataObject(
functionId, MetadataObject.Type.FUNCTION.name()));
}
+ },
+ () -> {
+ if (functionDeletedCount.get() > 0) {
+ SessionUtils.doWithoutCommit(
+ EntityChangeLogMapper.class,
+ mapper ->
+ mapper.insertEntityChange(
+ metalakeName,
+ Entity.EntityType.FUNCTION.name(),
+ functionFullName,
+ OperateType.DROP));
+ }
});
return functionDeletedCount.get() > 0;
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 38e325f85f..01136cc7f7 100644
---
a/core/src/test/java/org/apache/gravitino/authorization/TestAuthorizationUtils.java
+++
b/core/src/test/java/org/apache/gravitino/authorization/TestAuthorizationUtils.java
@@ -358,6 +358,31 @@ class TestAuthorizationUtils {
}
}
+ @Test
+ void testRemovePrivilegesNotifiesEntityNameIdMappingChange() {
+ GravitinoAuthorizer authorizer = Mockito.mock(GravitinoAuthorizer.class);
+ AccessControlDispatcher accessControlDispatcher =
Mockito.mock(AccessControlDispatcher.class);
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ BaseCatalog<?> baseCatalog = Mockito.mock(BaseCatalog.class);
+
Mockito.when(catalogManager.loadCatalog(Mockito.any())).thenReturn(baseCatalog);
+
+ GravitinoEnv envMock = Mockito.mock(GravitinoEnv.class);
+ Mockito.when(envMock.gravitinoAuthorizer()).thenReturn(authorizer);
+
Mockito.when(envMock.internalAccessControlDispatcher()).thenReturn(accessControlDispatcher);
+ Mockito.when(envMock.catalogManager()).thenReturn(catalogManager);
+
+ try (MockedStatic<GravitinoEnv> envStatic =
Mockito.mockStatic(GravitinoEnv.class)) {
+ envStatic.when(GravitinoEnv::getInstance).thenReturn(envMock);
+
+ NameIdentifier ident = NameIdentifier.of("metalake", "catalog",
"schema", "table");
+ AuthorizationUtils.authorizationPluginRemovePrivileges(
+ ident, Entity.EntityType.TABLE, Collections.emptyList());
+
+ Mockito.verify(authorizer)
+ .handleEntityNameIdMappingChange("metalake", ident,
Entity.EntityType.TABLE);
+ }
+ }
+
@Test
void
testRenameTablePrivilegesNotifiesAuthorizationPluginWithExpectedChange() {
NameIdentifier ident = NameIdentifier.of("metalake", "catalog", "schema",
"table");
diff --git
a/core/src/test/java/org/apache/gravitino/hook/TestFunctionHookDispatcher.java
b/core/src/test/java/org/apache/gravitino/hook/TestFunctionHookDispatcher.java
index 7b892d0ced..104bcb6847 100644
---
a/core/src/test/java/org/apache/gravitino/hook/TestFunctionHookDispatcher.java
+++
b/core/src/test/java/org/apache/gravitino/hook/TestFunctionHookDispatcher.java
@@ -19,20 +19,26 @@
package org.apache.gravitino.hook;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
-import org.apache.commons.lang3.reflect.FieldUtils;
+import org.apache.gravitino.Entity;
import org.apache.gravitino.GravitinoEnv;
import org.apache.gravitino.MetadataObject;
import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.auth.AuthConstants;
+import org.apache.gravitino.authorization.AccessControlDispatcher;
+import org.apache.gravitino.authorization.GravitinoAuthorizer;
import org.apache.gravitino.authorization.Owner;
import org.apache.gravitino.authorization.OwnerDispatcher;
import org.apache.gravitino.catalog.CatalogManager;
import org.apache.gravitino.catalog.FunctionDispatcher;
+import org.apache.gravitino.connector.BaseCatalog;
+import org.apache.gravitino.connector.authorization.AuthorizationPlugin;
import org.apache.gravitino.connector.capability.Capability;
import org.apache.gravitino.connector.capability.CapabilityResult;
import org.apache.gravitino.function.Function;
@@ -40,17 +46,13 @@ import org.apache.gravitino.function.FunctionDefinition;
import org.apache.gravitino.function.FunctionType;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;
+import org.mockito.MockedStatic;
import org.mockito.Mockito;
public class TestFunctionHookDispatcher {
@Test
public void testRegisterFunctionSetOwnerAfterRegister() throws Exception {
- GravitinoEnv gravitinoEnv = GravitinoEnv.getInstance();
- Object originalOwnerDispatcher =
- FieldUtils.readField(gravitinoEnv, "internalOwnerDispatcher", true);
- Object originalCatalogManager = FieldUtils.readField(gravitinoEnv,
"catalogManager", true);
-
NameIdentifier functionIdentifier =
NameIdentifier.of("metalake1", "catalog1", "schema1", "func1");
FunctionDefinition[] definitions = new FunctionDefinition[] {};
@@ -58,13 +60,7 @@ public class TestFunctionHookDispatcher {
Function registeredFunction = Mockito.mock(Function.class);
OwnerDispatcher ownerDispatcher = Mockito.mock(OwnerDispatcher.class);
- // Wire a case-sensitive capability so the un-normalized identifier
reaches setOwner unchanged,
- // while still exercising the normalization codepath in the hook.
- CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
- CatalogManager.CatalogWrapper catalogWrapper =
- Mockito.mock(CatalogManager.CatalogWrapper.class);
- Mockito.when(catalogWrapper.capabilities()).thenReturn(Capability.DEFAULT);
-
Mockito.when(catalogManager.loadCatalogAndWrap(any())).thenReturn(catalogWrapper);
+ CatalogManager catalogManager = catalogManagerWith(Capability.DEFAULT);
Mockito.when(
dispatcher.registerFunction(
@@ -75,38 +71,28 @@ public class TestFunctionHookDispatcher {
Mockito.eq(definitions)))
.thenReturn(registeredFunction);
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher",
ownerDispatcher, true);
- FieldUtils.writeField(gravitinoEnv, "catalogManager", catalogManager,
true);
- try {
- FunctionHookDispatcher hookDispatcher = new
FunctionHookDispatcher(dispatcher);
- Function result =
- hookDispatcher.registerFunction(
- functionIdentifier, "comment", FunctionType.SCALAR, true,
definitions);
-
- assertSame(registeredFunction, result);
-
- ArgumentCaptor<MetadataObject> metadataObjectCaptor =
- ArgumentCaptor.forClass(MetadataObject.class);
- Mockito.verify(ownerDispatcher)
- .setOwner(
- Mockito.eq("metalake1"),
- metadataObjectCaptor.capture(),
- Mockito.eq(AuthConstants.ANONYMOUS_USER),
- Mockito.eq(Owner.Type.USER));
- assertEquals(MetadataObject.Type.FUNCTION,
metadataObjectCaptor.getValue().type());
- assertEquals("catalog1.schema1.func1",
metadataObjectCaptor.getValue().fullName());
- } finally {
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher",
originalOwnerDispatcher, true);
- FieldUtils.writeField(gravitinoEnv, "catalogManager",
originalCatalogManager, true);
- }
+ FunctionHookDispatcher hookDispatcher =
+ new FunctionHookDispatcher(dispatcher, () -> ownerDispatcher,
catalogManager);
+ Function result =
+ hookDispatcher.registerFunction(
+ functionIdentifier, "comment", FunctionType.SCALAR, true,
definitions);
+
+ assertSame(registeredFunction, result);
+
+ ArgumentCaptor<MetadataObject> metadataObjectCaptor =
+ ArgumentCaptor.forClass(MetadataObject.class);
+ Mockito.verify(ownerDispatcher)
+ .setOwner(
+ Mockito.eq("metalake1"),
+ metadataObjectCaptor.capture(),
+ Mockito.eq(AuthConstants.ANONYMOUS_USER),
+ Mockito.eq(Owner.Type.USER));
+ assertEquals(MetadataObject.Type.FUNCTION,
metadataObjectCaptor.getValue().type());
+ assertEquals("catalog1.schema1.func1",
metadataObjectCaptor.getValue().fullName());
}
@Test
- public void testRegisterFunctionSucceedsWhenOwnerDispatcherIsDisabled()
throws Exception {
- GravitinoEnv gravitinoEnv = GravitinoEnv.getInstance();
- Object originalOwnerDispatcher =
- FieldUtils.readField(gravitinoEnv, "internalOwnerDispatcher", true);
-
+ public void testRegisterFunctionSucceedsWhenOwnerDispatcherIsDisabled() {
NameIdentifier functionIdentifier =
NameIdentifier.of("metalake1", "catalog1", "schema1", "func1");
FunctionDefinition[] definitions = new FunctionDefinition[] {};
@@ -122,34 +108,24 @@ public class TestFunctionHookDispatcher {
Mockito.eq(definitions)))
.thenReturn(registeredFunction);
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher", null, true);
- try {
- FunctionHookDispatcher hookDispatcher = new
FunctionHookDispatcher(dispatcher);
- Function result =
- hookDispatcher.registerFunction(
- functionIdentifier, "comment", FunctionType.SCALAR, true,
definitions);
-
- assertSame(registeredFunction, result);
- Mockito.verify(dispatcher)
- .registerFunction(functionIdentifier, "comment",
FunctionType.SCALAR, true, definitions);
- } finally {
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher",
originalOwnerDispatcher, true);
- }
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ FunctionHookDispatcher hookDispatcher =
+ new FunctionHookDispatcher(dispatcher, () -> null, catalogManager);
+ Function result =
+ hookDispatcher.registerFunction(
+ functionIdentifier, "comment", FunctionType.SCALAR, true,
definitions);
+
+ assertSame(registeredFunction, result);
+ Mockito.verify(dispatcher)
+ .registerFunction(functionIdentifier, "comment", FunctionType.SCALAR,
true, definitions);
+ Mockito.verifyNoInteractions(catalogManager);
}
@Test
public void testRegisterFunctionSetsOwnerWithNormalizedIdentifier() throws
Exception {
// Verifies the hook applies Capability.Scope.FUNCTION normalization
before setOwner, so the
// owner relation references the same identifier that NormalizeDispatcher
persists under.
- GravitinoEnv gravitinoEnv = GravitinoEnv.getInstance();
- Object originalOwnerDispatcher =
- FieldUtils.readField(gravitinoEnv, "internalOwnerDispatcher", true);
- Object originalCatalogManager = FieldUtils.readField(gravitinoEnv,
"catalogManager", true);
-
- CatalogManager mockCatalogManager = Mockito.mock(CatalogManager.class);
- CatalogManager.CatalogWrapper mockWrapper =
Mockito.mock(CatalogManager.CatalogWrapper.class);
- Mockito.when(mockWrapper.capabilities()).thenReturn(new
CaseInsensitiveCapability());
-
Mockito.when(mockCatalogManager.loadCatalogAndWrap(any())).thenReturn(mockWrapper);
+ CatalogManager catalogManager = catalogManagerWith(new
CaseInsensitiveCapability());
OwnerDispatcher mockOwnerDispatcher = Mockito.mock(OwnerDispatcher.class);
FunctionDispatcher mockFunctionDispatcher =
Mockito.mock(FunctionDispatcher.class);
@@ -160,50 +136,27 @@ public class TestFunctionHookDispatcher {
any(), any(), any(), Mockito.anyBoolean(), any()))
.thenReturn(mockFunction);
- FieldUtils.writeField(gravitinoEnv, "catalogManager", mockCatalogManager,
true);
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher",
mockOwnerDispatcher, true);
-
- try {
- FunctionHookDispatcher hook = new
FunctionHookDispatcher(mockFunctionDispatcher);
- NameIdentifier ident = NameIdentifier.of("metalake1", "catalog1",
"SCHEMA_NORM", "MY_FUNC");
- hook.registerFunction(ident, "comment", FunctionType.SCALAR, true,
definitions);
-
- ArgumentCaptor<MetadataObject> captor =
ArgumentCaptor.forClass(MetadataObject.class);
- Mockito.verify(mockOwnerDispatcher)
- .setOwner(eq("metalake1"), captor.capture(), any(),
eq(Owner.Type.USER));
- assertEquals(
- "my_func",
- captor.getValue().name(),
- "Function name passed to setOwner must be lowercased by
Capability.Scope.FUNCTION"
- + " normalization");
- assertEquals(
- "catalog1.schema_norm",
- captor.getValue().parent(),
- "Function parent (catalog.schema) must have its schema component
lowercased by"
- + " Capability.Scope.FUNCTION namespace normalization");
- } finally {
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher",
originalOwnerDispatcher, true);
- FieldUtils.writeField(gravitinoEnv, "catalogManager",
originalCatalogManager, true);
- }
+ FunctionHookDispatcher hook =
+ new FunctionHookDispatcher(
+ mockFunctionDispatcher, () -> mockOwnerDispatcher, catalogManager);
+ NameIdentifier ident = NameIdentifier.of("metalake1", "catalog1",
"SCHEMA_NORM", "MY_FUNC");
+ hook.registerFunction(ident, "comment", FunctionType.SCALAR, true,
definitions);
+
+ ArgumentCaptor<MetadataObject> captor =
ArgumentCaptor.forClass(MetadataObject.class);
+ Mockito.verify(mockOwnerDispatcher)
+ .setOwner(eq("metalake1"), captor.capture(), any(),
eq(Owner.Type.USER));
+ assertEquals("my_func", captor.getValue().name());
+ assertEquals("catalog1.schema_norm", captor.getValue().parent());
}
@Test
public void testRegisterFunctionThrowsWhenSetOwnerFails() throws Exception {
- GravitinoEnv gravitinoEnv = GravitinoEnv.getInstance();
- Object originalOwnerDispatcher =
- FieldUtils.readField(gravitinoEnv, "internalOwnerDispatcher", true);
- Object originalCatalogManager = FieldUtils.readField(gravitinoEnv,
"catalogManager", true);
-
OwnerDispatcher mockOwnerDispatcher = Mockito.mock(OwnerDispatcher.class);
Mockito.doThrow(new RuntimeException("Set owner failed"))
.when(mockOwnerDispatcher)
.setOwner(any(), any(), any(), any());
- CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
- CatalogManager.CatalogWrapper catalogWrapper =
- Mockito.mock(CatalogManager.CatalogWrapper.class);
- Mockito.when(catalogWrapper.capabilities()).thenReturn(Capability.DEFAULT);
-
Mockito.when(catalogManager.loadCatalogAndWrap(any())).thenReturn(catalogWrapper);
+ CatalogManager catalogManager = catalogManagerWith(Capability.DEFAULT);
FunctionDispatcher mockFunctionDispatcher =
Mockito.mock(FunctionDispatcher.class);
Function mockFunction = Mockito.mock(Function.class);
@@ -213,25 +166,74 @@ public class TestFunctionHookDispatcher {
any(), any(), any(), Mockito.anyBoolean(), any()))
.thenReturn(mockFunction);
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher",
mockOwnerDispatcher, true);
- FieldUtils.writeField(gravitinoEnv, "catalogManager", catalogManager,
true);
+ FunctionHookDispatcher hook =
+ new FunctionHookDispatcher(
+ mockFunctionDispatcher, () -> mockOwnerDispatcher, catalogManager);
+ NameIdentifier ident =
+ NameIdentifier.of("metalake1", "catalog1", "schema_owner_fail",
"func_owner_fail");
+ RuntimeException thrown =
+ assertThrows(
+ RuntimeException.class,
+ () -> hook.registerFunction(ident, "comment", FunctionType.SCALAR,
true, definitions));
+ assertEquals("Set owner failed", thrown.getMessage());
+ }
- try {
- FunctionHookDispatcher hook = new
FunctionHookDispatcher(mockFunctionDispatcher);
- NameIdentifier ident =
- NameIdentifier.of("metalake1", "catalog1", "schema_owner_fail",
"func_owner_fail");
+ @Test
+ public void
testDropFunctionInvalidatesOnlyInternalCacheWithNormalizedIdentifier()
+ throws Exception {
+ NameIdentifier functionIdentifier =
+ NameIdentifier.of("metalake1", "catalog1", "SCHEMA1", "FUNC1");
+ NameIdentifier normalizedIdentifier =
+ NameIdentifier.of("metalake1", "catalog1", "schema1", "func1");
+ FunctionDispatcher dispatcher = Mockito.mock(FunctionDispatcher.class);
+ Mockito.when(dispatcher.dropFunction(functionIdentifier))
+ .thenReturn(true, false)
+ .thenThrow(new RuntimeException("Drop failed"));
+ CatalogManager catalogManager = catalogManagerWith(new
CaseInsensitiveCapability());
+
+ GravitinoAuthorizer authorizer = Mockito.mock(GravitinoAuthorizer.class);
+ AuthorizationPlugin catalogAuthorizationPlugin =
Mockito.mock(AuthorizationPlugin.class);
+ BaseCatalog<?> catalog = Mockito.mock(BaseCatalog.class);
+
Mockito.when(catalog.getAuthorizationPlugin()).thenReturn(catalogAuthorizationPlugin);
+ Mockito.when(catalogManager.loadCatalog(any())).thenReturn(catalog);
+
+ GravitinoEnv env = Mockito.mock(GravitinoEnv.class);
+ Mockito.when(env.gravitinoAuthorizer()).thenReturn(authorizer);
+ Mockito.when(env.accessControlDispatcher())
+ .thenReturn(Mockito.mock(AccessControlDispatcher.class));
+ Mockito.when(env.catalogManager()).thenReturn(catalogManager);
+
+ try (MockedStatic<GravitinoEnv> envStatic =
Mockito.mockStatic(GravitinoEnv.class)) {
+ envStatic.when(GravitinoEnv::getInstance).thenReturn(env);
+ FunctionHookDispatcher hookDispatcher =
+ new FunctionHookDispatcher(dispatcher, () -> null, catalogManager);
+
+ assertTrue(hookDispatcher.dropFunction(functionIdentifier));
+ assertFalse(hookDispatcher.dropFunction(functionIdentifier));
RuntimeException thrown =
assertThrows(
- RuntimeException.class,
- () ->
- hook.registerFunction(ident, "comment", FunctionType.SCALAR,
true, definitions));
- assertEquals("Set owner failed", thrown.getMessage());
- } finally {
- FieldUtils.writeField(gravitinoEnv, "internalOwnerDispatcher",
originalOwnerDispatcher, true);
- FieldUtils.writeField(gravitinoEnv, "catalogManager",
originalCatalogManager, true);
+ RuntimeException.class, () ->
hookDispatcher.dropFunction(functionIdentifier));
+ assertEquals("Drop failed", thrown.getMessage());
+ Mockito.verify(authorizer)
+ .handleEntityNameIdMappingChange(
+ "metalake1", normalizedIdentifier, Entity.EntityType.FUNCTION);
+ // Ranger HadoopSQL is one catalog plugin that rejects FUNCTION
metadata. The generic plugin
+ // mock represents that integration boundary and must not receive a
removal callback.
+ Mockito.verifyNoInteractions(catalogAuthorizationPlugin);
+ Mockito.verify(catalogManager, Mockito.never()).loadCatalog(any());
+ Mockito.verify(catalogManager,
Mockito.times(1)).loadCatalogAndWrap(any());
}
}
+ private static CatalogManager catalogManagerWith(Capability capability)
throws Exception {
+ CatalogManager catalogManager = Mockito.mock(CatalogManager.class);
+ CatalogManager.CatalogWrapper catalogWrapper =
+ Mockito.mock(CatalogManager.CatalogWrapper.class);
+ Mockito.when(catalogWrapper.capabilities()).thenReturn(capability);
+
Mockito.when(catalogManager.loadCatalogAndWrap(any())).thenReturn(catalogWrapper);
+ return catalogManager;
+ }
+
private static class CaseInsensitiveCapability implements Capability {
@Override
public CapabilityResult caseSensitiveOnName(Scope scope) {
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestEntityChangeLogService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestEntityChangeLogService.java
index f58c6b9da3..1ab16e5b91 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestEntityChangeLogService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestEntityChangeLogService.java
@@ -19,6 +19,9 @@
package org.apache.gravitino.storage.relational.service;
import java.io.IOException;
+import java.sql.Connection;
+import java.sql.SQLException;
+import java.sql.Statement;
import java.util.List;
import java.util.Map;
import org.apache.gravitino.Catalog;
@@ -27,6 +30,7 @@ import org.apache.gravitino.Namespace;
import org.apache.gravitino.meta.BaseMetalake;
import org.apache.gravitino.meta.CatalogEntity;
import org.apache.gravitino.meta.FilesetEntity;
+import org.apache.gravitino.meta.FunctionEntity;
import org.apache.gravitino.meta.ModelEntity;
import org.apache.gravitino.meta.SchemaEntity;
import org.apache.gravitino.meta.TableEntity;
@@ -37,9 +41,11 @@ import
org.apache.gravitino.storage.relational.TestJDBCBackend;
import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
import org.apache.gravitino.storage.relational.po.cache.OperateType;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
import org.apache.gravitino.storage.relational.utils.SessionUtils;
import org.apache.gravitino.utils.NameIdentifierUtil;
import org.apache.gravitino.utils.NamespaceUtil;
+import org.apache.ibatis.session.SqlSession;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.TestTemplate;
@@ -361,5 +367,57 @@ public class TestEntityChangeLogService extends
TestJDBCBackend {
Entity.EntityType.MODEL,
NameIdentifierUtil.ofModel(METALAKE_NAME, CATALOG_NAME, SCHEMA_NAME,
"model2").toString(),
OperateType.DROP);
+
+ FunctionEntity function =
+ createFunctionEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofFunction(METALAKE_NAME, CATALOG_NAME, SCHEMA_NAME),
+ "function1",
+ AUDIT_INFO);
+ FunctionMetaService.getInstance().insertFunction(function, false);
+ long maxIdBeforeFunctionDrop = maxEntityChangeId();
+ Assertions.assertTrue(
+ backend.delete(function.nameIdentifier(), Entity.EntityType.FUNCTION,
false));
+ assertEntityChange(
+ maxIdBeforeFunctionDrop,
+ METALAKE_NAME,
+ Entity.EntityType.FUNCTION,
+ NameIdentifierUtil.ofFunction(METALAKE_NAME, CATALOG_NAME,
SCHEMA_NAME, "function1")
+ .toString(),
+ OperateType.DROP);
+ }
+
+ @TestTemplate
+ void testFunctionDropRollsBackWhenChangeLogInsertFails() throws Exception {
+ createParentEntities(METALAKE_NAME, CATALOG_NAME, SCHEMA_NAME, AUDIT_INFO);
+ FunctionEntity function =
+ createFunctionEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofFunction(METALAKE_NAME, CATALOG_NAME, SCHEMA_NAME),
+ "function_rolled_back",
+ AUDIT_INFO);
+ FunctionMetaService.getInstance().insertFunction(function, false);
+
+ renameTable("entity_change_log", "entity_change_log_bak");
+ try {
+ Assertions.assertThrows(
+ Exception.class,
+ () -> backend.delete(function.nameIdentifier(),
Entity.EntityType.FUNCTION, false));
+ } finally {
+ renameTable("entity_change_log_bak", "entity_change_log");
+ }
+
+ FunctionEntity persistedFunction =
+
FunctionMetaService.getInstance().getFunctionByIdentifier(function.nameIdentifier());
+ Assertions.assertEquals(function.id(), persistedFunction.id());
+ }
+
+ private void renameTable(String from, String to) throws SQLException {
+ try (SqlSession session =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = session.getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute(String.format("ALTER TABLE %s RENAME TO %s", from,
to));
+ }
}
}
diff --git
a/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinChangePoller.java
b/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinChangePoller.java
index 6bd330fc4c..33f1046a45 100644
---
a/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinChangePoller.java
+++
b/server-common/src/test/java/org/apache/gravitino/server/authorization/jcasbin/TestJcasbinChangePoller.java
@@ -102,6 +102,20 @@ public class TestJcasbinChangePoller {
Assertions.assertEquals(List.of(), metadataIdCache.invalidatedKeys);
}
+ @Test
+ void testPollEntityChangesInvalidatesFunctionByExactKey() {
+ RecordingCache<String, Long> metadataIdCache = new RecordingCache<>();
+ RecordingCache<Long, Optional<OwnerInfo>> ownerRelCache = new
RecordingCache<>();
+
+ JcasbinChangeListener poller = new JcasbinChangeListener(metadataIdCache,
ownerRelCache, 1);
+ poller.onEntityChange(List.of(change(1L, MetadataObject.Type.FUNCTION,
"ml1.cat1.sch1.func1")));
+
+ Assertions.assertEquals(
+ List.of(key("ml1", "CATALOG", "cat1", "SCHEMA", "sch1", "FUNCTION",
"func1")),
+ metadataIdCache.invalidatedKeys);
+ Assertions.assertEquals(List.of(), metadataIdCache.invalidatedPrefixes);
+ }
+
@Test
void testPollCursorAdvancementIsSynchronized() throws NoSuchMethodException {
Method pollOwnerChanges =
JcasbinChangeListener.class.getDeclaredMethod("pollOwnerChanges");