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

CRZbulabula 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 c8bf0bcfbf3 Validate region operation requests before submitting 
procedures (#18696)
c8bf0bcfbf3 is described below

commit c8bf0bcfbf36e89f745075e61224652dc0b202ed
Author: Yongzao <[email protected]>
AuthorDate: Thu Sep 24 16:40:53 2026 +0800

    Validate region operation requests before submitting procedures (#18696)
    
    * Improve region operation submission results in CLI
    
    * Validate complete region operation requests before submitting procedures
    
    * Reject invalid multi-region submissions in ITs
---
 .../commit/IoTDBMigrateMultiRegionForIoTV1IT.java  |   9 +-
 .../IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java | 123 +++---
 .../java/org/apache/iotdb/cli/AbstractCliTest.java |  44 ++
 .../org/apache/iotdb/jdbc/IoTDBStatementTest.java  |  39 ++
 .../iotdb/confignode/i18n/ManagerMessages.java     |  14 +
 .../iotdb/confignode/i18n/ManagerMessages.java     |  14 +
 .../iotdb/confignode/manager/ProcedureManager.java | 482 ++++++++-------------
 .../ProcedureManagerReconstructRegionTest.java     | 190 --------
 .../ProcedureManagerRegionOperationTest.java       | 401 +++++++++++++++++
 .../queryengine/execution/ConfigExecutionTest.java |  27 ++
 10 files changed, 770 insertions(+), 573 deletions(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBMigrateMultiRegionForIoTV1IT.java
 
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBMigrateMultiRegionForIoTV1IT.java
index da361b32a7c..b915f5f83f0 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBMigrateMultiRegionForIoTV1IT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBMigrateMultiRegionForIoTV1IT.java
@@ -36,6 +36,7 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.sql.Connection;
+import java.sql.SQLException;
 import java.sql.Statement;
 import java.util.ArrayList;
 import java.util.List;
@@ -125,14 +126,8 @@ public class IoTDBMigrateMultiRegionForIoTV1IT extends 
IoTDBRegionOperationRelia
                 try {
                   statement.execute(command);
                   return true;
-                } catch (Exception e) {
+                } catch (SQLException e) {
                   String errorMessage = e.getMessage();
-                  if (errorMessage != null
-                      && errorMessage.contains("successfully submitted")
-                      && errorMessage.contains("failed to submit")) {
-                    LOGGER.warn("Multi-region migrate partially succeeded: 
{}", errorMessage);
-                    return true;
-                  }
                   LOGGER.warn("Multi-region migrate failed, retrying: {}", 
errorMessage);
                   return false;
                 }
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
 
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
index 8bb1671b2d3..c4902cbc931 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
@@ -26,6 +26,7 @@ import org.apache.iotdb.consensus.ConsensusFactory;
 import org.apache.iotdb.it.env.EnvFactory;
 import org.apache.iotdb.it.framework.IoTDBTestRunner;
 import org.apache.iotdb.itbase.category.ClusterIT;
+import org.apache.iotdb.rpc.TSStatusCode;
 
 import org.awaitility.Awaitility;
 import org.junit.Assert;
@@ -168,14 +169,13 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
         Assert.assertFalse(
             "ConfigNode should not throw NullPointerException, but got: " + 
message,
             message.contains("NullPointerException"));
-        // ... and the submission must be rejected cleanly. "extend region" 
wraps every region's
-        // result, so the top-level message only reports the aggregate counts; 
the concrete "does
-        // not
-        // exist in the cluster" reason is carried in the per-region 
sub-status.
+        // The operation-specific error and the concrete reason must reach the 
JDBC client.
+        Assert.assertEquals(TSStatusCode.EXTEND_REGION_ERROR.getStatusCode(), 
e.getErrorCode());
         Assert.assertTrue(
             "Expected the extend submission to be rejected but got: " + 
message,
-            message.contains("failed to submit: 1"));
+            message.contains("Target DataNode " + invalidDataNodeId + " does 
not exist"));
       }
+      assertRegionMapUnchanged(statement, regionMap);
     }
   }
 
@@ -314,7 +314,7 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
     }
   }
 
