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

Caideyipi pushed a commit to branch fix/pipe-table-batch-isolation-master
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 158c11dc873641da732b927eb572b9a12e86c7af
Author: Caideyipi <[email protected]>
AuthorDate: Thu Sep 3 11:40:37 2026 +0800

    Fix table-model Pipe batch isolation
---
 .../request/PipeTransferTabletBatchReqV2.java      |  39 +++---
 .../db/pipe/sink/util/cacher/LeaderCacheUtils.java |   6 +-
 .../plan/analyze/schema/SchemaValidator.java       |  37 +++++-
 .../plan/relational/planner/TableModelPlanner.java |   6 +-
 .../pipe/sink/PipeDataNodeThriftRequestTest.java   | 133 +++++++++++++++++++++
 .../sink/util/cacher/LeaderCacheUtilsTest.java     |  62 ++++++++++
 .../plan/analyze/schema/SchemaValidatorTest.java   | 101 ++++++++++++++++
 .../relational/planner/TableModelPlannerTest.java  |  78 ++++++++++++
 8 files changed, 438 insertions(+), 24 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
index c79cab7d88f..2cec219fa75 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/payload/evolvable/request/PipeTransferTabletBatchReqV2.java
@@ -64,8 +64,9 @@ public class PipeTransferTabletBatchReqV2 extends 
TPipeTransferReq {
     final List<InsertBaseStatement> statements =
         new ArrayList<>(insertNodeReqs.size() + tabletReqs.size());
 
-    final Map<String, List<InsertRowStatement>> 
tableModelDatabaseInsertRowStatementMap =
-        new LinkedHashMap<>();
+    // Keep permission checks, schema validation, and redirect metadata scoped 
to one table.
+    final Map<String, Map<String, List<InsertRowStatement>>>
+        tableModelDatabaseInsertRowStatementMap = new LinkedHashMap<>();
     final Map<String, List<InsertRowStatement>> 
treeModelDatabaseInsertRowStatementMap =
         new LinkedHashMap<>();
     final Map<String, List<InsertTabletStatement>> 
treeModelDatabaseInsertTabletStatementMap =
@@ -78,17 +79,15 @@ public class PipeTransferTabletBatchReqV2 extends 
TPipeTransferReq {
       }
       if (statement.isWriteToTable()) {
         if (statement instanceof InsertRowStatement) {
-          tableModelDatabaseInsertRowStatementMap
-              .computeIfAbsent(statement.getDatabaseName().get(), k -> new 
ArrayList<>())
-              .add((InsertRowStatement) statement);
+          addTableModelInsertRowStatement(
+              tableModelDatabaseInsertRowStatementMap, (InsertRowStatement) 
statement);
         } else if (statement instanceof InsertTabletStatement) {
           statements.add(statement);
         } else if (statement instanceof InsertRowsStatement) {
           for (final InsertRowStatement insertRowStatement :
               ((InsertRowsStatement) statement).getInsertRowStatementList()) {
-            tableModelDatabaseInsertRowStatementMap
-                .computeIfAbsent(insertRowStatement.getDatabaseName().get(), k 
-> new ArrayList<>())
-                .add(insertRowStatement);
+            addTableModelInsertRowStatement(
+                tableModelDatabaseInsertRowStatementMap, insertRowStatement);
           }
         } else {
           throw new UnsupportedOperationException(
@@ -141,18 +140,30 @@ public class PipeTransferTabletBatchReqV2 extends 
TPipeTransferReq {
     addTreeModelInsertRowsStatements(statements, 
treeModelDatabaseInsertRowStatementMap);
     addTreeModelInsertTabletsStatements(statements, 
treeModelDatabaseInsertTabletStatementMap);
 
-    for (final Map.Entry<String, List<InsertRowStatement>> insertRows :
+    for (final Map.Entry<String, Map<String, List<InsertRowStatement>>> 
insertRows :
         tableModelDatabaseInsertRowStatementMap.entrySet()) {
-      final InsertRowsStatement statement = new InsertRowsStatement();
-      statement.setWriteToTable(true);
-      statement.setDatabaseName(insertRows.getKey());
-      statement.setInsertRowStatementList(insertRows.getValue());
-      statements.add(statement);
+      for (final Map.Entry<String, List<InsertRowStatement>> tableInsertRows :
+          insertRows.getValue().entrySet()) {
+        final InsertRowsStatement statement = new InsertRowsStatement();
+        statement.setWriteToTable(true);
+        statement.setDatabaseName(insertRows.getKey());
+        statement.setInsertRowStatementList(tableInsertRows.getValue());
+        statements.add(statement);
+      }
     }
 
     return statements;
   }
 
+  private static void addTableModelInsertRowStatement(
+      final Map<String, Map<String, List<InsertRowStatement>>> 
databaseInsertRowStatementMap,
+      final InsertRowStatement insertRowStatement) {
+    databaseInsertRowStatementMap
+        .computeIfAbsent(insertRowStatement.getDatabaseName().get(), k -> new 
LinkedHashMap<>())
+        .computeIfAbsent(insertRowStatement.getTableName(), k -> new 
ArrayList<>())
+        .add(insertRowStatement);
+  }
+
   private void addTreeModelInsertRowsStatements(
       final List<InsertBaseStatement> statements,
       final Map<String, List<InsertRowStatement>> 
databaseInsertRowStatementMap) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
index 754b406dcd7..0f6beade80d 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtils.java
@@ -41,11 +41,11 @@ public class LeaderCacheUtils {
    * @return a list of pairs, each pair contains a device path and its 
redirect endpoint.
    */
   public static List<Pair<String, TEndPoint>> 
parseRecommendedRedirections(TSStatus status) {
-    // If there is no exception, there should be 2 sub-statuses, one for 
InsertRowsStatement and one
-    // for InsertMultiTabletsStatement (see 
IoTDBDataNodeReceiver#handleTransferTabletBatch).
+    // Each top-level sub-status corresponds to one statement constructed by 
the receiver. V2 batch
+    // requests may contain any number of statements because rows are grouped 
by database and table.
     final List<Pair<String, TEndPoint>> redirectList = new ArrayList<>();
 
-    if (status.getSubStatusSize() != 2) {
+    if (!status.isSetSubStatus()) {
       return redirectList;
     }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
index 6123afb1d9e..06905feead4 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidator.java
@@ -25,11 +25,14 @@ import org.apache.iotdb.commons.path.PartialPath;
 import 
org.apache.iotdb.commons.queryengine.plan.relational.metadata.QualifiedObjectName;
 import org.apache.iotdb.db.queryengine.common.MPPQueryContext;
 import org.apache.iotdb.db.queryengine.common.schematree.ISchemaTree;
+import org.apache.iotdb.db.queryengine.plan.analyze.AnalyzeUtils;
 import org.apache.iotdb.db.queryengine.plan.relational.metadata.Metadata;
 import org.apache.iotdb.db.queryengine.plan.relational.security.AccessControl;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.InsertRows;
 import 
org.apache.iotdb.db.queryengine.plan.relational.sql.ast.WrappedInsertStatement;
 import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
 import 
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertMultiTabletsStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement;
 import 
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsOfOneDeviceStatement;
 import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
 
@@ -39,9 +42,12 @@ import org.apache.tsfile.file.metadata.enums.TSEncoding;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.util.LinkedHashSet;
 import java.util.List;
+import java.util.Set;
 
 import static org.apache.iotdb.commons.utils.PathUtils.unQualifyDatabaseName;
+import static 
org.apache.iotdb.db.queryengine.plan.execution.config.TableConfigTaskVisitor.DATABASE_NOT_SPECIFIED;
 
 public class SchemaValidator {
 
@@ -71,11 +77,10 @@ public class SchemaValidator {
       final MPPQueryContext context,
       AccessControl accessControl) {
     try {
-      accessControl.checkCanInsertIntoTable(
-          context.getSession().getUserName(),
-          new QualifiedObjectName(
-              unQualifyDatabaseName(insertStatement.getDatabase()), 
insertStatement.getTableName()),
-          context);
+      for (final QualifiedObjectName targetTable : 
getTargetTables(insertStatement, context)) {
+        accessControl.checkCanInsertIntoTable(
+            context.getSession().getUserName(), targetTable, context);
+      }
       insertStatement.validateTableSchema(metadata, context);
       insertStatement.updateAfterSchemaValidation(context);
       insertStatement.validateDeviceSchema(metadata, context);
@@ -85,6 +90,28 @@ public class SchemaValidator {
     }
   }
 
+  private static Set<QualifiedObjectName> getTargetTables(
+      final WrappedInsertStatement insertStatement, final MPPQueryContext 
context) {
+    final Set<QualifiedObjectName> targetTables = new LinkedHashSet<>();
+    if (insertStatement instanceof InsertRows) {
+      for (final InsertRowStatement rowStatement :
+          ((InsertRows) 
insertStatement).getInnerTreeStatement().getInsertRowStatementList()) {
+        final String database = AnalyzeUtils.getDatabaseName(rowStatement, 
context);
+        if (database == null) {
+          throw new SemanticException(DATABASE_NOT_SPECIFIED);
+        }
+        targetTables.add(
+            new QualifiedObjectName(unQualifyDatabaseName(database), 
rowStatement.getTableName()));
+      }
+    } else {
+      targetTables.add(
+          new QualifiedObjectName(
+              unQualifyDatabaseName(insertStatement.getDatabase()),
+              insertStatement.getTableName()));
+    }
+    return targetTables;
+  }
+
   public static ISchemaTree validate(
       ISchemaFetcher schemaFetcher,
       List<PartialPath> devicePaths,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
index c8433d95416..cca616e0a3b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlanner.java
@@ -55,6 +55,7 @@ import 
org.apache.iotdb.db.queryengine.plan.scheduler.ClusterScheduler;
 import org.apache.iotdb.db.queryengine.plan.scheduler.IScheduler;
 import org.apache.iotdb.db.queryengine.plan.scheduler.load.LoadTsFileScheduler;
 import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
 import 
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertTabletStatement;
 import org.apache.iotdb.rpc.RpcUtils;
 import org.apache.iotdb.rpc.TSStatusCode;
@@ -238,8 +239,9 @@ public class TableModelPlanner implements IPlanner {
         ((WrappedInsertStatement) statementToRedirect).getInnerTreeStatement();
 
     if (!analysis.isFinishQueryAfterAnalyze()) {
-      // Table Model Session only supports insertTablet
-      if (insertStatement instanceof InsertTabletStatement) {
+      // Table Model Session supports insertTablet and pipe-generated 
insertRows statements.
+      if (insertStatement instanceof InsertTabletStatement
+          || insertStatement instanceof InsertRowsStatement) {
         if (tsstatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) 
{
           boolean needRedirect = false;
           List<TEndPoint> redirectNodeList = analysis.getRedirectNodeList();
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
index 143d195d284..b780d19f01c 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/PipeDataNodeThriftRequestTest.java
@@ -53,6 +53,7 @@ import 
org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferTsFil
 import 
org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferTsFileSealWithModReq;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.CreateAlignedTimeSeriesNode;
 import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
+import 
org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowsNode;
 import org.apache.iotdb.db.queryengine.plan.statement.Statement;
 import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertBaseStatement;
 import 
org.apache.iotdb.db.queryengine.plan.statement.crud.InsertMultiTabletsStatement;
@@ -1048,6 +1049,138 @@ public class PipeDataNodeThriftRequestTest {
         new HashSet<>(java.util.Arrays.asList("root.db1", "root.db2")), 
insertTabletsDatabases);
   }
 
+  @Test
+  public void testPipeTransferTabletBatchReqV2SeparatesTableModelTables() 
throws IOException {
+    final List<ByteBuffer> insertNodeBuffers = new ArrayList<>();
+    final List<String> insertNodeDataBases = new ArrayList<>();
+
+    insertNodeBuffers.add(
+        new InsertRowNode(
+                new PlanNodeId(""),
+                new PartialPath("table1", false),
+                false,
+                new String[] {"s"},
+                new TSDataType[] {TSDataType.INT32},
+                1,
+                new Object[] {1},
+                false)
+            .serializeToByteBuffer());
+    insertNodeDataBases.add("db1");
+
+    insertNodeBuffers.add(
+        new InsertRowNode(
+                new PlanNodeId(""),
+                new PartialPath("table2", false),
+                false,
+                new String[] {"s"},
+                new TSDataType[] {TSDataType.INT32},
+                2,
+                new Object[] {2},
+                false)
+            .serializeToByteBuffer());
+    insertNodeDataBases.add("db1");
+
+    insertNodeBuffers.add(
+        new InsertRowNode(
+                new PlanNodeId(""),
+                new PartialPath("table1", false),
+                false,
+                new String[] {"s"},
+                new TSDataType[] {TSDataType.INT32},
+                3,
+                new Object[] {3},
+                false)
+            .serializeToByteBuffer());
+    insertNodeDataBases.add("db1");
+
+    final PipeTransferTabletBatchReqV2 request =
+        PipeTransferTabletBatchReqV2.fromTPipeTransferReq(
+            PipeTransferTabletBatchReqV2.toTPipeTransferReq(
+                insertNodeBuffers,
+                Collections.emptyList(),
+                insertNodeDataBases,
+                Collections.emptyList()));
+
+    final List<InsertBaseStatement> statements = request.constructStatements();
+
+    Assert.assertEquals(2, statements.size());
+    final InsertRowsStatement table1Statement = (InsertRowsStatement) 
statements.get(0);
+    final InsertRowsStatement table2Statement = (InsertRowsStatement) 
statements.get(1);
+    Assert.assertTrue(table1Statement.isWriteToTable());
+    Assert.assertTrue(table2Statement.isWriteToTable());
+    Assert.assertEquals("db1", table1Statement.getDatabaseName().get());
+    Assert.assertEquals("db1", table2Statement.getDatabaseName().get());
+    Assert.assertEquals(2, table1Statement.getInsertRowStatementList().size());
+    Assert.assertEquals(1, table2Statement.getInsertRowStatementList().size());
+    Assert.assertEquals(
+        "table1", 
table1Statement.getInsertRowStatementList().get(0).getTableName());
+    Assert.assertEquals(
+        "table1", 
table1Statement.getInsertRowStatementList().get(1).getTableName());
+    Assert.assertEquals(
+        "table2", 
table2Statement.getInsertRowStatementList().get(0).getTableName());
+  }
+
+  @Test
+  public void 
testPipeTransferTabletBatchReqV2SeparatesTablesWithinInsertRowsNode()
+      throws IOException {
+    final InsertRowsNode insertRowsNode = new InsertRowsNode(new 
PlanNodeId("rows"));
+    insertRowsNode.addOneInsertRowNode(
+        new InsertRowNode(
+            new PlanNodeId("row1"),
+            new PartialPath("table1", false),
+            false,
+            new String[] {"s"},
+            new TSDataType[] {TSDataType.INT32},
+            1,
+            new Object[] {1},
+            false),
+        0);
+    insertRowsNode.addOneInsertRowNode(
+        new InsertRowNode(
+            new PlanNodeId("row2"),
+            new PartialPath("table2", false),
+            false,
+            new String[] {"s"},
+            new TSDataType[] {TSDataType.INT32},
+            2,
+            new Object[] {2},
+            false),
+        1);
+    insertRowsNode.addOneInsertRowNode(
+        new InsertRowNode(
+            new PlanNodeId("row3"),
+            new PartialPath("table1", false),
+            false,
+            new String[] {"s"},
+            new TSDataType[] {TSDataType.INT32},
+            3,
+            new Object[] {3},
+            false),
+        2);
+
+    final PipeTransferTabletBatchReqV2 request =
+        PipeTransferTabletBatchReqV2.fromTPipeTransferReq(
+            PipeTransferTabletBatchReqV2.toTPipeTransferReq(
+                
Collections.singletonList(insertRowsNode.serializeToByteBuffer()),
+                Collections.emptyList(),
+                Collections.singletonList("db1"),
+                Collections.emptyList()));
+
+    final List<InsertBaseStatement> statements = request.constructStatements();
+
+    Assert.assertEquals(2, statements.size());
+    final InsertRowsStatement table1Statement = (InsertRowsStatement) 
statements.get(0);
+    final InsertRowsStatement table2Statement = (InsertRowsStatement) 
statements.get(1);
+    Assert.assertEquals(2, table1Statement.getInsertRowStatementList().size());
+    Assert.assertEquals(1, table2Statement.getInsertRowStatementList().size());
+    Assert.assertEquals(
+        "table1", 
table1Statement.getInsertRowStatementList().get(0).getTableName());
+    Assert.assertEquals(
+        "table1", 
table1Statement.getInsertRowStatementList().get(1).getTableName());
+    Assert.assertEquals(
+        "table2", 
table2Statement.getInsertRowStatementList().get(0).getTableName());
+  }
+
   @Test
   public void testPipeTransferFilePieceReq() throws IOException {
     final byte[] body = "testPipeTransferFilePieceReq".getBytes();
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
new file mode 100644
index 00000000000..76c6bef2467
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/util/cacher/LeaderCacheUtilsTest.java
@@ -0,0 +1,62 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.pipe.sink.util.cacher;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.apache.tsfile.utils.Pair;
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+public class LeaderCacheUtilsTest {
+
+  @Test
+  public void testParseRecommendedRedirectionsFromVariableStatementCount() {
+    final TEndPoint redirectEndPoint = new TEndPoint("127.0.0.2", 6667);
+    final TSStatus redirectedRowStatus =
+        RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)
+            .setMessage("table1.device1")
+            .setRedirectNode(redirectEndPoint);
+    final TSStatus redirectedStatementStatus =
+        RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+            .setSubStatus(Collections.singletonList(redirectedRowStatus));
+    final TSStatus batchStatus =
+        RpcUtils.getStatus(TSStatusCode.REDIRECTION_RECOMMEND)
+            .setSubStatus(
+                Arrays.asList(
+                    RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS),
+                    redirectedStatementStatus,
+                    RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS)));
+
+    final List<Pair<String, TEndPoint>> redirects =
+        LeaderCacheUtils.parseRecommendedRedirections(batchStatus);
+
+    Assert.assertEquals(1, redirects.size());
+    Assert.assertEquals("table1.device1", redirects.get(0).getLeft());
+    Assert.assertEquals(redirectEndPoint, redirects.get(0).getRight());
+  }
+}
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidatorTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidatorTest.java
new file mode 100644
index 00000000000..8cb3cfb1ee1
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/analyze/schema/SchemaValidatorTest.java
@@ -0,0 +1,101 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.queryengine.plan.analyze.schema;
+
+import org.apache.iotdb.commons.exception.auth.AccessDeniedException;
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.queryengine.common.SessionInfo;
+import org.apache.iotdb.commons.queryengine.common.SqlDialect;
+import 
org.apache.iotdb.commons.queryengine.plan.relational.metadata.QualifiedObjectName;
+import org.apache.iotdb.db.queryengine.common.MPPQueryContext;
+import org.apache.iotdb.db.queryengine.common.QueryId;
+import org.apache.iotdb.db.queryengine.plan.relational.metadata.Metadata;
+import org.apache.iotdb.db.queryengine.plan.relational.security.AccessControl;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.InsertRows;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowStatement;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
+
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.InOrder;
+import org.mockito.Mockito;
+
+import java.time.ZoneId;
+import java.util.Arrays;
+
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.same;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+import static org.mockito.Mockito.verifyZeroInteractions;
+
+public class SchemaValidatorTest {
+
+  private static final String USER = "pipe-user";
+
+  @Test
+  public void testAllInsertRowsTablesCheckedBeforeSchemaValidation() {
+    final MPPQueryContext context =
+        new MPPQueryContext(
+            "",
+            new QueryId("check_all_insert_rows_tables"),
+            new SessionInfo(1L, USER, ZoneId.systemDefault(), "db1", 
SqlDialect.TABLE),
+            null,
+            null);
+    final InsertRowsStatement statement = new InsertRowsStatement();
+    statement.setWriteToTable(true);
+    statement.setInsertRowStatementList(
+        Arrays.asList(
+            createInsertRowStatement("db1", "table1"),
+            createInsertRowStatement("db1", "table2"),
+            createInsertRowStatement("db1", "table1")));
+
+    final Metadata metadata = Mockito.mock(Metadata.class);
+    final AccessControl accessControl = Mockito.mock(AccessControl.class);
+    final QualifiedObjectName table1 = new QualifiedObjectName("db1", 
"table1");
+    final QualifiedObjectName table2 = new QualifiedObjectName("db1", 
"table2");
+    Mockito.doThrow(new AccessDeniedException("denied"))
+        .when(accessControl)
+        .checkCanInsertIntoTable(eq(USER), eq(table2), same(context));
+
+    Assert.assertThrows(
+        AccessDeniedException.class,
+        () ->
+            SchemaValidator.validate(
+                metadata, new InsertRows(statement, context), context, 
accessControl));
+
+    final InOrder inOrder = Mockito.inOrder(accessControl);
+    inOrder.verify(accessControl).checkCanInsertIntoTable(eq(USER), 
eq(table1), same(context));
+    inOrder.verify(accessControl).checkCanInsertIntoTable(eq(USER), 
eq(table2), same(context));
+    verify(accessControl, times(1)).checkCanInsertIntoTable(eq(USER), 
eq(table1), same(context));
+    verifyNoMoreInteractions(accessControl);
+    verifyZeroInteractions(metadata);
+  }
+
+  private static InsertRowStatement createInsertRowStatement(
+      final String database, final String table) {
+    final InsertRowStatement statement = new InsertRowStatement();
+    statement.setWriteToTable(true);
+    statement.setDatabaseName(database);
+    statement.setDevicePath(new PartialPath(table, false));
+    return statement;
+  }
+}
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlannerTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlannerTest.java
new file mode 100644
index 00000000000..99c0ed02357
--- /dev/null
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/TableModelPlannerTest.java
@@ -0,0 +1,78 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.queryengine.plan.relational.planner;
+
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.db.queryengine.plan.relational.analyzer.Analysis;
+import org.apache.iotdb.db.queryengine.plan.relational.sql.ast.InsertRows;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.InsertRowsStatement;
+import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+
+public class TableModelPlannerTest {
+
+  @Test
+  public void testSetRedirectInfoForInsertRows() {
+    final InsertRowsStatement insertRowsStatement = new InsertRowsStatement();
+    insertRowsStatement.setInsertRowStatementList(Collections.emptyList());
+    final Analysis analysis =
+        new Analysis(new InsertRows(insertRowsStatement, null), 
Collections.emptyMap());
+    final TEndPoint localEndPoint = new TEndPoint("127.0.0.1", 6667);
+    final TEndPoint remoteEndPoint = new TEndPoint("127.0.0.2", 6667);
+    analysis.setRedirectNodeList(Arrays.asList(localEndPoint, remoteEndPoint));
+    final TSStatus status = RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS);
+
+    createPlanner().setRedirectInfo(analysis, localEndPoint, status);
+
+    Assert.assertEquals(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode(), 
status.getCode());
+    final List<TSStatus> subStatus = status.getSubStatus();
+    Assert.assertEquals(2, subStatus.size());
+    Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(), 
subStatus.get(0).getCode());
+    Assert.assertFalse(subStatus.get(0).isSetRedirectNode());
+    Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(), 
subStatus.get(1).getCode());
+    Assert.assertEquals(remoteEndPoint, subStatus.get(1).getRedirectNode());
+  }
+
+  private static TableModelPlanner createPlanner() {
+    return new TableModelPlanner(
+        null,
+        null,
+        null,
+        null,
+        null,
+        null,
+        null,
+        Collections.emptyList(),
+        Collections.emptyList(),
+        null,
+        null,
+        Collections.emptyList(),
+        Collections.emptyMap(),
+        null);
+  }
+}

Reply via email to