This is an automated email from the ASF dual-hosted git repository.
CRZbulabula pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/dev/1.3 by this push:
new a177a93017c [To dev/1.3] Handle invalid DataNode id in region
operations (#18478)
a177a93017c is described below
commit a177a93017c1741ecc65809fc68976bdb8341cbe
Author: Yongzao <[email protected]>
AuthorDate: Sun Aug 16 16:18:30 2026 +0800
[To dev/1.3] Handle invalid DataNode id in region operations (#18478)
---
.../IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java | 48 +++++
.../commit/IoTDBRegionReconstructForIoTV1IT.java | 50 +++++
.../iotdb/confignode/manager/ProcedureManager.java | 35 +++-
.../ProcedureManagerRegionOperationTest.java | 216 +++++++++++++++++++++
4 files changed, 342 insertions(+), 7 deletions(-)
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 afc1163fef0..32302eb0e0c 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
@@ -19,13 +19,16 @@
package org.apache.iotdb.confignode.it.regionmigration.pass.commit;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.client.sync.SyncConfigNodeIServiceClient;
import
org.apache.iotdb.confignode.it.regionmigration.IoTDBRegionOperationReliabilityITFramework;
+import org.apache.iotdb.confignode.rpc.thrift.TExtendRegionReq;
import org.apache.iotdb.confignode.rpc.thrift.TShowRegionResp;
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;
@@ -38,6 +41,7 @@ import org.slf4j.LoggerFactory;
import java.sql.Connection;
import java.sql.Statement;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.Optional;
@@ -120,6 +124,50 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
}
}
+ @Test
+ public void rejectInvalidTargetDataNodeTest() throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
+
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
+ .setDataReplicationFactor(1)
+ .setSchemaReplicationFactor(1);
+
+ EnvFactory.getEnv().initClusterEnvironment(1, 3);
+
+ try (final Connection connection =
makeItCloseQuietly(EnvFactory.getEnv().getConnection());
+ final Statement statement =
makeItCloseQuietly(connection.createStatement());
+ SyncConfigNodeIServiceClient client =
+ (SyncConfigNodeIServiceClient)
EnvFactory.getEnv().getLeaderConfigNodeConnection()) {
+ statement.execute(INSERTION1);
+ statement.execute(FLUSH_COMMAND);
+
+ Map<Integer, Set<Integer>> regionMap = getAllRegionMap(statement);
+ Set<Integer> allDataNodeIds = getAllDataNodes(statement);
+ Assert.assertFalse(regionMap.isEmpty());
+
+ int selectedRegion = regionMap.keySet().iterator().next();
+ int unknownDataNodeId = Collections.max(allDataNodeIds) + 1000;
+ int configNodeId =
client.showCluster().getConfigNodeList().get(0).getConfigNodeId();
+ Assert.assertFalse(allDataNodeIds.contains(configNodeId));
+
+ assertExtendRegionRejected(client, selectedRegion, unknownDataNodeId);
+ assertExtendRegionRejected(client, selectedRegion, configNodeId);
+ Assert.assertEquals(regionMap, getAllRegionMap(statement));
+ }
+ }
+
+ private void assertExtendRegionRejected(
+ SyncConfigNodeIServiceClient client, int regionId, int dataNodeId)
throws Exception {
+ TSStatus status =
+ client.extendRegion(new
TExtendRegionReq(Collections.singletonList(regionId), dataNodeId));
+ Assert.assertEquals(TSStatusCode.EXTEND_REGION_ERROR.getStatusCode(),
status.getCode());
+ Assert.assertEquals(
+ String.format("Target DataNode %s does not exist in the cluster",
dataNodeId),
+ status.getMessage());
+ }
+
private void regionGroupExpand(
Statement statement,
SyncConfigNodeIServiceClient client,
diff --git
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionReconstructForIoTV1IT.java
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionReconstructForIoTV1IT.java
index 2957b93dd12..a48f5515bff 100644
---
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionReconstructForIoTV1IT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionReconstructForIoTV1IT.java
@@ -19,9 +19,11 @@
package org.apache.iotdb.confignode.it.regionmigration.pass.commit;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.client.sync.SyncConfigNodeIServiceClient;
import org.apache.iotdb.commons.cluster.NodeStatus;
import
org.apache.iotdb.confignode.it.regionmigration.IoTDBRegionOperationReliabilityITFramework;
+import org.apache.iotdb.confignode.rpc.thrift.TReconstructRegionReq;
import org.apache.iotdb.consensus.ConsensusFactory;
import org.apache.iotdb.isession.SessionDataSet;
import org.apache.iotdb.it.env.EnvFactory;
@@ -29,12 +31,14 @@ import org.apache.iotdb.it.env.cluster.node.DataNodeWrapper;
import org.apache.iotdb.it.framework.IoTDBTestRunner;
import org.apache.iotdb.itbase.category.ClusterIT;
import org.apache.iotdb.rpc.StatementExecutionException;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.session.Session;
import org.apache.commons.io.FileUtils;
import org.apache.tsfile.read.common.RowRecord;
import org.awaitility.Awaitility;
import org.junit.Assert;
+import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
import org.slf4j.Logger;
@@ -43,6 +47,7 @@ import org.slf4j.LoggerFactory;
import java.io.File;
import java.sql.Connection;
import java.sql.Statement;
+import java.util.Collections;
import java.util.Iterator;
import java.util.Map;
import java.util.Set;
@@ -145,4 +150,49 @@ public class IoTDBRegionReconstructForIoTV1IT extends
IoTDBRegionOperationReliab
Assert.assertEquals("1.0",
rowRecord.getFields().get(1).getStringValue());
}
}
+
+ @Test
+ public void rejectInvalidTargetDataNodeTest() throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
+
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
+ .setDataReplicationFactor(1)
+ .setSchemaReplicationFactor(1);
+
+ EnvFactory.getEnv().initClusterEnvironment(1, 3);
+
+ try (Connection connection =
makeItCloseQuietly(EnvFactory.getEnv().getConnection());
+ Statement statement = makeItCloseQuietly(connection.createStatement());
+ SyncConfigNodeIServiceClient client =
+ (SyncConfigNodeIServiceClient)
EnvFactory.getEnv().getLeaderConfigNodeConnection()) {
+ statement.execute(INSERTION1);
+ statement.execute(FLUSH_COMMAND);
+
+ Map<Integer, Set<Integer>> regionMap = getAllRegionMap(statement);
+ Set<Integer> allDataNodeIds = getAllDataNodes(statement);
+ Assert.assertFalse(regionMap.isEmpty());
+
+ int selectedRegion = regionMap.keySet().iterator().next();
+ int unknownDataNodeId = Collections.max(allDataNodeIds) + 1000;
+ int configNodeId =
client.showCluster().getConfigNodeList().get(0).getConfigNodeId();
+ Assert.assertFalse(allDataNodeIds.contains(configNodeId));
+
+ assertReconstructRegionRejected(client, selectedRegion,
unknownDataNodeId);
+ assertReconstructRegionRejected(client, selectedRegion, configNodeId);
+ Assert.assertEquals(regionMap, getAllRegionMap(statement));
+ }
+ }
+
+ private void assertReconstructRegionRejected(
+ SyncConfigNodeIServiceClient client, int regionId, int dataNodeId)
throws Exception {
+ TSStatus status =
+ client.reconstructRegion(
+ new TReconstructRegionReq(Collections.singletonList(regionId),
dataNodeId));
+ Assert.assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(),
status.getCode());
+ Assert.assertEquals(
+ String.format("Target DataNode %s does not exist in the cluster",
dataNodeId),
+ status.getMessage());
+ }
}
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 f656444ca19..364f3fc9c7e 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
@@ -1011,9 +1011,13 @@ public class ProcedureManager {
}
public TSStatus reconstructRegion(TReconstructRegionReq req) {
- RegionMaintainHandler handler = env.getRegionMaintainHandler();
final TDataNodeLocation targetDataNode =
-
configManager.getNodeManager().getRegisteredDataNode(req.getDataNodeId()).getLocation();
+ getRegisteredDataNodeLocationOrNull(req.getDataNodeId());
+ if (targetDataNode == null) {
+ return targetDataNodeNotExistStatus(
+ req.getDataNodeId(), TSStatusCode.RECONSTRUCT_REGION_ERROR);
+ }
+ RegionMaintainHandler handler = env.getRegionMaintainHandler();
try (AutoCloseableLock ignoredLock =
AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
List<ReconstructRegionProcedure> procedures = new ArrayList<>();
@@ -1051,8 +1055,16 @@ public class ProcedureManager {
}
public TSStatus extendRegions(TExtendRegionReq req) {
+ final TDataNodeLocation targetDataNode =
+ getRegisteredDataNodeLocationOrNull(req.getDataNodeId());
+ if (targetDataNode == null) {
+ return targetDataNodeNotExistStatus(req.getDataNodeId(),
TSStatusCode.EXTEND_REGION_ERROR);
+ }
return processExtendOrRemoveRegions(
- req.getRegionId(), req, this::extendOneRegion,
TSStatusCode.EXTEND_REGION_ERROR);
+ req.getRegionId(),
+ req,
+ (regionId, request) -> extendOneRegion(regionId, request,
targetDataNode),
+ TSStatusCode.EXTEND_REGION_ERROR);
}
public TSStatus removeRegions(TRemoveRegionReq req) {
@@ -1098,7 +1110,8 @@ public class ProcedureManager {
return resp;
}
- private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) {
+ private TSStatus extendOneRegion(
+ int theRegionId, TExtendRegionReq req, TDataNodeLocation targetDataNode)
{
try (AutoCloseableLock ignoredLock =
AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
TConsensusGroupId regionId;
@@ -1112,9 +1125,6 @@ public class ProcedureManager {
.setMessage("get region group id fail");
}
- // find target dn
- final TDataNodeLocation targetDataNode =
-
configManager.getNodeManager().getRegisteredDataNode(req.getDataNodeId()).getLocation();
// select coordinator for adding peer
RegionMaintainHandler handler = env.getRegionMaintainHandler();
final TDataNodeLocation coordinator =
@@ -1141,6 +1151,17 @@ public class ProcedureManager {
}
}
+ private TDataNodeLocation getRegisteredDataNodeLocationOrNull(int
dataNodeId) {
+ TDataNodeConfiguration dataNodeConfiguration =
+ configManager.getNodeManager().getRegisteredDataNode(dataNodeId);
+ return dataNodeConfiguration == null ? null :
dataNodeConfiguration.getLocation();
+ }
+
+ private TSStatus targetDataNodeNotExistStatus(int dataNodeId, TSStatusCode
statusCode) {
+ return new TSStatus(statusCode.getStatusCode())
+ .setMessage(String.format("Target DataNode %s does not exist in the
cluster", dataNodeId));
+ }
+
private TSStatus removeOneRegion(int theRegionId, TRemoveRegionReq req) {
try (AutoCloseableLock ignoredLock =
AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
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..bc1be58289f
--- /dev/null
+++
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerRegionOperationTest.java
@@ -0,0 +1,216 @@
+/*
+ * 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.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.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.AddRegionPeerProcedure;
+import
org.apache.iotdb.confignode.procedure.impl.region.ReconstructRegionProcedure;
+import org.apache.iotdb.confignode.rpc.thrift.TExtendRegionReq;
+import org.apache.iotdb.confignode.rpc.thrift.TReconstructRegionReq;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.Optional;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.ReentrantLock;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyInt;
+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;
+
+public class ProcedureManagerRegionOperationTest {
+
+ private static final int REGION_ID = 1;
+ private static final int CONFIG_NODE_ID = 0;
+ private static final int TARGET_DATA_NODE_ID = 2;
+ private static final int COORDINATOR_DATA_NODE_ID = 3;
+ private static final int UNKNOWN_DATA_NODE_ID = 9999;
+
+ private ConfigManager configManager;
+ private NodeManager nodeManager;
+ private PartitionManager partitionManager;
+ private ProcedureExecutor<ConfigNodeProcedureEnv> executor;
+ private ConfigNodeProcedureEnv env;
+ private RegionMaintainHandler regionMaintainHandler;
+ private ProcedureManager procedureManager;
+
+ @Before
+ @SuppressWarnings("unchecked")
+ public void setUp() {
+ configManager = mock(ConfigManager.class);
+ nodeManager = mock(NodeManager.class);
+ partitionManager = mock(PartitionManager.class);
+ executor = mock(ProcedureExecutor.class);
+ env = mock(ConfigNodeProcedureEnv.class);
+ regionMaintainHandler = mock(RegionMaintainHandler.class);
+
+ when(configManager.getNodeManager()).thenReturn(nodeManager);
+ when(configManager.getPartitionManager()).thenReturn(partitionManager);
+ when(executor.getProcedures()).thenReturn(new ConcurrentHashMap<>());
+ when(env.getRegionMaintainHandler()).thenReturn(regionMaintainHandler);
+
+ procedureManager = new ProcedureManager(configManager,
mock(ProcedureInfo.class));
+ procedureManager.setExecutor(executor);
+ procedureManager.setEnv(env);
+ }
+
+ @Test
+ public void reconstructRegionRejectsUnknownDataNodeId() {
+ assertReconstructRejected(UNKNOWN_DATA_NODE_ID);
+ }
+
+ @Test
+ public void reconstructRegionRejectsConfigNodeId() {
+ assertReconstructRejected(CONFIG_NODE_ID);
+ }
+
+ @Test
+ public void reconstructRegionSubmitsProcedureForRegisteredDataNode() {
+ configureRegisteredTarget(true);
+
+ TSStatus status =
+ procedureManager.reconstructRegion(
+ new TReconstructRegionReq(Collections.singletonList(REGION_ID),
TARGET_DATA_NODE_ID));
+
+ Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(),
status.getCode());
+ verify(executor,
times(1)).submitProcedure(any(ReconstructRegionProcedure.class));
+ }
+
+ @Test
+ public void extendRegionRejectsUnknownDataNodeId() {
+ assertExtendRejected(UNKNOWN_DATA_NODE_ID);
+ }
+
+ @Test
+ public void extendRegionRejectsConfigNodeId() {
+ assertExtendRejected(CONFIG_NODE_ID);
+ }
+
+ @Test
+ public void extendRegionSubmitsProcedureForRegisteredDataNode() {
+ configureRegisteredTarget(false);
+
+ TSStatus status =
+ procedureManager.extendRegions(
+ new TExtendRegionReq(Collections.singletonList(REGION_ID),
TARGET_DATA_NODE_ID));
+
+ Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(),
status.getCode());
+ verify(executor,
times(1)).submitProcedure(any(AddRegionPeerProcedure.class));
+ }
+
+ private void assertReconstructRejected(int dataNodeId) {
+ when(nodeManager.getRegisteredDataNode(dataNodeId)).thenReturn(new
TDataNodeConfiguration());
+
+ TSStatus status =
+ procedureManager.reconstructRegion(
+ new TReconstructRegionReq(Collections.singletonList(REGION_ID),
dataNodeId));
+
+ assertRejected(status, dataNodeId, TSStatusCode.RECONSTRUCT_REGION_ERROR);
+ }
+
+ private void assertExtendRejected(int dataNodeId) {
+ when(nodeManager.getRegisteredDataNode(dataNodeId)).thenReturn(new
TDataNodeConfiguration());
+
+ TSStatus status =
+ procedureManager.extendRegions(
+ new TExtendRegionReq(Collections.singletonList(REGION_ID),
dataNodeId));
+
+ assertRejected(status, dataNodeId, TSStatusCode.EXTEND_REGION_ERROR);
+ }
+
+ private void assertRejected(TSStatus status, int dataNodeId, TSStatusCode
expectedCode) {
+ Assert.assertEquals(expectedCode.getStatusCode(), status.getCode());
+ Assert.assertEquals(
+ String.format("Target DataNode %s does not exist in the cluster",
dataNodeId),
+ status.getMessage());
+ verify(partitionManager,
never()).generateTConsensusGroupIdByRegionId(anyInt());
+ verify(executor, never()).submitProcedure(any(Procedure.class));
+ }
+
+ private void configureRegisteredTarget(boolean reconstruct) {
+ TConsensusGroupId consensusGroupId =
+ new TConsensusGroupId(TConsensusGroupType.DataRegion, REGION_ID);
+ TDataNodeLocation targetDataNode = dataNodeLocation(TARGET_DATA_NODE_ID,
7000);
+ TDataNodeLocation coordinatorDataNode =
dataNodeLocation(COORDINATOR_DATA_NODE_ID, 7100);
+ TDataNodeConfiguration targetDataNodeConfiguration =
+ new TDataNodeConfiguration().setLocation(targetDataNode);
+ TDataNodeConfiguration coordinatorDataNodeConfiguration =
+ new TDataNodeConfiguration().setLocation(coordinatorDataNode);
+ TRegionReplicaSet replicaSet =
+ new TRegionReplicaSet(
+ consensusGroupId,
+ reconstruct
+ ? Arrays.asList(targetDataNode, coordinatorDataNode)
+ : Collections.singletonList(coordinatorDataNode));
+
+ when(nodeManager.getRegisteredDataNode(TARGET_DATA_NODE_ID))
+ .thenReturn(targetDataNodeConfiguration);
+ when(nodeManager.filterDataNodeThroughStatus(NodeStatus.Running))
+ .thenReturn(Arrays.asList(targetDataNodeConfiguration,
coordinatorDataNodeConfiguration));
+ when(partitionManager.generateTConsensusGroupIdByRegionId(REGION_ID))
+ .thenReturn(Optional.of(consensusGroupId));
+ when(partitionManager.getAllReplicaSets(TARGET_DATA_NODE_ID))
+ .thenReturn(reconstruct ? Collections.singletonList(replicaSet) :
Collections.emptyList());
+ when(partitionManager.getAllReplicaSetsMap(TConsensusGroupType.DataRegion))
+ .thenReturn(Collections.singletonMap(consensusGroupId, replicaSet));
+ when(regionMaintainHandler.filterDataNodeWithOtherRegionReplica(
+ consensusGroupId,
+ targetDataNode,
+ NodeStatus.Running,
+ NodeStatus.Removing,
+ NodeStatus.ReadOnly))
+ .thenReturn(Optional.of(coordinatorDataNode));
+ when(env.getSubmitRegionMigrateLock()).thenReturn(new ReentrantLock());
+ when(executor.submitProcedure(any(Procedure.class))).thenReturn(1L);
+ }
+
+ private TDataNodeLocation dataNodeLocation(int dataNodeId, int basePort) {
+ return new TDataNodeLocation(
+ dataNodeId,
+ new TEndPoint("127.0.0.1", basePort),
+ new TEndPoint("127.0.0.1", basePort + 1),
+ new TEndPoint("127.0.0.1", basePort + 2),
+ new TEndPoint("127.0.0.1", basePort + 3),
+ new TEndPoint("127.0.0.1", basePort + 4));
+ }
+}