-  /** Test multi-region expand with partial regions already in target DataNode 
*/
+  /** Reject the entire expansion when a later region already exists on the 
target DataNode. */
   @Test
   public void multiRegionExpandPartialExistTest() throws Exception {
     EnvFactory.getEnv()
@@ -338,40 +338,33 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
       Map<Integer, Set<Integer>> regionMap = getAllRegionMap(statement);
       Set<Integer> allDataNodeId = getAllDataNodes(statement);
 
-      List<Integer> allRegions = new ArrayList<>(regionMap.keySet());
-      List<Integer> selectedRegions = allRegions.subList(0, Math.min(3, 
allRegions.size()));
+      Assert.assertEquals(2, regionMap.size());
+      List<Integer> selectedRegions = new ArrayList<>(regionMap.keySet());
 
       int targetDataNode =
           findDataNodeNotContainsAnyRegion(allDataNodeId, regionMap, 
selectedRegions);
 
-      // first expand some regions individually
-      List<Integer> preExpandRegions =
-          selectedRegions.subList(0, Math.min(2, selectedRegions.size()));
-      for (int regionId : preExpandRegions) {
-        regionGroupExpand(statement, client, regionId, targetDataNode);
-      }
-
-      // now try to expand all regions (including already expanded ones)
-      LOGGER.info(
-          "Testing multi-expand with regions {} to DataNode {}, where {} 
already exist",
-          selectedRegions,
-          targetDataNode,
-          preExpandRegions);
-
-      multiRegionGroupExpand(statement, client, selectedRegions, 
targetDataNode);
-
-      // verify all regions are in target DataNode
+      // Keep the first region valid so this also catches submission during 
validation.
+      regionGroupExpand(statement, client, selectedRegions.get(1), 
targetDataNode);
       regionMap = getAllRegionMap(statement);
-      for (int regionId : selectedRegions) {
-        Assert.assertTrue(
-            "Region " + regionId + " should contain target DataNode " + 
targetDataNode,
-            regionMap.get(regionId).contains(targetDataNode));
-      }
-      LOGGER.info("Multi-region expand partial exist test passed");
+      
Assert.assertFalse(regionMap.get(selectedRegions.get(0)).contains(targetDataNode));
+      
Assert.assertTrue(regionMap.get(selectedRegions.get(1)).contains(targetDataNode));
+
+      SQLException exception =
+          Assert.assertThrows(
+              SQLException.class,
+              () ->
+                  statement.execute(
+                      buildMultiRegionCommand(
+                          MULTI_EXPAND_FORMAT, selectedRegions, 
targetDataNode)));
+      Assert.assertEquals(
+          TSStatusCode.EXTEND_REGION_ERROR.getStatusCode(), 
exception.getErrorCode());
+      Assert.assertTrue(exception.getMessage().contains("already contains 
region"));
+      assertRegionMapUnchanged(statement, regionMap);
     }
   }
 
-  /** Test multi-region shrink with partial regions not in target DataNode */
+  /** Reject the entire removal when a later region does not exist on the 
target DataNode. */
   @Test
   public void multiRegionShrinkPartialNotExistTest() throws Exception {
     EnvFactory.getEnv()
@@ -379,8 +372,8 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
         .getCommonConfig()
         .setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
         
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
-        .setDataReplicationFactor(1)
-        .setSchemaReplicationFactor(1);
+        .setDataReplicationFactor(2)
+        .setSchemaReplicationFactor(2);
 
     EnvFactory.getEnv().initClusterEnvironment(1, 5);
 
@@ -395,8 +388,8 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
       Map<Integer, Set<Integer>> regionMap = getAllRegionMap(statement);
       Set<Integer> allDataNodeId = getAllDataNodes(statement);
 
-      List<Integer> allRegions = new ArrayList<>(regionMap.keySet());
-      List<Integer> selectedRegions = allRegions.subList(0, Math.min(3, 
allRegions.size()));
+      Assert.assertEquals(2, regionMap.size());
+      List<Integer> selectedRegions = new ArrayList<>(regionMap.keySet());
 
       int targetDataNode =
           findDataNodeNotContainsAnyRegion(allDataNodeId, regionMap, 
selectedRegions);
@@ -404,33 +397,34 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
       // first expand all regions to target DataNode
       multiRegionGroupExpand(statement, client, selectedRegions, 
targetDataNode);
 
-      // then shrink some regions individually
-      List<Integer> preShrinkRegions =
-          selectedRegions.subList(0, Math.min(2, selectedRegions.size()));
-      for (int regionId : preShrinkRegions) {
-        regionGroupShrink(statement, client, regionId, targetDataNode);
-      }
-
-      // now try to shrink all regions (including already shrunk ones)
-      LOGGER.info(
-          "Testing multi-shrink with regions {} from DataNode {}, where {} 
already removed",
-          selectedRegions,
-          targetDataNode,
-          preShrinkRegions);
-
-      multiRegionGroupShrink(statement, client, selectedRegions, 
targetDataNode);
-
-      // verify all regions are not in target DataNode
+      // Keep the first region valid and leave two replicas of the second 
region elsewhere.
+      regionGroupShrink(statement, client, selectedRegions.get(1), 
targetDataNode);
       regionMap = getAllRegionMap(statement);
-      for (int regionId : selectedRegions) {
-        Assert.assertFalse(
-            "Region " + regionId + " should not contain target DataNode " + 
targetDataNode,
-            regionMap.get(regionId).contains(targetDataNode));
-      }
-      LOGGER.info("Multi-region shrink partial not exist test passed");
+      
Assert.assertTrue(regionMap.get(selectedRegions.get(0)).contains(targetDataNode));
+      
Assert.assertFalse(regionMap.get(selectedRegions.get(1)).contains(targetDataNode));
+
+      SQLException exception =
+          Assert.assertThrows(
+              SQLException.class,
+              () ->
+                  statement.execute(
+                      buildMultiRegionCommand(
+                          MULTI_SHRINK_FORMAT, selectedRegions, 
targetDataNode)));
+      Assert.assertEquals(
+          TSStatusCode.REMOVE_REGION_PEER_ERROR.getStatusCode(), 
exception.getErrorCode());
+      Assert.assertTrue(exception.getMessage().contains("doesn't contain 
Region"));
+      assertRegionMapUnchanged(statement, regionMap);
     }
   }
 
+  private void assertRegionMapUnchanged(
+      Statement statement, Map<Integer, Set<Integer>> expectedRegionMap) {
+    Awaitility.await()
+        .during(3, TimeUnit.SECONDS)
+        .atMost(10, TimeUnit.SECONDS)
+        .untilAsserted(() -> Assert.assertEquals(expectedRegionMap, 
getAllRegionMap(statement)));
+  }
+
   private void multiRegionGroupExpand(
       Statement statement,
       SyncConfigNodeIServiceClient client,
@@ -519,17 +513,8 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
               try {
                 statement.execute(command);
                 return true;
-              } catch (Exception e) {
+              } catch (SQLException e) {
                 String errorMessage = e.getMessage();
-                // If error message contains both "successfully submitted" and 
"failed to submit",
-                // consider it as partial success and continue
-                if (errorMessage != null
-                    && errorMessage.contains("successfully submitted")
-                    && errorMessage.contains("failed to submit")) {
-                  LOGGER.warn(
-                      "Multi-region {} partially succeeded: {}", 
operationType, errorMessage);
-                  return true;
-                }
                 LOGGER.warn(
                     "Multi-region {} command execution failed, retrying: {}",
                     operationType,
diff --git 
a/iotdb-client/cli/src/test/java/org/apache/iotdb/cli/AbstractCliTest.java 
b/iotdb-client/cli/src/test/java/org/apache/iotdb/cli/AbstractCliTest.java
index 01644d58412..deba69d50be 100644
--- a/iotdb-client/cli/src/test/java/org/apache/iotdb/cli/AbstractCliTest.java
+++ b/iotdb-client/cli/src/test/java/org/apache/iotdb/cli/AbstractCliTest.java
@@ -22,9 +22,13 @@ package org.apache.iotdb.cli;
 import org.apache.iotdb.cli.AbstractCli.OperationResult;
 import org.apache.iotdb.cli.type.ExitType;
 import org.apache.iotdb.cli.utils.CliContext;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.exception.ArgsErrorException;
 import org.apache.iotdb.jdbc.IoTDBConnection;
+import org.apache.iotdb.jdbc.IoTDBConnectionParams;
 import org.apache.iotdb.jdbc.IoTDBDatabaseMetadata;
+import org.apache.iotdb.jdbc.IoTDBSQLException;
+import org.apache.iotdb.rpc.TSStatusCode;
 
 import org.apache.commons.cli.CommandLine;
 import org.apache.commons.cli.CommandLineParser;
@@ -44,6 +48,7 @@ import java.io.PrintStream;
 import java.lang.reflect.Field;
 import java.lang.reflect.Method;
 import java.nio.charset.StandardCharsets;
+import java.sql.Statement;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.List;
@@ -52,6 +57,7 @@ import static org.junit.Assert.assertEquals;
 import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertTrue;
 import static org.junit.Assert.fail;
+import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
 
 public class AbstractCliTest {
@@ -71,6 +77,44 @@ public class AbstractCliTest {
   public void tearDown() throws Exception {
     setStaticField("lineCount", 0);
     setStaticField("isReachEnd", false);
+    AbstractCli.lastProcessStatus = AbstractCli.CODE_OK;
+  }
+
+  @Test
+  public void testRegionValidationErrorIsPrintedAndReturnsErrorStatus() throws 
Exception {
+    String[] statements = {
+      "MIGRATE REGION 12,99 FROM 6 TO 7",
+      "RECONSTRUCT REGION 12,99 ON 7",
+      "EXTEND REGION 12,99 TO 7",
+      "REMOVE REGION 12,99 FROM 7"
+    };
+    TSStatusCode[] codes = {
+      TSStatusCode.MIGRATE_REGION_ERROR,
+      TSStatusCode.RECONSTRUCT_REGION_ERROR,
+      TSStatusCode.EXTEND_REGION_ERROR,
+      TSStatusCode.REMOVE_REGION_PEER_ERROR
+    };
+    when(connection.getParams())
+        .thenReturn(new IoTDBConnectionParams("jdbc:iotdb://localhost:6667/"));
+    for (int i = 0; i < statements.length; i++) {
+      ByteArrayOutputStream out = new ByteArrayOutputStream();
+      try (PrintStream printer = new PrintStream(out, true, 
StandardCharsets.UTF_8.name())) {
+        CliContext ctx = new CliContext(System.in, printer, System.err, 
ExitType.EXCEPTION);
+        Statement statement = mock(Statement.class);
+        when(connection.createStatement()).thenReturn(statement);
+        TSStatus status =
+            new TSStatus(codes[i].getStatusCode()).setMessage("Region 99 does 
not exist");
+        when(statement.execute(statements[i]))
+            .thenThrow(new IoTDBSQLException(status.getMessage(), status));
+
+        AbstractCli.handleInputCmd(ctx, statements[i], connection);
+
+        assertEquals(AbstractCli.CODE_ERROR, AbstractCli.lastProcessStatus);
+        String output = new String(out.toByteArray(), StandardCharsets.UTF_8);
+        assertTrue(output.contains(status.getMessage()));
+        assertFalse(output.contains("The statement is executed successfully"));
+      }
+    }
   }
 
   @Test
diff --git 
a/iotdb-client/jdbc/src/test/java/org/apache/iotdb/jdbc/IoTDBStatementTest.java 
b/iotdb-client/jdbc/src/test/java/org/apache/iotdb/jdbc/IoTDBStatementTest.java
index 5cfd5342e05..2f7d993e37d 100644
--- 
a/iotdb-client/jdbc/src/test/java/org/apache/iotdb/jdbc/IoTDBStatementTest.java
+++ 
b/iotdb-client/jdbc/src/test/java/org/apache/iotdb/jdbc/IoTDBStatementTest.java
@@ -19,8 +19,11 @@
 
 package org.apache.iotdb.jdbc;
 
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.rpc.RpcUtils;
+import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.service.rpc.thrift.IClientRPCService.Iface;
+import org.apache.iotdb.service.rpc.thrift.TSExecuteStatementResp;
 import org.apache.iotdb.service.rpc.thrift.TSFetchMetadataReq;
 import org.apache.iotdb.service.rpc.thrift.TSFetchMetadataResp;
 
@@ -106,4 +109,40 @@ public class IoTDBStatementTest {
     statement.setQueryTimeout(100);
     Assert.assertEquals(100, statement.getQueryTimeout());
   }
+
+  @SuppressWarnings("resource")
+  @Test
+  public void regionValidationErrorsArePropagatedByExecuteAndExecuteUpdate() 
throws Exception {
+    String[] statements = {
+      "MIGRATE REGION 1,1 FROM 2 TO 3",
+      "RECONSTRUCT REGION 1,1 ON 2",
+      "EXTEND REGION 1,1 TO 2",
+      "REMOVE REGION 1,1 FROM 2"
+    };
+    TSStatusCode[] codes = {
+      TSStatusCode.MIGRATE_REGION_ERROR,
+      TSStatusCode.RECONSTRUCT_REGION_ERROR,
+      TSStatusCode.EXTEND_REGION_ERROR,
+      TSStatusCode.REMOVE_REGION_PEER_ERROR
+    };
+    for (int i = 0; i < statements.length; i++) {
+      final String sql = statements[i];
+      TSStatus status =
+          new TSStatus(codes[i].getStatusCode()).setMessage("Duplicate Region 
ID 1 in the request");
+      TSExecuteStatementResp response = new 
TSExecuteStatementResp().setStatus(status);
+      when(client.executeStatementV2(any())).thenReturn(response);
+      when(client.executeUpdateStatement(any())).thenReturn(response);
+      IoTDBStatement statement = new IoTDBStatement(connection, client, 
sessionId, zoneID, 0, 1L);
+
+      SQLException executeError =
+          Assert.assertThrows(SQLException.class, () -> 
statement.execute(sql));
+      assertEquals(status.getCode(), executeError.getErrorCode());
+      
Assert.assertTrue(executeError.getMessage().contains(status.getMessage()));
+      SQLException updateError =
+          Assert.assertThrows(SQLException.class, () -> 
statement.executeUpdate(sql));
+      assertEquals(status.getCode(), updateError.getErrorCode());
+      
Assert.assertTrue(updateError.getMessage().contains(status.getMessage()));
+      Assert.assertNull(statement.getWarnings());
+    }
+  }
 }
diff --git 
a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java
 
b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java
index 03ec8003251..6a36ac71302 100644
--- 
a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java
+++ 
b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java
@@ -295,6 +295,20 @@ public final class ManagerMessages {
   public static final String
       
LOG_SKIP_NON_EXISTENT_REGION_ID_ARG_IN_RECONSTRUCTREGION_REQUEST_TO_DATANODE_ARG_7F76D789
 =
           "Skip non-existent Region ID {} in ReconstructRegion request to 
DataNode {}.";
+  public static final String 
MESSAGE_DUPLICATE_REGION_ID_ARG_IN_THE_REQUEST_B6FFCCFC =
+      "Duplicate Region ID %d in the request";
+  public static final String MESSAGE_REGION_IDS_MUST_NOT_BE_EMPTY_B42DAAFD =
+      "Region IDs must not be empty";
+  public static final String 
MESSAGE_SOURCE_AND_TARGET_DATANODE_IDS_MUST_BE_DIFFERENT_ARG_286D3838 =
+      "Source and target DataNode IDs must be different: %d";
+  public static final String 
LOG_SUBMIT_REGION_OPERATION_PROCEDURE_SUCCESSFULLY_ARG_90468B38 =
+      "Submit region operation procedure successfully: {}";
+  public static final String MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9 =
+      "Region %d does not exist";
+  public static final String 
MESSAGE_SOURCE_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_2255633C =
+      "Source DataNode %s does not exist in the cluster";
+  public static final String 
MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF =
+      "Target DataNode %s does not exist in the cluster";
   public static final String 
MIGRATEREGION_SUBMIT_REGIONMIGRATEPROCEDURE_SUCCESSFULLY_REGION_ORIGIN_DATANODE 
=
       "[MigrateRegion] Submit RegionMigrateProcedure successfully, Region: {}, 
Origin DataNode: {}, Dest DataNode: {}, Add Coordinator: {}, Remove 
Coordinator: {}";
   public static final String 
SUBMIT_REGIONMIGRATEPROCEDURE_FAILED_BECAUSE_REGIONGROUP_DOESN_T_EXIST =
diff --git 
a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java
 
b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java
index f145e04052e..3c35171dd06 100644
--- 
a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java
+++ 
b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java
@@ -293,6 +293,20 @@ public final class ManagerMessages {
   public static final String
       
LOG_SKIP_NON_EXISTENT_REGION_ID_ARG_IN_RECONSTRUCTREGION_REQUEST_TO_DATANODE_ARG_7F76D789
 =
           "跳过 ReconstructRegion 请求中不存在的 Region ID {},目标 DataNode 为 {}。";
+  public static final String 
MESSAGE_DUPLICATE_REGION_ID_ARG_IN_THE_REQUEST_B6FFCCFC =
+      "请求中包含重复的 Region ID %d";
+  public static final String MESSAGE_REGION_IDS_MUST_NOT_BE_EMPTY_B42DAAFD =
+      "Region ID 列表不能为空";
+  public static final String 
MESSAGE_SOURCE_AND_TARGET_DATANODE_IDS_MUST_BE_DIFFERENT_ARG_286D3838 =
+      "源和目标 DataNode ID 不能相同:%d";
+  public static final String 
LOG_SUBMIT_REGION_OPERATION_PROCEDURE_SUCCESSFULLY_ARG_90468B38 =
+      "成功提交 Region 运维 procedure:{}";
+  public static final String MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9 =
+      "Region %d 不存在";
+  public static final String 
MESSAGE_SOURCE_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_2255633C =
+      "源 DataNode %s 不存在于集群中";
+  public static final String 
MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF =
+      "目标 DataNode %s 不存在于集群中";
   public static final String 
MIGRATEREGION_SUBMIT_REGIONMIGRATEPROCEDURE_SUCCESSFULLY_REGION_ORIGIN_DATANODE 
=
       "[MigrateRegion] 成功提交 RegionMigrateProcedure,Region:{},原 DataNode:{},目标 
DataNode:{},新增 Coordinator:{},移除 Coordinator:{}";
   public static final String 
SUBMIT_REGIONMIGRATEPROCEDURE_FAILED_BECAUSE_REGIONGROUP_DOESN_T_EXIST =
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
index dfbb440d9d6..e6efa881da1 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
@@ -25,7 +25,6 @@ import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
 import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
 import org.apache.iotdb.common.rpc.thrift.TDataNodeConfiguration;
 import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
-import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.cluster.NodeStatus;
 import org.apache.iotdb.commons.conf.CommonConfig;
@@ -188,7 +187,6 @@ import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.HashMap;
 import java.util.HashSet;
-import java.util.LinkedHashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
@@ -196,7 +194,6 @@ import java.util.Optional;
 import java.util.Set;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.locks.ReentrantLock;
-import java.util.function.BiFunction;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
 
@@ -963,7 +960,7 @@ public class ProcedureManager {
 
     if (failMessage != null) {
       LOGGER.warn(failMessage);
-      TSStatus failStatus = new 
TSStatus(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode());
+      TSStatus failStatus = new 
TSStatus(TSStatusCode.EXTEND_REGION_ERROR.getStatusCode());
       failStatus.setMessage(failMessage);
       return failStatus;
     }
@@ -973,7 +970,7 @@ public class ProcedureManager {
   private TSStatus checkRemoveRegion(
       TRemoveRegionReq req,
       TConsensusGroupId regionId,
-      @Nullable TDataNodeLocation targetDataNode,
+      TDataNodeLocation targetDataNode,
       TDataNodeLocation coordinator) {
     String failMessage =
         regionOperationCommonCheck(
@@ -991,12 +988,11 @@ public class ProcedureManager {
             .getDataNodeLocationsSize()
         == 1) {
       failMessage = String.format("%s only has 1 replica, it cannot be 
removed", regionId);
-    } else if (targetDataNode != null
-        && configManager
-            .getPartitionManager()
-            .getAllReplicaSets(targetDataNode.getDataNodeId())
-            .stream()
-            .noneMatch(replicaSet -> 
replicaSet.getRegionId().equals(regionId))) {
+    } else if (configManager
+        .getPartitionManager()
+        .getAllReplicaSets(targetDataNode.getDataNodeId())
+        .stream()
+        .noneMatch(replicaSet -> replicaSet.getRegionId().equals(regionId))) {
       failMessage =
           String.format(
               "Target DataNode %s doesn't contain Region %s", 
req.getDataNodeId(), regionId);
@@ -1152,132 +1148,67 @@ public class ProcedureManager {
 
   // end region
 
-  public TSStatus migrateRegion(TMigrateRegionReq migrateRegionReq) {
+  public TSStatus migrateRegion(TMigrateRegionReq req) {
     try (AutoCloseableLock ignoredLock =
         AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
-      // The source and destination DataNodes are fixed for the whole 
statement, so resolve them
-      // once and reuse them for every region.
-      final TDataNodeConfiguration originalDataNodeConfiguration =
-          
configManager.getNodeManager().getRegisteredDataNode(migrateRegionReq.getFromId());
-      final TDataNodeConfiguration destDataNodeConfiguration =
-          
configManager.getNodeManager().getRegisteredDataNode(migrateRegionReq.getToId());
-      if (originalDataNodeConfiguration == null) {
+      final TDataNodeLocation originalDataNode =
+          getRegisteredDataNodeLocationOrNull(req.getFromId());
+      final TDataNodeLocation destDataNode = 
getRegisteredDataNodeLocationOrNull(req.getToId());
+      if (originalDataNode == null) {
         return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode())
             .setMessage(
                 String.format(
-                    "Source DataNode %s does not exist in the cluster",
-                    migrateRegionReq.getFromId()));
+                    ManagerMessages
+                        
.MESSAGE_SOURCE_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_2255633C,
+                    req.getFromId()));
       }
-      if (destDataNodeConfiguration == null) {
+      if (destDataNode == null) {
         return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode())
             .setMessage(
                 String.format(
-                    "Target DataNode %s does not exist in the cluster",
-                    migrateRegionReq.getToId()));
+                    ManagerMessages
+                        
.MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF,
+                    req.getToId()));
+      }
+      if (req.getFromId() == req.getToId()) {
+        return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode())
+            .setMessage(
+                String.format(
+                    ManagerMessages
+                        
.MESSAGE_SOURCE_AND_TARGET_DATANODE_IDS_MUST_BE_DIFFERENT_ARG_286D3838,
+                    req.getFromId()));
       }
-      final TDataNodeLocation originalDataNode = 
originalDataNodeConfiguration.getLocation();
-      final TDataNodeLocation destDataNode = 
destDataNodeConfiguration.getLocation();
-      final RegionMaintainHandler handler = env.getRegionMaintainHandler();
 
-      TSStatus resp = new TSStatus();
-      StringBuilder messageBuilder = new StringBuilder();
-      int total = 0, success = 0;
-      // dedup region ids while preserving the user-specified order
-      for (int theRegionId : new 
LinkedHashSet<>(migrateRegionReq.getRegionIds())) {
-        total++;
-        TSStatus subStatus =
-            migrateOneRegion(
-                migrateRegionReq, theRegionId, originalDataNode, destDataNode, 
handler);
-        if (subStatus.getCode() == 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-          messageBuilder.append("region ").append(theRegionId).append(": 
Successfully submitted\n");
-          success++;
-        } else {
-          messageBuilder
-              .append("region ")
-              .append(theRegionId)
-              .append(": ")
-              .append(subStatus.getMessage())
-              .append('\n');
+      List<TConsensusGroupId> regionIds = new ArrayList<>();
+      TSStatus status =
+          checkRegionIds(req.getRegionIds(), regionIds, 
TSStatusCode.MIGRATE_REGION_ERROR);
+      if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+        return status;
+      }
+      final RegionMaintainHandler handler = env.getRegionMaintainHandler();
+      List<RegionMigrateProcedure> procedures = new ArrayList<>();
+      for (TConsensusGroupId regionId : regionIds) {
+        final TDataNodeLocation coordinator =
+            handler
+                .filterDataNodeWithOtherRegionReplica(
+                    regionId,
+                    destDataNode,
+                    NodeStatus.Running,
+                    NodeStatus.Removing,
+                    NodeStatus.ReadOnly)
+                .orElse(null);
+        status =
+            checkMigrateRegion(
+                req, regionId.getId(), regionId, originalDataNode, 
destDataNode, coordinator);
+        if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+          return status;
         }
-        resp.addToSubStatus(subStatus);
+        procedures.add(
+            new RegionMigrateProcedure(
+                regionId, originalDataNode, destDataNode, coordinator, 
destDataNode));
       }
-
-      messageBuilder.insert(
-          0,
-          String.format(
-              "Total regions: %d, successfully submitted: %d, failed to 
submit: %d\n",
-              total, success, total - success));
-      resp.setCode(
-          total == success
-              ? TSStatusCode.SUCCESS_STATUS.getStatusCode()
-              : TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
-      resp.setMessage(messageBuilder.toString());
-      return resp;
-    }
-  }
-
-  private TSStatus migrateOneRegion(
-      TMigrateRegionReq migrateRegionReq,
-      int theRegionId,
-      TDataNodeLocation originalDataNode,
-      TDataNodeLocation destDataNode,
-      RegionMaintainHandler handler) {
-    TConsensusGroupId regionGroupId;
-    Optional<TConsensusGroupId> optional =
-        
configManager.getPartitionManager().generateTConsensusGroupIdByRegionId(theRegionId);
-    if (optional.isPresent()) {
-      regionGroupId = optional.get();
-    } else {
-      LOGGER.error(ManagerMessages.GET_REGION_GROUP_ID_FAIL);
-      return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode())
-          .setMessage(ManagerMessages.GET_REGION_GROUP_ID_FAIL);
-    }
-
-    // select coordinator for adding peer
-    // (future improvement: choose the DataNode which has the lowest load)
-    final TDataNodeLocation coordinatorForAddPeer =
-        handler
-            .filterDataNodeWithOtherRegionReplica(
-                regionGroupId,
-                destDataNode,
-                NodeStatus.Running,
-                NodeStatus.Removing,
-                NodeStatus.ReadOnly)
-            .orElse(null);
-    // Select coordinator for removing peer
-    // For now, destDataNode temporarily acts as the coordinatorForRemovePeer
-    final TDataNodeLocation coordinatorForRemovePeer = destDataNode;
-
-    TSStatus status =
-        checkMigrateRegion(
-            migrateRegionReq,
-            theRegionId,
-            regionGroupId,
-            originalDataNode,
-            destDataNode,
-            coordinatorForAddPeer);
-    if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-      return status;
+      return submitRegionOperationProcedures(procedures);
     }
-
-    // finally, submit procedure
-    this.executor.submitProcedure(
-        new RegionMigrateProcedure(
-            regionGroupId,
-            originalDataNode,
-            destDataNode,
-            coordinatorForAddPeer,
-            coordinatorForRemovePeer));
-    LOGGER.info(
-        ManagerMessages
-            
.MIGRATEREGION_SUBMIT_REGIONMIGRATEPROCEDURE_SUCCESSFULLY_REGION_ORIGIN_DATANODE,
-        regionGroupId,
-        originalDataNode,
-        destDataNode,
-        coordinatorForAddPeer,
-        coordinatorForRemovePeer);
-
-    return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
   }
 
   /**
@@ -1288,45 +1219,34 @@ public class ProcedureManager {
    * dereference the result blindly.
    */
   private TDataNodeLocation getRegisteredDataNodeLocationOrNull(int 
dataNodeId) {
-    return 
configManager.getNodeManager().getRegisteredDataNode(dataNodeId).getLocation();
+    final TDataNodeConfiguration dataNodeConfiguration =
+        configManager.getNodeManager().getRegisteredDataNode(dataNodeId);
+    return dataNodeConfiguration == null ? null : 
dataNodeConfiguration.getLocation();
   }
 
   public TSStatus reconstructRegion(TReconstructRegionReq req) {
-    RegionMaintainHandler handler = env.getRegionMaintainHandler();
-    final TDataNodeLocation targetDataNode =
-        getRegisteredDataNodeLocationOrNull(req.getDataNodeId());
-    if (targetDataNode == null) {
-      // The target id is not a registered DataNode. Reject here instead of 
pushing a null down into
-      // checkReconstructRegion, which would otherwise throw a 
NullPointerException.
-      return new 
TSStatus(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode())
-          .setMessage(
-              String.format(
-                  "Target DataNode %s does not exist in the cluster", 
req.getDataNodeId()));
-    }
     try (AutoCloseableLock ignoredLock =
         AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
+      final TDataNodeLocation targetDataNode =
+          getRegisteredDataNodeLocationOrNull(req.getDataNodeId());
+      if (targetDataNode == null) {
+        return new 
TSStatus(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode())
+            .setMessage(
+                String.format(
+                    ManagerMessages
+                        
.MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF,
+                    req.getDataNodeId()));
+      }
+
+      List<TConsensusGroupId> regionIds = new ArrayList<>();
+      TSStatus status =
+          checkRegionIds(req.getRegionIds(), regionIds, 
TSStatusCode.RECONSTRUCT_REGION_ERROR);
+      if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+        return status;
+      }
+      final RegionMaintainHandler handler = env.getRegionMaintainHandler();
       List<ReconstructRegionProcedure> procedures = new ArrayList<>();
-      Set<Integer> seenRegionIds = new HashSet<>();
-      for (int x : req.getRegionIds()) {
-        if (!seenRegionIds.add(x)) {
-          LOGGER.info(
-              ManagerMessages
-                  
.LOG_SKIP_DUPLICATE_REGION_ID_ARG_IN_RECONSTRUCTREGION_REQUEST_TO_DATANODE_ARG_ED195F69,
-              x,
-              req.getDataNodeId());
-          continue;
-        }
-        Optional<TConsensusGroupId> regionIdOptional =
-            
configManager.getPartitionManager().findTConsensusGroupIdByRegionId(x);
-        if (!regionIdOptional.isPresent()) {
-          LOGGER.info(
-              ManagerMessages
-                  
.LOG_SKIP_NON_EXISTENT_REGION_ID_ARG_IN_RECONSTRUCTREGION_REQUEST_TO_DATANODE_ARG_7F76D789,
-              x,
-              req.getDataNodeId());
-          continue;
-        }
-        TConsensusGroupId regionId = regionIdOptional.get();
+      for (TConsensusGroupId regionId : regionIds) {
         final TDataNodeLocation coordinator =
             handler
                 .filterDataNodeWithOtherRegionReplica(
@@ -1336,194 +1256,142 @@ public class ProcedureManager {
                     NodeStatus.Removing,
                     NodeStatus.ReadOnly)
                 .orElse(null);
-        TSStatus status = checkReconstructRegion(req, regionId, 
targetDataNode, coordinator);
+        status = checkReconstructRegion(req, regionId, targetDataNode, 
coordinator);
         if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
           return status;
         }
         procedures.add(new ReconstructRegionProcedure(regionId, 
targetDataNode, coordinator));
       }
-      // all checks pass, submit all procedures
-      procedures.forEach(
-          reconstructRegionProcedure -> {
-            this.executor.submitProcedure(reconstructRegionProcedure);
-            LOGGER.info(
-                
ManagerMessages.RECONSTRUCTREGION_SUBMIT_RECONSTRUCTREGIONPROCEDURE_SUCCESSFULLY,
-                reconstructRegionProcedure);
-          });
+      return submitRegionOperationProcedures(procedures);
     }
-    return RpcUtils.SUCCESS_STATUS;
   }
 
   public TSStatus extendRegions(TExtendRegionReq req) {
-    return processExtendOrRemoveRegions(
-        req.getRegionId(), req, this::extendOneRegion, 
TSStatusCode.EXTEND_REGION_ERROR);
-  }
-
-  public TSStatus removeRegions(TRemoveRegionReq req) {
-    return processExtendOrRemoveRegions(
-        req.getRegionId(), req, this::removeOneRegion, 
TSStatusCode.REMOVE_REGION_PEER_ERROR);
-  }
-
-  private <R> TSStatus processExtendOrRemoveRegions(
-      Iterable<Integer> regionIds,
-      R req,
-      BiFunction<Integer, R, TSStatus> regionAction,
-      TSStatusCode errorCode) {
-    TSStatus resp = new TSStatus();
-    StringBuilder messageBuilder = new StringBuilder();
-
-    int total = 0, success = 0;
-    for (int regionId : regionIds) {
-      total++;
-      TSStatus subStatus = regionAction.apply(regionId, req);
-      if (subStatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
-        messageBuilder.append("region ").append(regionId).append(": 
Successfully submitted\n");
-        success++;
-      } else {
-        messageBuilder
-            .append("region ")
-            .append(regionId)
-            .append(": ")
-            .append(subStatus.getMessage())
-            .append('\n');
-      }
-      resp.addToSubStatus(subStatus);
-    }
-
-    messageBuilder.insert(
-        0,
-        String.format(
-            "Total regions: %d, successfully submitted: %d, failed to submit: 
%d\n",
-            total, success, total - success));
-
-    resp.setCode(
-        total == success ? TSStatusCode.SUCCESS_STATUS.getStatusCode() : 
errorCode.getStatusCode());
-    resp.setMessage(messageBuilder.toString());
-    return resp;
-  }
-
-  private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) {
     try (AutoCloseableLock ignoredLock =
         AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
-      TConsensusGroupId regionId;
-      Optional<TConsensusGroupId> optional =
-          
configManager.getPartitionManager().generateTConsensusGroupIdByRegionId(theRegionId);
-      if (optional.isPresent()) {
-        regionId = optional.get();
-      } else {
-        LOGGER.error(ManagerMessages.GET_REGION_GROUP_ID_FAIL);
-        return new TSStatus(TSStatusCode.EXTEND_REGION_ERROR.getStatusCode())
-            .setMessage(ManagerMessages.GET_REGION_GROUP_ID_FAIL);
-      }
-
-      // find target dn
       final TDataNodeLocation targetDataNode =
           getRegisteredDataNodeLocationOrNull(req.getDataNodeId());
       if (targetDataNode == null) {
-        // The target id is not a registered DataNode. Reject here instead of 
pushing a null down
-        // into checkExtendRegion, which would otherwise throw a 
NullPointerException.
         return new TSStatus(TSStatusCode.EXTEND_REGION_ERROR.getStatusCode())
             .setMessage(
                 String.format(
-                    "Target DataNode %s does not exist in the cluster", 
req.getDataNodeId()));
-      }
-      // select coordinator for adding peer
-      RegionMaintainHandler handler = env.getRegionMaintainHandler();
-      // TODO: choose the DataNode which has lowest load
-      final TDataNodeLocation coordinator =
-          handler
-              .filterDataNodeWithOtherRegionReplica(
-                  regionId,
-                  targetDataNode,
-                  NodeStatus.Running,
-                  NodeStatus.Removing,
-                  NodeStatus.ReadOnly)
-              .orElse(null);
-      // do the check
-      TSStatus status = checkExtendRegion(req, regionId, targetDataNode, 
coordinator);
+                    ManagerMessages
+                        
.MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF,
+                    req.getDataNodeId()));
+      }
+
+      List<TConsensusGroupId> regionIds = new ArrayList<>();
+      TSStatus status =
+          checkRegionIds(req.getRegionId(), regionIds, 
TSStatusCode.EXTEND_REGION_ERROR);
       if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
         return status;
       }
-      // submit procedure
-      AddRegionPeerProcedure procedure =
-          new AddRegionPeerProcedure(regionId, coordinator, targetDataNode);
-      this.executor.submitProcedure(procedure);
-      LOGGER.info(
-          
ManagerMessages.EXTENDREGION_SUBMIT_ADDREGIONPEERPROCEDURE_SUCCESSFULLY, 
procedure);
-
-      return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+      final RegionMaintainHandler handler = env.getRegionMaintainHandler();
+      List<AddRegionPeerProcedure> procedures = new ArrayList<>();
+      for (TConsensusGroupId regionId : regionIds) {
+        final TDataNodeLocation coordinator =
+            handler
+                .filterDataNodeWithOtherRegionReplica(
+                    regionId,
+                    targetDataNode,
+                    NodeStatus.Running,
+                    NodeStatus.Removing,
+                    NodeStatus.ReadOnly)
+                .orElse(null);
+        status = checkExtendRegion(req, regionId, targetDataNode, coordinator);
+        if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+          return status;
+        }
+        procedures.add(new AddRegionPeerProcedure(regionId, coordinator, 
targetDataNode));
+      }
+      return submitRegionOperationProcedures(procedures);
     }
   }
 
-  private TSStatus removeOneRegion(int theRegionId, TRemoveRegionReq req) {
+  public TSStatus removeRegions(TRemoveRegionReq req) {
     try (AutoCloseableLock ignoredLock =
         AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
-      TConsensusGroupId regionId;
-      Optional<TConsensusGroupId> optional =
-          
configManager.getPartitionManager().generateTConsensusGroupIdByRegionId(theRegionId);
-      if (optional.isPresent()) {
-        regionId = optional.get();
-      } else {
-        LOGGER.error(ManagerMessages.GET_REGION_GROUP_ID_FAIL);
+      final TDataNodeLocation targetDataNode =
+          getRegisteredDataNodeLocationOrNull(req.getDataNodeId());
+      if (targetDataNode == null) {
         return new 
TSStatus(TSStatusCode.REMOVE_REGION_PEER_ERROR.getStatusCode())
-            .setMessage(ManagerMessages.GET_REGION_GROUP_ID_FAIL);
+            .setMessage(
+                String.format(
+                    ManagerMessages
+                        
.MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF,
+                    req.getDataNodeId()));
       }
 
-      // find target dn
-      final TDataNodeLocation targetDataNode =
-          
configManager.getNodeManager().getRegisteredDataNode(req.getDataNodeId()).getLocation();
-
-      // select coordinator for removing peer
-      RegionMaintainHandler handler = env.getRegionMaintainHandler();
-      final TDataNodeLocation coordinator =
-          handler
-              .filterDataNodeWithOtherRegionReplica(
-                  regionId,
-                  targetDataNode,
-                  NodeStatus.Running,
-                  NodeStatus.Removing,
-                  NodeStatus.ReadOnly)
-              .orElse(null);
-
-      // do the check
-      TSStatus status = checkRemoveRegion(req, regionId, targetDataNode, 
coordinator);
+      List<TConsensusGroupId> regionIds = new ArrayList<>();
+      TSStatus status =
+          checkRegionIds(req.getRegionId(), regionIds, 
TSStatusCode.REMOVE_REGION_PEER_ERROR);
       if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
         return status;
       }
+      final RegionMaintainHandler handler = env.getRegionMaintainHandler();
+      List<RemoveRegionPeerProcedure> procedures = new ArrayList<>();
+      for (TConsensusGroupId regionId : regionIds) {
+        final TDataNodeLocation coordinator =
+            handler
+                .filterDataNodeWithOtherRegionReplica(
+                    regionId,
+                    targetDataNode,
+                    NodeStatus.Running,
+                    NodeStatus.Removing,
+                    NodeStatus.ReadOnly)
+                .orElse(null);
+        status = checkRemoveRegion(req, regionId, targetDataNode, coordinator);
+        if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+          return status;
+        }
+        procedures.add(new RemoveRegionPeerProcedure(regionId, coordinator, 
targetDataNode));
+      }
+      return submitRegionOperationProcedures(procedures);
+    }
+  }
 
-      // SPECIAL CASE
-      if (targetDataNode == null) {
-        // If targetDataNode is null, it means the target DataNode does not 
exist in the
-        // NodeManager.
-        // In this case, simply clean up the partition table once and do 
nothing else.
-        LOGGER.warn(
-            
ManagerMessages.REMOVE_REGION_TARGET_DATANODE_NOT_FOUND_WILL_SIMPLY_CLEAN_UP,
-            req.getDataNodeId(),
-            req.getRegionId());
-        this.executor
-            .getEnvironment()
-            .getRegionMaintainHandler()
-            .removeRegionLocation(
-                regionId, buildFakeDataNodeLocation(req.getDataNodeId(), 
"FakeIpForRemoveRegion"));
-        return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+  /** Resolve every region ID before preparing or submitting any region 
operation. */
+  private TSStatus checkRegionIds(
+      List<Integer> requestedRegionIds, List<TConsensusGroupId> regionIds, 
TSStatusCode errorCode) {
+    if (requestedRegionIds == null || requestedRegionIds.isEmpty()) {
+      return new TSStatus(errorCode.getStatusCode())
+          
.setMessage(ManagerMessages.MESSAGE_REGION_IDS_MUST_NOT_BE_EMPTY_B42DAAFD);
+    }
+    Set<Integer> seenRegionIds = new HashSet<>();
+    for (int regionId : requestedRegionIds) {
+      if (!seenRegionIds.add(regionId)) {
+        return new TSStatus(errorCode.getStatusCode())
+            .setMessage(
+                String.format(
+                    
ManagerMessages.MESSAGE_DUPLICATE_REGION_ID_ARG_IN_THE_REQUEST_B6FFCCFC,
+                    regionId));
       }
+      Optional<TConsensusGroupId> resolvedRegionId =
+          
configManager.getPartitionManager().findTConsensusGroupIdByRegionId(regionId);
+      if (!resolvedRegionId.isPresent()) {
+        return new TSStatus(errorCode.getStatusCode())
+            .setMessage(
+                String.format(
+                    
ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, regionId));
+      }
+      regionIds.add(resolvedRegionId.get());
+    }
+    return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+  }
 
-      // submit procedure
-      RemoveRegionPeerProcedure procedure =
-          new RemoveRegionPeerProcedure(regionId, coordinator, targetDataNode);
-      this.executor.submitProcedure(procedure);
+  /**
+   * Called with the submission lock held, only after every region has passed 
validation. This
+   * prevents validation failures from leaving a partially submitted request.
+   */
+  private TSStatus submitRegionOperationProcedures(
+      List<? extends RegionOperationProcedure<?>> procedures) {
+    for (RegionOperationProcedure<?> procedure : procedures) {
+      executor.submitProcedure(procedure);
       LOGGER.info(
-          
ManagerMessages.REMOVEREGIONPEER_SUBMIT_REMOVEREGIONPEERPROCEDURE_SUCCESSFULLY,
+          
ManagerMessages.LOG_SUBMIT_REGION_OPERATION_PROCEDURE_SUCCESSFULLY_ARG_90468B38,
           procedure);
-
-      return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     }
-  }
-
-  private static TDataNodeLocation buildFakeDataNodeLocation(int dataNodeId, 
String message) {
-    TEndPoint fakeEndPoint = new TEndPoint(message, -1);
-    return new TDataNodeLocation(
-        dataNodeId, fakeEndPoint, fakeEndPoint, fakeEndPoint, fakeEndPoint, 
fakeEndPoint);
+    return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
   }
 
   // endregion
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerReconstructRegionTest.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerReconstructRegionTest.java
deleted file mode 100644
index 4e2b58daa63..00000000000
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerReconstructRegionTest.java
+++ /dev/null
@@ -1,190 +0,0 @@
-/*
- * 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.confignode.manager;
-
-import org.apache.iotdb.common.rpc.thrift.Model;
-import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
-import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
-import org.apache.iotdb.common.rpc.thrift.TDataNodeConfiguration;
-import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
-import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
-import org.apache.iotdb.common.rpc.thrift.TSStatus;
-import org.apache.iotdb.commons.cluster.NodeStatus;
-import org.apache.iotdb.confignode.manager.node.NodeManager;
-import org.apache.iotdb.confignode.manager.partition.PartitionManager;
-import org.apache.iotdb.confignode.persistence.ProcedureInfo;
-import org.apache.iotdb.confignode.procedure.Procedure;
-import org.apache.iotdb.confignode.procedure.ProcedureExecutor;
-import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
-import org.apache.iotdb.confignode.procedure.env.RegionMaintainHandler;
-import 
org.apache.iotdb.confignode.procedure.impl.region.ReconstructRegionProcedure;
-import org.apache.iotdb.confignode.rpc.thrift.TReconstructRegionReq;
-import org.apache.iotdb.confignode.rpc.thrift.TRemoveRegionReq;
-import org.apache.iotdb.rpc.TSStatusCode;
-
-import org.junit.Before;
-import org.junit.Test;
-import org.mockito.ArgumentCaptor;
-
-import java.lang.reflect.Field;
-import java.util.Arrays;
-import java.util.Collections;
-import java.util.HashMap;
-import java.util.Map;
-import java.util.Optional;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.locks.ReentrantLock;
-
-import static org.junit.Assert.assertEquals;
-import static org.junit.Assert.assertTrue;
-import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.ArgumentMatchers.eq;
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.times;
-import static org.mockito.Mockito.verify;
-import static org.mockito.Mockito.when;
-
-public class ProcedureManagerReconstructRegionTest {
-
-  private final TConsensusGroupId firstRegion =
-      new TConsensusGroupId(TConsensusGroupType.DataRegion, 12);
-  private final TConsensusGroupId secondRegion =
-      new TConsensusGroupId(TConsensusGroupType.DataRegion, 14);
-  private final TDataNodeLocation target = new 
TDataNodeLocation().setDataNodeId(7);
-  private final TDataNodeLocation coordinator = new 
TDataNodeLocation().setDataNodeId(8);
-
-  private ProcedureManager manager;
-  private ProcedureExecutor<ConfigNodeProcedureEnv> executor;
-  private NodeManager nodeManager;
-  private PartitionManager partitionManager;
-  private final ConcurrentHashMap<Long, Procedure<ConfigNodeProcedureEnv>> 
procedures =
-      new ConcurrentHashMap<>();
-
-  @Before
-  public void setUp() throws Exception {
-    ConfigManager configManager = mock(ConfigManager.class);
-    nodeManager = mock(NodeManager.class);
-    partitionManager = mock(PartitionManager.class);
-    ConfigNodeProcedureEnv env = mock(ConfigNodeProcedureEnv.class);
-    RegionMaintainHandler handler = mock(RegionMaintainHandler.class);
-    executor = mock(ProcedureExecutor.class);
-
-    when(configManager.getNodeManager()).thenReturn(nodeManager);
-    when(configManager.getPartitionManager()).thenReturn(partitionManager);
-    when(nodeManager.getRegisteredDataNode(target.getDataNodeId()))
-        .thenReturn(new TDataNodeConfiguration().setLocation(target));
-    when(nodeManager.filterDataNodeThroughStatus(NodeStatus.Running))
-        .thenReturn(Collections.singletonList(new 
TDataNodeConfiguration().setLocation(target)));
-    when(nodeManager.filterDataNodeThroughStatus(NodeStatus.Running, 
NodeStatus.ReadOnly))
-        .thenReturn(Collections.singletonList(new 
TDataNodeConfiguration().setLocation(target)));
-    
when(partitionManager.findTConsensusGroupIdByRegionId(12)).thenReturn(Optional.of(firstRegion));
-    when(partitionManager.findTConsensusGroupIdByRegionId(14))
-        .thenReturn(Optional.of(secondRegion));
-    
when(partitionManager.findTConsensusGroupIdByRegionId(99)).thenReturn(Optional.empty());
-    when(partitionManager.generateTConsensusGroupIdByRegionId(12))
-        .thenReturn(Optional.of(firstRegion));
-    when(partitionManager.generateTConsensusGroupIdByRegionId(14))
-        .thenReturn(Optional.of(secondRegion));
-    
when(partitionManager.getRegionDatabase(any(TConsensusGroupId.class))).thenReturn("root.sg");
-
-    Map<TConsensusGroupId, TRegionReplicaSet> replicaSets = new HashMap<>();
-    replicaSets.put(
-        firstRegion, new TRegionReplicaSet(firstRegion, Arrays.asList(target, 
coordinator)));
-    replicaSets.put(
-        secondRegion, new TRegionReplicaSet(secondRegion, 
Arrays.asList(target, coordinator)));
-    when(partitionManager.getAllReplicaSetsMap(TConsensusGroupType.DataRegion))
-        .thenReturn(replicaSets);
-    when(partitionManager.getAllReplicaSets(target.getDataNodeId()))
-        .thenReturn(Arrays.asList(replicaSets.get(firstRegion), 
replicaSets.get(secondRegion)));
-
-    when(env.getSubmitRegionMigrateLock()).thenReturn(new ReentrantLock());
-    when(env.getRegionMaintainHandler()).thenReturn(handler);
-    when(handler.filterDataNodeWithOtherRegionReplica(
-            any(TConsensusGroupId.class),
-            eq(target),
-            eq(NodeStatus.Running),
-            eq(NodeStatus.Removing),
-            eq(NodeStatus.ReadOnly)))
-        .thenReturn(Optional.of(coordinator));
-    when(executor.getProcedures()).thenReturn(procedures);
-
-    manager = new ProcedureManager(configManager, mock(ProcedureInfo.class));
-    Field envField = ProcedureManager.class.getDeclaredField("env");
-    envField.setAccessible(true);
-    envField.set(manager, env);
-    Field executorField = ProcedureManager.class.getDeclaredField("executor");
-    executorField.setAccessible(true);
-    executorField.set(manager, executor);
-  }
-
-  @Test
-  public void testDuplicateAndNonExistentRegionIdsAreSkippedInInputOrder() {
-    TReconstructRegionReq request =
-        new TReconstructRegionReq(Arrays.asList(12, 99, 14, 12, 99, 14), 7, 
Model.TREE);
-
-    assertEquals(
-        TSStatusCode.SUCCESS_STATUS.getStatusCode(), 
manager.reconstructRegion(request).getCode());
-
-    ArgumentCaptor<ReconstructRegionProcedure> captor =
-        ArgumentCaptor.forClass(ReconstructRegionProcedure.class);
-    verify(executor, times(2)).submitProcedure(captor.capture());
-    assertEquals(firstRegion, captor.getAllValues().get(0).getRegionId());
-    assertEquals(secondRegion, captor.getAllValues().get(1).getRegionId());
-    verify(partitionManager, times(1)).findTConsensusGroupIdByRegionId(12);
-    verify(partitionManager, times(1)).findTConsensusGroupIdByRegionId(14);
-    verify(partitionManager, times(1)).findTConsensusGroupIdByRegionId(99);
-  }
-
-  @Test
-  public void 
testRequestWithNoUsableRegionIdsSucceedsWithoutSubmittingProcedure() {
-    TReconstructRegionReq request = new 
TReconstructRegionReq(Arrays.asList(99, 99), 7, Model.TREE);
-
-    assertEquals(
-        TSStatusCode.SUCCESS_STATUS.getStatusCode(), 
manager.reconstructRegion(request).getCode());
-    verify(executor, times(0)).submitProcedure(any());
-  }
-
-  @Test
-  public void testAnotherRequestCannotReconstructRegionWithActiveProcedure() {
-    TReconstructRegionReq request =
-        new TReconstructRegionReq(Collections.singletonList(12), 7, 
Model.TREE);
-    ReconstructRegionProcedure activeProcedure =
-        new ReconstructRegionProcedure(firstRegion, target, coordinator);
-    procedures.put(1L, activeProcedure);
-
-    TSStatus status = manager.reconstructRegion(request);
-    assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(), 
status.getCode());
-    assertTrue(status.getMessage().contains("in progress"));
-    verify(executor, times(0)).submitProcedure(any());
-  }
-
-  @Test
-  public void testRemoveRegionAllowsReadOnlyTargetDataNode() {
-    procedures.clear();
-    when(nodeManager.filterDataNodeThroughStatus(NodeStatus.Running))
-        .thenReturn(
-            Collections.singletonList(new 
TDataNodeConfiguration().setLocation(coordinator)));
-    TRemoveRegionReq request = new 
TRemoveRegionReq(Collections.singletonList(12), 7, Model.TREE);
-
-    assertEquals(
-        TSStatusCode.SUCCESS_STATUS.getStatusCode(), 
manager.removeRegions(request).getCode());
-    verify(executor, times(1)).submitProcedure(any());
-  }
-}
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerRegionOperationTest.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerRegionOperationTest.java
new file mode 100644
index 00000000000..d0fb6efbb0a
--- /dev/null
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerRegionOperationTest.java
@@ -0,0 +1,401 @@
+/*
+ * 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.confignode.manager;
+
+import org.apache.iotdb.common.rpc.thrift.Model;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
+import org.apache.iotdb.common.rpc.thrift.TDataNodeConfiguration;
+import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.cluster.NodeStatus;
+import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
+import org.apache.iotdb.confignode.i18n.ManagerMessages;
+import org.apache.iotdb.confignode.manager.node.NodeManager;
+import org.apache.iotdb.confignode.manager.partition.PartitionManager;
+import org.apache.iotdb.confignode.persistence.ProcedureInfo;
+import org.apache.iotdb.confignode.procedure.Procedure;
+import org.apache.iotdb.confignode.procedure.ProcedureExecutor;
+import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
+import org.apache.iotdb.confignode.procedure.env.RegionMaintainHandler;
+import 
org.apache.iotdb.confignode.procedure.impl.region.AddRegionPeerProcedure;
+import 
org.apache.iotdb.confignode.procedure.impl.region.ReconstructRegionProcedure;
+import 
org.apache.iotdb.confignode.procedure.impl.region.RegionMigrateProcedure;
+import 
org.apache.iotdb.confignode.procedure.impl.region.RegionOperationProcedure;
+import 
org.apache.iotdb.confignode.procedure.impl.region.RemoveRegionPeerProcedure;
+import org.apache.iotdb.confignode.rpc.thrift.TExtendRegionReq;
+import org.apache.iotdb.confignode.rpc.thrift.TMigrateRegionReq;
+import org.apache.iotdb.confignode.rpc.thrift.TReconstructRegionReq;
+import org.apache.iotdb.confignode.rpc.thrift.TRemoveRegionReq;
+import org.apache.iotdb.consensus.ConsensusFactory;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.Parameterized;
+import org.mockito.ArgumentCaptor;
+
+import java.lang.reflect.Field;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.ReentrantLock;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertTrue;
+import static org.junit.Assume.assumeTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+@RunWith(Parameterized.class)
+public class ProcedureManagerRegionOperationTest {
+
+  private enum Operation {
+    MIGRATE(TSStatusCode.MIGRATE_REGION_ERROR, RegionMigrateProcedure.class),
+    RECONSTRUCT(TSStatusCode.RECONSTRUCT_REGION_ERROR, 
ReconstructRegionProcedure.class),
+    EXTEND(TSStatusCode.EXTEND_REGION_ERROR, AddRegionPeerProcedure.class),
+    REMOVE(TSStatusCode.REMOVE_REGION_PEER_ERROR, 
RemoveRegionPeerProcedure.class);
+
+    private final TSStatusCode errorCode;
+    private final Class<?> procedureClass;
+
+    Operation(TSStatusCode errorCode, Class<?> procedureClass) {
+      this.errorCode = errorCode;
+      this.procedureClass = procedureClass;
+    }
+  }
+
+  @Parameterized.Parameters(name = "{0}-{1}-{2}")
+  public static Iterable<Object[]> parameters() {
+    List<Object[]> parameters = new ArrayList<>();
+    for (Operation operation : Operation.values()) {
+      for (Model model : Arrays.asList(Model.TREE, Model.TABLE)) {
+        for (TConsensusGroupType type :
+            Arrays.asList(TConsensusGroupType.DataRegion, 
TConsensusGroupType.SchemaRegion)) {
+          parameters.add(new Object[] {operation, model, type});
+        }
+      }
+    }
+    return parameters;
+  }
+
+  private final Operation operation;
+  private final Model model;
+  private final TConsensusGroupId firstRegion;
+  private final TConsensusGroupId secondRegion;
+  private final TDataNodeLocation source = new 
TDataNodeLocation().setDataNodeId(6);
+  private final TDataNodeLocation target = new 
TDataNodeLocation().setDataNodeId(7);
+  private final TDataNodeLocation coordinator = new 
TDataNodeLocation().setDataNodeId(8);
+  private final ConcurrentHashMap<Long, Procedure<ConfigNodeProcedureEnv>> 
procedures =
+      new ConcurrentHashMap<>();
+
+  private ProcedureManager manager;
+  private ProcedureExecutor<ConfigNodeProcedureEnv> executor;
+  private NodeManager nodeManager;
+  private PartitionManager partitionManager;
+  private RegionMaintainHandler handler;
+  private ReentrantLock submissionLock;
+  private String originalDataConsensus;
+  private String originalSchemaConsensus;
+
+  public ProcedureManagerRegionOperationTest(
+      Operation operation, Model model, TConsensusGroupType type) {
+    this.operation = operation;
+    this.model = model;
+    firstRegion = new TConsensusGroupId(type, 12);
+    secondRegion = new TConsensusGroupId(type, 14);
+  }
+
+  @Before
+  @SuppressWarnings("unchecked")
+  public void setUp() throws Exception {
+    originalDataConsensus =
+        
ConfigNodeDescriptor.getInstance().getConf().getDataRegionConsensusProtocolClass();
+    originalSchemaConsensus =
+        
ConfigNodeDescriptor.getInstance().getConf().getSchemaRegionConsensusProtocolClass();
+    ConfigNodeDescriptor.getInstance()
+        .getConf()
+        .setDataRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS);
+    ConfigNodeDescriptor.getInstance()
+        .getConf()
+        
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS);
+
+    ConfigManager configManager = mock(ConfigManager.class);
+    nodeManager = mock(NodeManager.class);
+    partitionManager = mock(PartitionManager.class);
+    ConfigNodeProcedureEnv env = mock(ConfigNodeProcedureEnv.class);
+    handler = mock(RegionMaintainHandler.class);
+    executor = mock(ProcedureExecutor.class);
+    submissionLock = new ReentrantLock();
+    when(configManager.getNodeManager()).thenReturn(nodeManager);
+    when(configManager.getPartitionManager()).thenReturn(partitionManager);
+
+    // NodeManager returns an empty configuration for IDs that are not 
registered DataNodes.
+    when(nodeManager.getRegisteredDataNode(anyInt())).thenReturn(new 
TDataNodeConfiguration());
+    when(nodeManager.getRegisteredDataNode(source.getDataNodeId()))
+        .thenReturn(new TDataNodeConfiguration().setLocation(source));
+    when(nodeManager.getRegisteredDataNode(target.getDataNodeId()))
+        .thenReturn(new TDataNodeConfiguration().setLocation(target));
+    when(nodeManager.filterDataNodeThroughStatus(NodeStatus.Running))
+        .thenReturn(Collections.singletonList(new 
TDataNodeConfiguration().setLocation(target)));
+    when(nodeManager.filterDataNodeThroughStatus(NodeStatus.Running, 
NodeStatus.ReadOnly))
+        .thenReturn(Collections.singletonList(new 
TDataNodeConfiguration().setLocation(target)));
+
+    
when(partitionManager.findTConsensusGroupIdByRegionId(anyInt())).thenReturn(Optional.empty());
+    
when(partitionManager.findTConsensusGroupIdByRegionId(12)).thenReturn(Optional.of(firstRegion));
+    when(partitionManager.findTConsensusGroupIdByRegionId(14))
+        .thenReturn(Optional.of(secondRegion));
+    when(partitionManager.getRegionDatabase(any(TConsensusGroupId.class)))
+        .thenReturn(model == Model.TREE ? "root.sg" : "db");
+
+    TDataNodeLocation replica =
+        operation == Operation.MIGRATE || operation == Operation.EXTEND ? 
source : target;
+    Map<TConsensusGroupId, TRegionReplicaSet> replicaSets = new HashMap<>();
+    replicaSets.put(
+        firstRegion, new TRegionReplicaSet(firstRegion, Arrays.asList(replica, 
coordinator)));
+    replicaSets.put(
+        secondRegion, new TRegionReplicaSet(secondRegion, 
Arrays.asList(replica, coordinator)));
+    
when(partitionManager.getAllReplicaSetsMap(firstRegion.getType())).thenReturn(replicaSets);
+    when(partitionManager.getAllReplicaSets(replica.getDataNodeId()))
+        .thenReturn(Arrays.asList(replicaSets.get(firstRegion), 
replicaSets.get(secondRegion)));
+
+    when(env.getSubmitRegionMigrateLock()).thenReturn(submissionLock);
+    when(env.getRegionMaintainHandler()).thenReturn(handler);
+    when(handler.filterDataNodeWithOtherRegionReplica(
+            any(TConsensusGroupId.class),
+            eq(target),
+            eq(NodeStatus.Running),
+            eq(NodeStatus.Removing),
+            eq(NodeStatus.ReadOnly)))
+        .thenReturn(Optional.of(coordinator));
+    when(executor.getProcedures()).thenReturn(procedures);
+
+    manager = new ProcedureManager(configManager, mock(ProcedureInfo.class));
+    Field envField = ProcedureManager.class.getDeclaredField("env");
+    envField.setAccessible(true);
+    envField.set(manager, env);
+    Field executorField = ProcedureManager.class.getDeclaredField("executor");
+    executorField.setAccessible(true);
+    executorField.set(manager, executor);
+  }
+
+  @After
+  public void tearDown() {
+    ConfigNodeDescriptor.getInstance()
+        .getConf()
+        .setDataRegionConsensusProtocolClass(originalDataConsensus);
+    ConfigNodeDescriptor.getInstance()
+        .getConf()
+        .setSchemaRegionConsensusProtocolClass(originalSchemaConsensus);
+    assertFalse(submissionLock.isLocked());
+  }
+
+  private TSStatus execute(List<Integer> regionIds, int fromId, int toId) {
+    switch (operation) {
+      case MIGRATE:
+        return manager.migrateRegion(new TMigrateRegionReq(regionIds, fromId, 
toId, model));
+      case RECONSTRUCT:
+        return manager.reconstructRegion(new TReconstructRegionReq(regionIds, 
toId, model));
+      case EXTEND:
+        return manager.extendRegions(new TExtendRegionReq(regionIds, toId, 
model));
+      case REMOVE:
+        return manager.removeRegions(new TRemoveRegionReq(regionIds, toId, 
model));
+      default:
+        throw new AssertionError(operation);
+    }
+  }
+
+  private TSStatus execute(List<Integer> regionIds) {
+    return execute(regionIds, source.getDataNodeId(), target.getDataNodeId());
+  }
+
+  private void assertRejected(TSStatus status) {
+    assertEquals(operation.errorCode.getStatusCode(), status.getCode());
+    assertNotNull(status.getMessage());
+    assertFalse(status.getMessage().isEmpty());
+    verify(executor, never()).submitProcedure(any());
+    verify(handler, never()).removeRegionLocation(any(), any());
+  }
+
+  @Test
+  public void testDuplicateRegionIdRejectsWholeRequest() {
+    TSStatus status = execute(Arrays.asList(12, 14, 12));
+    assertRejected(status);
+    assertEquals(
+        
String.format(ManagerMessages.MESSAGE_DUPLICATE_REGION_ID_ARG_IN_THE_REQUEST_B6FFCCFC,
 12),
+        status.getMessage());
+  }
+
+  @Test
+  public void testMissingLastRegionRejectsWholeRequest() {
+    TSStatus status = execute(Arrays.asList(12, 14, 99));
+    assertRejected(status);
+    assertEquals(
+        
String.format(ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, 99),
+        status.getMessage());
+  }
+
+  @Test
+  public void testMissingFirstRegionRejectsWholeRequest() {
+    assertRejected(execute(Arrays.asList(99, 12, 14)));
+  }
+
+  @Test
+  public void testAllMissingRegionsRejectWholeRequest() {
+    assertRejected(execute(Arrays.asList(98, 99)));
+  }
+
+  @Test
+  public void testEmptyRegionIdsRejectWholeRequest() {
+    TSStatus status = execute(Collections.emptyList());
+    assertRejected(status);
+    assertEquals(
+        ManagerMessages.MESSAGE_REGION_IDS_MUST_NOT_BE_EMPTY_B42DAAFD, 
status.getMessage());
+  }
+
+  @Test
+  public void testMissingTargetDataNodeRejectsWholeRequest() {
+    TSStatus status = execute(Arrays.asList(12, 14), source.getDataNodeId(), 
99);
+    assertRejected(status);
+    assertEquals(
+        String.format(
+            
ManagerMessages.MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF,
 99),
+        status.getMessage());
+  }
+
+  @Test
+  public void testNullTargetDataNodeConfigurationRejectsWholeRequest() {
+    when(nodeManager.getRegisteredDataNode(99)).thenReturn(null);
+    assertRejected(execute(Arrays.asList(12, 14), source.getDataNodeId(), 99));
+  }
+
+  @Test
+  public void testMissingMigrationSourceRejectsWholeRequest() {
+    assumeTrue(operation == Operation.MIGRATE);
+    TSStatus status = execute(Arrays.asList(12, 14), 99, 
target.getDataNodeId());
+    assertRejected(status);
+    assertEquals(
+        String.format(
+            
ManagerMessages.MESSAGE_SOURCE_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_2255633C,
 99),
+        status.getMessage());
+  }
+
+  @Test
+  public void testIdenticalMigrationDataNodeIdsRejectWholeRequest() {
+    assumeTrue(operation == Operation.MIGRATE);
+    TSStatus status =
+        execute(Arrays.asList(12, 14), target.getDataNodeId(), 
target.getDataNodeId());
+    assertRejected(status);
+    assertEquals(
+        String.format(
+            
ManagerMessages.MESSAGE_SOURCE_AND_TARGET_DATANODE_IDS_MUST_BE_DIFFERENT_ARG_286D3838,
+            target.getDataNodeId()),
+        status.getMessage());
+  }
+
+  @Test
+  public void testConflictOnLastRegionRejectsWholeRequest() {
+    procedures.put(1L, new ReconstructRegionProcedure(secondRegion, target, 
coordinator));
+    assertRejected(execute(Arrays.asList(12, 14)));
+  }
+
+  @Test
+  public void testInvalidReplicaPlacementOnLastRegionRejectsWholeRequest() {
+    Map<TConsensusGroupId, TRegionReplicaSet> replicaSets =
+        partitionManager.getAllReplicaSetsMap(firstRegion.getType());
+    if (operation == Operation.EXTEND) {
+      when(partitionManager.getAllReplicaSets(target.getDataNodeId()))
+          
.thenReturn(Collections.singletonList(replicaSets.get(secondRegion)));
+    } else {
+      int dataNodeId =
+          operation == Operation.MIGRATE ? source.getDataNodeId() : 
target.getDataNodeId();
+      when(partitionManager.getAllReplicaSets(dataNodeId))
+          .thenReturn(Collections.singletonList(replicaSets.get(firstRegion)));
+    }
+    assertRejected(execute(Arrays.asList(12, 14)));
+  }
+
+  @Test
+  public void testMissingCoordinatorOnLastRegionRejectsWholeRequest() {
+    when(handler.filterDataNodeWithOtherRegionReplica(
+            eq(secondRegion),
+            eq(target),
+            eq(NodeStatus.Running),
+            eq(NodeStatus.Removing),
+            eq(NodeStatus.ReadOnly)))
+        .thenReturn(Optional.empty());
+    assertRejected(execute(Arrays.asList(12, 14)));
+  }
+
+  @Test
+  public void testValidRegionsAreSubmittedInInputOrderAfterAllChecks() {
+    doAnswer(
+            invocation -> {
+              assertTrue(submissionLock.isHeldByCurrentThread());
+              verify(partitionManager).findTConsensusGroupIdByRegionId(14);
+              verify(partitionManager).findTConsensusGroupIdByRegionId(12);
+              verify(partitionManager).getRegionDatabase(secondRegion);
+              verify(partitionManager).getRegionDatabase(firstRegion);
+              return 1L;
+            })
+        .when(executor)
+        .submitProcedure(any());
+
+    assertEquals(
+        TSStatusCode.SUCCESS_STATUS.getStatusCode(), execute(Arrays.asList(14, 
12)).getCode());
+    ArgumentCaptor<Procedure> captor = 
ArgumentCaptor.forClass(Procedure.class);
+    verify(executor, times(2)).submitProcedure(captor.capture());
+    assertEquals(operation.procedureClass, 
captor.getAllValues().get(0).getClass());
+    assertEquals(operation.procedureClass, 
captor.getAllValues().get(1).getClass());
+    assertEquals(
+        secondRegion, ((RegionOperationProcedure<?>) 
captor.getAllValues().get(0)).getRegionId());
+    assertEquals(
+        firstRegion, ((RegionOperationProcedure<?>) 
captor.getAllValues().get(1)).getRegionId());
+  }
+
+  @Test
+  public void testRemoveRegionAllowsReadOnlyTargetDataNode() {
+    assumeTrue(operation == Operation.REMOVE);
+    when(nodeManager.filterDataNodeThroughStatus(NodeStatus.Running))
+        .thenReturn(
+            Collections.singletonList(new 
TDataNodeConfiguration().setLocation(coordinator)));
+    assertEquals(
+        TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+        execute(Collections.singletonList(12)).getCode());
+    verify(executor).submitProcedure(any());
+  }
+}
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/ConfigExecutionTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/ConfigExecutionTest.java
index e2522b2730d..3d11defae7b 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/ConfigExecutionTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/execution/ConfigExecutionTest.java
@@ -19,7 +19,9 @@
 
 package org.apache.iotdb.db.queryengine.execution;
 
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
 import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
+import org.apache.iotdb.commons.exception.IoTDBException;
 import org.apache.iotdb.commons.schema.column.ColumnHeader;
 import org.apache.iotdb.db.queryengine.common.MPPQueryContext;
 import org.apache.iotdb.db.queryengine.common.QueryId;
@@ -60,6 +62,31 @@ public class ConfigExecutionTest {
     assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(), 
result.status.code);
   }
 
+  @Test
+  public void regionValidationErrorRetainsStatusAndMessage() {
+    ExecutorService executor = getExecutor();
+    try {
+      for (TSStatusCode code :
+          new TSStatusCode[] {
+            TSStatusCode.MIGRATE_REGION_ERROR,
+            TSStatusCode.RECONSTRUCT_REGION_ERROR,
+            TSStatusCode.EXTEND_REGION_ERROR,
+            TSStatusCode.REMOVE_REGION_PEER_ERROR
+          }) {
+        TSStatus status = new 
TSStatus(code.getStatusCode()).setMessage("Region 99 does not exist");
+        SettableFuture<ConfigTaskResult> future = SettableFuture.create();
+        future.setException(new IoTDBException(status));
+        IConfigTask task = clientManager -> future;
+        ConfigExecution execution = new ConfigExecution(genMPPQueryContext(), 
executor, task);
+        execution.start();
+
+        assertEquals(status, execution.getStatus().status);
+      }
+    } finally {
+      executor.shutdownNow();
+    }
+  }
+
   @Test
   public void normalConfigTaskWithResultTest() {
     TsBlock tsBlock =

Reply via email to