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

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


The following commit(s) were added to refs/heads/master by this push:
     new 46fc1e0a6fa Pipe: Fixed the tsFile parsing & write-back-sink auto 
create db bug (#15240)
46fc1e0a6fa is described below

commit 46fc1e0a6fad9614d871dbe67b3ce137be4ea33a
Author: Caideyipi <[email protected]>
AuthorDate: Tue Apr 1 00:06:14 2025 +0800

    Pipe: Fixed the tsFile parsing & write-back-sink auto create db bug (#15240)
---
 .../pipe/it/dual/tablemodel/TableModelUtils.java   |  8 ++--
 .../pipe/it/single/IoTDBPipePermissionIT.java      | 43 ++++++++++++++++++++++
 .../agent/task/connection/PipeEventCollector.java  |  4 +-
 .../subtask/processor/PipeProcessorSubtask.java    | 15 +++++++-
 .../protocol/writeback/WriteBackConnector.java     | 10 +++++
 .../common/tsfile/PipeTsFileInsertionEvent.java    |  4 +-
 6 files changed, 75 insertions(+), 9 deletions(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
index 17b470296cb..cc154a9d943 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/tablemodel/TableModelUtils.java
@@ -100,11 +100,11 @@ public class TableModelUtils {
   public static boolean insertData(
       final String dataBaseName,
       final String tableName,
-      final int start,
-      final int end,
+      final int startInclusive,
+      final int endExclusive,
       final BaseEnv baseEnv) {
-    List<String> list = new ArrayList<>(end - start + 1);
-    for (int i = start; i < end; ++i) {
+    List<String> list = new ArrayList<>(endExclusive - startInclusive + 1);
+    for (int i = startInclusive; i < endExclusive; ++i) {
       list.add(
           String.format(
               "insert into %s (s0, s3, s2, s1, s4, s5, s6, s7, s8, s9, s10, 
s11, time) values ('t%s','t%s','t%s','t%s','%s', %s.0, %s, %s, %d, %d.0, '%s', 
'%s', %s)",
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
index 005eb49afef..64afeeb594d 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
@@ -32,6 +32,7 @@ import org.junit.runner.RunWith;
 import java.sql.Connection;
 import java.sql.SQLException;
 import java.sql.Statement;
+import java.util.Arrays;
 
 import static org.junit.Assert.fail;
 
@@ -154,4 +155,46 @@ public class IoTDBPipePermissionIT extends 
AbstractPipeSingleIT {
 
     TableModelUtils.assertCountData("test", "test", 100, env);
   }
+
+  @Test
+  public void testSinkPermissionWithHistoricalDataAndTablePattern() {
+    TableModelUtils.createDataBaseAndTable(env, "test", "test1");
+    TableModelUtils.createDataBaseAndTable(env, "test1", "test1");
+    TableModelUtils.createDataBaseAndTable(env, "test", "test");
+    TableModelUtils.createDataBaseAndTable(env, "test1", "test");
+
+    if (!TestUtils.tryExecuteNonQueriesWithRetry(
+        "test",
+        BaseEnv.TABLE_SQL_DIALECT,
+        env,
+        Arrays.asList(
+            "create user thulab 'passwd'", "grant INSERT on test.test1 to user 
thulab"))) {
+      return;
+    }
+
+    // Write some data
+    if (!TableModelUtils.insertData("test1", "test", 0, 100, env)) {
+      return;
+    }
+
+    if (!TableModelUtils.insertData("test1", "test1", 0, 100, env)) {
+      return;
+    }
+
+    // Use current session, user is root
+    try (final Connection connection = 
env.getConnection(BaseEnv.TABLE_SQL_DIALECT);
+        final Statement statement = connection.createStatement()) {
+      statement.execute(
+          "create pipe a2b "
+              + "with source ('database'='test1', 'table'='test1') "
+              + "with processor('processor'='rename-database-processor', 
'processor.new-db-name'='test') "
+              + "with sink ('sink'='write-back-sink', 'username'='thulab', 
'password'='passwd')");
+    } catch (final SQLException e) {
+      e.printStackTrace();
+      fail("Create pipe without user shall succeed if use the current 
session");
+    }
+
+    TableModelUtils.assertCountData("test", "test", 0, env);
+    TableModelUtils.assertCountData("test", "test1", 100, env);
+  }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
index 6e265711ad0..3bc4553c852 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
@@ -135,9 +135,7 @@ public class PipeEventCollector implements EventCollector {
       return;
     }
 
-    if (!forceTabletFormat
-        && !sourceEvent.shouldParse4Privilege()
-        && canSkipParsing4TsFileEvent(sourceEvent)) {
+    if (!forceTabletFormat && canSkipParsing4TsFileEvent(sourceEvent)) {
       collectEvent(sourceEvent);
       return;
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
index f611b18ad61..40352766630 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java
@@ -31,6 +31,7 @@ import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
 import org.apache.iotdb.db.pipe.agent.task.connection.PipeEventCollector;
 import org.apache.iotdb.db.pipe.event.UserDefinedEnrichedEvent;
 import org.apache.iotdb.db.pipe.event.common.heartbeat.PipeHeartbeatEvent;
+import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
 import 
org.apache.iotdb.db.pipe.metric.overview.PipeDataNodeRemainingEventAndTimeMetrics;
 import org.apache.iotdb.db.pipe.metric.processor.PipeProcessorMetrics;
 import org.apache.iotdb.db.pipe.processor.pipeconsensus.PipeConsensusProcessor;
@@ -141,7 +142,19 @@ public class PipeProcessorSubtask extends 
PipeReportableSubtask {
           pipeProcessor.process((TabletInsertionEvent) event, 
outputEventCollector);
           PipeProcessorMetrics.getInstance().markTabletEvent(taskID);
         } else if (event instanceof TsFileInsertionEvent) {
-          pipeProcessor.process((TsFileInsertionEvent) event, 
outputEventCollector);
+          // We have to parse the privilege first, to avoid passing 
no-privilege data to processor
+          if (event instanceof PipeTsFileInsertionEvent
+              && ((PipeTsFileInsertionEvent) event).shouldParse4Privilege()) {
+            try (final PipeTsFileInsertionEvent tsFileInsertionEvent =
+                (PipeTsFileInsertionEvent) event) {
+              for (final TabletInsertionEvent tabletInsertionEvent :
+                  tsFileInsertionEvent.toTabletInsertionEvents()) {
+                pipeProcessor.process(tabletInsertionEvent, 
outputEventCollector);
+              }
+            }
+          } else {
+            pipeProcessor.process((TsFileInsertionEvent) event, 
outputEventCollector);
+          }
           PipeProcessorMetrics.getInstance().markTsFileEvent(taskID);
           PipeDataNodeRemainingEventAndTimeMetrics.getInstance()
               .markTsFileCollectInvocationCount(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
index 201f776b9e6..0bc4f76d25c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/writeback/WriteBackConnector.java
@@ -374,6 +374,16 @@ public class WriteBackConnector implements PipeConnector {
       return;
     }
 
+    try {
+      Coordinator.getInstance()
+          .getAccessControl()
+          .checkCanCreateDatabase(session.getUsername(), database);
+    } catch (final AccessDeniedException e) {
+      // Auto create failed, we still check if there are existing databases
+      // If there are not, this will be removed by catching database not 
exists exception
+      ALREADY_CREATED_DATABASES.add(database);
+      return;
+    }
     final TDatabaseSchema schema = new TDatabaseSchema(new 
TDatabaseSchema(database));
     schema.setIsTableModel(true);
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
index 04cf1d58e59..3077b05237e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
@@ -672,7 +672,9 @@ public class PipeTsFileInsertionEvent extends 
PipeInsertionEvent
                   startTime,
                   endTime,
                   pipeTaskMeta,
-                  userName,
+                  // Do not parse privilege if it should not be parsed
+                  // To avoid renaming of the tsFile database
+                  shouldParse4Privilege ? userName : null,
                   this)
               .provide());
       return eventParser.get();

Reply via email to