This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new b35e19f38d [core] Add privilege checks for partition and consumer ops 
(#8929)
b35e19f38d is described below

commit b35e19f38d055267712b5fc13d57f45e219970d9
Author: cxzl25 <[email protected]>
AuthorDate: Mon Aug 3 19:13:34 2026 +0800

    [core] Add privilege checks for partition and consumer ops (#8929)
---
 .../apache/paimon/privilege/PrivilegedCatalog.java | 30 +++++++++++++++++++
 .../paimon/privilege/PrivilegedFileStoreTable.java |  9 ++++++
 .../paimon/privilege/PrivilegedCatalogTest.java    | 34 ++++++++++++++++++++++
 .../flink/procedure/ClearConsumersProcedure.java   |  6 +---
 .../flink/procedure/ResetConsumerProcedure.java    | 12 ++------
 .../paimon/flink/action/ClearConsumerAction.java   |  6 +---
 .../paimon/flink/action/ResetConsumerAction.java   |  6 +---
 .../flink/procedure/ClearConsumersProcedure.java   |  6 +---
 .../flink/procedure/ResetConsumerProcedure.java    |  6 +---
 .../spark/procedure/ClearConsumersProcedure.java   |  6 +---
 .../spark/procedure/ResetConsumerProcedure.java    |  6 +---
 11 files changed, 82 insertions(+), 45 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java 
b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java
index b408055e51..0aadcff860 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedCatalog.java
@@ -27,6 +27,7 @@ import org.apache.paimon.catalog.PropertyChange;
 import org.apache.paimon.options.ConfigOption;
 import org.apache.paimon.options.ConfigOptions;
 import org.apache.paimon.options.Options;
+import org.apache.paimon.partition.PartitionStatistics;
 import org.apache.paimon.schema.Schema;
 import org.apache.paimon.schema.SchemaChange;
 import org.apache.paimon.table.FileStoreTable;
@@ -165,6 +166,35 @@ public class PrivilegedCatalog extends DelegateCatalog {
         wrapped.markDonePartitions(identifier, partitions);
     }
 
+    @Override
+    public void createPartitions(Identifier identifier, List<Map<String, 
String>> partitions)
+            throws TableNotExistException {
+        privilegeManager.getPrivilegeChecker().assertCanInsert(identifier);
+        wrapped.createPartitions(identifier, partitions);
+    }
+
+    @Override
+    public void createPartitions(
+            Identifier identifier, List<Map<String, String>> partitions, 
boolean ignoreIfExists)
+            throws TableNotExistException {
+        privilegeManager.getPrivilegeChecker().assertCanInsert(identifier);
+        wrapped.createPartitions(identifier, partitions, ignoreIfExists);
+    }
+
+    @Override
+    public void dropPartitions(Identifier identifier, List<Map<String, 
String>> partitions)
+            throws TableNotExistException {
+        privilegeManager.getPrivilegeChecker().assertCanInsert(identifier);
+        wrapped.dropPartitions(identifier, partitions);
+    }
+
+    @Override
+    public void alterPartitions(Identifier identifier, 
List<PartitionStatistics> partitions)
+            throws TableNotExistException {
+        privilegeManager.getPrivilegeChecker().assertCanInsert(identifier);
+        wrapped.alterPartitions(identifier, partitions);
+    }
+
     public void createPrivilegedUser(String user, String password) {
         privilegeManager.createUser(user, password);
     }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java
 
b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java
index 36e8d03340..5a309ae102 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/privilege/PrivilegedFileStoreTable.java
@@ -21,6 +21,7 @@ package org.apache.paimon.privilege;
 import org.apache.paimon.FileStore;
 import org.apache.paimon.Snapshot;
 import org.apache.paimon.catalog.Identifier;
+import org.apache.paimon.consumer.ConsumerManager;
 import org.apache.paimon.schema.TableSchema;
 import org.apache.paimon.stats.Statistics;
 import org.apache.paimon.table.DelegatedFileStoreTable;
@@ -72,6 +73,14 @@ public class PrivilegedFileStoreTable extends 
DelegatedFileStoreTable {
         return wrapped.changelogManager();
     }
 
+    @Override
+    public ConsumerManager consumerManager() {
+        // Resetting/deleting a consumer's progress affects what data 
downstream streaming
+        // readers will (re)consume, so treat it as a write/insert-level 
operation.
+        privilegeChecker.assertCanInsert(identifier);
+        return wrapped.consumerManager();
+    }
+
     @Override
     public Optional<Snapshot> latestSnapshot() {
         privilegeChecker.assertCanSelectOrInsert(identifier);
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/privilege/PrivilegedCatalogTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/privilege/PrivilegedCatalogTest.java
index a6070d5fcc..745ae77761 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/privilege/PrivilegedCatalogTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/privilege/PrivilegedCatalogTest.java
@@ -27,6 +27,9 @@ import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.function.Executable;
 
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
 import java.util.Optional;
 
 import static org.assertj.core.api.Assertions.assertThat;
@@ -75,6 +78,37 @@ public class PrivilegedCatalogTest extends 
FileSystemCatalogTest {
         assertThatThrownBy(() -> dataTable2.snapshot(0)).isNotNull();
     }
 
+    @Test
+    public void testCreatePartitionsWithIgnoreIfExistsRequiresPrivilege() 
throws Exception {
+        catalog.createDatabase("test_db", false);
+
+        Identifier identifier = Identifier.create("test_db", "test_table");
+        catalog.createTable(identifier, DEFAULT_TABLE_SCHEMA, false);
+
+        PrivilegedCatalog rootCatalog = ((PrivilegedCatalog) catalog);
+        rootCatalog.createPrivilegedUser(USERNAME_TEST_USER, 
PASSWORD_TEST_USER);
+        Catalog userCatalogWithoutPrivilege =
+                create(rootCatalog.wrapped(), USERNAME_TEST_USER, 
PASSWORD_TEST_USER);
+
+        Map<String, String> partitionSpec = new HashMap<>();
+        partitionSpec.put("pk", "0");
+
+        // The 3-arg overload (identifier, partitions, ignoreIfExists) must 
also be guarded:
+        // without this, a user could bypass the INSERT privilege check that 
the 2-arg overload
+        // enforces, since DelegateCatalog#createPartitions(3 args) forwards 
straight to the
+        // wrapped catalog.
+        assertNoPrivilege(
+                () ->
+                        userCatalogWithoutPrivilege.createPartitions(
+                                identifier, 
Collections.singletonList(partitionSpec), true));
+
+        rootCatalog.grantPrivilegeOnTable(USERNAME_TEST_USER, identifier, 
PrivilegeType.INSERT);
+        Catalog userCatalogWithPrivilege =
+                create(rootCatalog.wrapped(), USERNAME_TEST_USER, 
PASSWORD_TEST_USER);
+        userCatalogWithPrivilege.createPartitions(
+                identifier, Collections.singletonList(partitionSpec), true);
+    }
+
     private PrivilegedCatalog create(Catalog catalog, String user, String 
password) {
         return new PrivilegedCatalog(
                 catalog, new FileBasedPrivilegeManagerLoader(warehouse, 
fileIO, user, password));
diff --git 
a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
 
b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
index a3f6713f4c..cd2090e375 100644
--- 
a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
+++ 
b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
@@ -56,11 +56,7 @@ public class ClearConsumersProcedure extends ProcedureBase {
             throws Catalog.TableNotExistException {
         FileStoreTable fileStoreTable =
                 (FileStoreTable) 
catalog.getTable(Identifier.fromString(tableId));
-        ConsumerManager consumerManager =
-                new ConsumerManager(
-                        fileStoreTable.fileIO(),
-                        fileStoreTable.location(),
-                        fileStoreTable.snapshotManager().branch());
+        ConsumerManager consumerManager = fileStoreTable.consumerManager();
 
         Pattern includingPattern =
                 StringUtils.isNullOrWhitespaceOnly(includingConsumers)
diff --git 
a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
 
b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
index 7777ccda19..85b0f39ed5 100644
--- 
a/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
+++ 
b/paimon-flink/paimon-flink-1.18/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
@@ -50,11 +50,7 @@ public class ResetConsumerProcedure extends ProcedureBase {
         FileStoreTable fileStoreTable =
                 (FileStoreTable) 
catalog.getTable(Identifier.fromString(tableId));
         fileStoreTable.snapshotManager().snapshot(nextSnapshotId);
-        ConsumerManager consumerManager =
-                new ConsumerManager(
-                        fileStoreTable.fileIO(),
-                        fileStoreTable.location(),
-                        fileStoreTable.snapshotManager().branch());
+        ConsumerManager consumerManager = fileStoreTable.consumerManager();
         consumerManager.resetConsumer(consumerId, new 
Consumer(nextSnapshotId));
 
         return new String[] {"Success"};
@@ -64,11 +60,7 @@ public class ResetConsumerProcedure extends ProcedureBase {
             throws Catalog.TableNotExistException {
         FileStoreTable fileStoreTable =
                 (FileStoreTable) 
catalog.getTable(Identifier.fromString(tableId));
-        ConsumerManager consumerManager =
-                new ConsumerManager(
-                        fileStoreTable.fileIO(),
-                        fileStoreTable.location(),
-                        fileStoreTable.snapshotManager().branch());
+        ConsumerManager consumerManager = fileStoreTable.consumerManager();
         consumerManager.deleteConsumer(consumerId);
 
         return new String[] {"Success"};
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java
index 124bf28fed..485b021a61 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ClearConsumerAction.java
@@ -51,11 +51,7 @@ public class ClearConsumerAction extends TableActionBase 
implements LocalAction
     @Override
     public void executeLocally() {
         FileStoreTable dataTable = (FileStoreTable) table;
-        ConsumerManager consumerManager =
-                new ConsumerManager(
-                        dataTable.fileIO(),
-                        dataTable.location(),
-                        dataTable.snapshotManager().branch());
+        ConsumerManager consumerManager = dataTable.consumerManager();
 
         Pattern includingPattern =
                 StringUtils.isNullOrWhitespaceOnly(includingConsumers)
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java
index cebc3b1959..418360f108 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/action/ResetConsumerAction.java
@@ -48,11 +48,7 @@ public class ResetConsumerAction extends TableActionBase 
implements LocalAction
     @Override
     public void executeLocally() throws Exception {
         FileStoreTable dataTable = (FileStoreTable) table;
-        ConsumerManager consumerManager =
-                new ConsumerManager(
-                        dataTable.fileIO(),
-                        dataTable.location(),
-                        dataTable.snapshotManager().branch());
+        ConsumerManager consumerManager = dataTable.consumerManager();
         if (Objects.isNull(nextSnapshotId)) {
             consumerManager.deleteConsumer(consumerId);
         } else {
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
index de4c371d30..3a510fc3cb 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ClearConsumersProcedure.java
@@ -71,11 +71,7 @@ public class ClearConsumersProcedure extends ProcedureBase {
             throws Catalog.TableNotExistException {
         FileStoreTable fileStoreTable =
                 (FileStoreTable) 
catalog.getTable(Identifier.fromString(tableId));
-        ConsumerManager consumerManager =
-                new ConsumerManager(
-                        fileStoreTable.fileIO(),
-                        fileStoreTable.location(),
-                        fileStoreTable.snapshotManager().branch());
+        ConsumerManager consumerManager = fileStoreTable.consumerManager();
 
         Pattern includingPattern =
                 StringUtils.isNullOrWhitespaceOnly(includingConsumers)
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
index 934ce182a0..2b35e4b1e4 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/ResetConsumerProcedure.java
@@ -61,11 +61,7 @@ public class ResetConsumerProcedure extends ProcedureBase {
             throws Catalog.TableNotExistException {
         FileStoreTable fileStoreTable =
                 (FileStoreTable) 
catalog.getTable(Identifier.fromString(tableId));
-        ConsumerManager consumerManager =
-                new ConsumerManager(
-                        fileStoreTable.fileIO(),
-                        fileStoreTable.location(),
-                        fileStoreTable.snapshotManager().branch());
+        ConsumerManager consumerManager = fileStoreTable.consumerManager();
         if (nextSnapshotId != null) {
             fileStoreTable.snapshotManager().snapshot(nextSnapshotId);
             consumerManager.resetConsumer(consumerId, new 
Consumer(nextSnapshotId));
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java
index cdde1c6f83..1847b9cf3e 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ClearConsumersProcedure.java
@@ -97,11 +97,7 @@ public class ClearConsumersProcedure extends BaseProcedure {
                 tableIdent,
                 table -> {
                     FileStoreTable fileStoreTable = (FileStoreTable) table;
-                    ConsumerManager consumerManager =
-                            new ConsumerManager(
-                                    fileStoreTable.fileIO(),
-                                    fileStoreTable.location(),
-                                    fileStoreTable.snapshotManager().branch());
+                    ConsumerManager consumerManager = 
fileStoreTable.consumerManager();
                     consumerManager.clearConsumers(includingPattern, 
excludingPattern);
 
                     InternalRow outputRow = newInternalRow(true);
diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java
index 0f7fabd05d..856693175d 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/ResetConsumerProcedure.java
@@ -82,11 +82,7 @@ public class ResetConsumerProcedure extends BaseProcedure {
                 tableIdent,
                 table -> {
                     FileStoreTable fileStoreTable = (FileStoreTable) table;
-                    ConsumerManager consumerManager =
-                            new ConsumerManager(
-                                    fileStoreTable.fileIO(),
-                                    fileStoreTable.location(),
-                                    fileStoreTable.snapshotManager().branch());
+                    ConsumerManager consumerManager = 
fileStoreTable.consumerManager();
                     if (nextSnapshotId == null) {
                         consumerManager.deleteConsumer(consumerId);
                     } else {

Reply via email to