Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -38,6 +41,7 @@
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;
Expand Down Expand Up @@ -120,6 +124,50 @@ public void singleRegionTest() throws Exception {
}
}

@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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,22 +19,26 @@

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;
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;
Expand All @@ -43,6 +47,7 @@
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;
Expand Down Expand Up @@ -145,4 +150,49 @@ public void normal1C3DTest() throws Exception {
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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -1011,9 +1011,13 @@ public TSStatus migrateRegion(TMigrateRegionReq migrateRegionReq) {
}

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<>();
Expand Down Expand Up @@ -1051,8 +1055,16 @@ public TSStatus reconstructRegion(TReconstructRegionReq req) {
}

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) {
Expand Down Expand Up @@ -1098,7 +1110,8 @@ private <R> TSStatus processExtendOrRemoveRegions(
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;
Expand All @@ -1112,9 +1125,6 @@ private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) {
.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 =
Expand All @@ -1141,6 +1151,17 @@ private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) {
}
}

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())) {
Expand Down
Loading
Loading