From 735edecd5ff7a79e8007929c001a087144756f87 Mon Sep 17 00:00:00 2001 From: Yongzao <532741407@qq.com> Date: Tue, 22 Sep 2026 15:48:29 +0800 Subject: [PATCH 1/3] Improve region operation submission results in CLI --- .../org/apache/iotdb/cli/AbstractCli.java | 6 +- .../org/apache/iotdb/jdbc/IoTDBStatement.java | 11 ++ .../apache/iotdb/jdbc/IoTDBStatementTest.java | 23 +++ .../confignode/i18n/ManagerMessages.java | 13 ++ .../confignode/i18n/ManagerMessages.java | 13 ++ .../confignode/manager/ProcedureManager.java | 131 ++++++++++++------ ...ProcedureManagerReconstructRegionTest.java | 26 +++- .../executor/ClusterConfigTaskExecutor.java | 8 +- 8 files changed, 181 insertions(+), 50 deletions(-) diff --git a/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java b/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java index aad1074860fbf..409d8c488bca3 100644 --- a/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java +++ b/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java @@ -46,6 +46,7 @@ import java.sql.ResultSet; import java.sql.ResultSetMetaData; import java.sql.SQLException; +import java.sql.SQLWarning; import java.sql.Statement; import java.time.ZoneId; import java.util.ArrayList; @@ -704,7 +705,10 @@ private static int executeQuery(CliContext ctx, IoTDBConnection connection, Stri } } } else { - ctx.getPrinter().println("Msg: " + SUCCESS_MESSAGE); + // Config statements may return a detailed submission result through JDBC warnings. + SQLWarning warning = statement.getWarnings(); + ctx.getPrinter() + .println("Msg: " + (warning == null ? SUCCESS_MESSAGE : warning.getMessage())); } } catch (Exception e) { ctx.getPrinter().println("Msg: " + e); diff --git a/iotdb-client/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java b/iotdb-client/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java index 2af6e970a9255..2d10ff3a4ba03 100644 --- a/iotdb-client/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java +++ b/iotdb-client/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java @@ -356,6 +356,7 @@ private T callWithRetryAndReconnect(TFunction rpc, Function */ private boolean executeSQL(String sql) throws TException, SQLException { isCancelled = false; + warningChain = null; TSExecuteStatementReq execReq = new TSExecuteStatementReq(sessionId, sql, stmtId); int rows = fetchSize; if (maxRows != 0 && fetchSize > maxRows) { @@ -382,6 +383,7 @@ private boolean executeSQL(String sql) throws TException, SQLException { } catch (StatementExecutionException e) { throw new IoTDBSQLException(e.getMessage(), execResp.getStatus()); } + setExecutionWarning(execResp.getStatus()); if (execResp.isSetDatabase()) { connection.changeDefaultDatabase(execResp.getDatabase()); @@ -439,6 +441,7 @@ public int[] executeBatch() throws SQLException { private int[] executeBatchSQL() throws TException, BatchUpdateException, SQLException { isCancelled = false; + warningChain = null; TSExecuteBatchStatementReq execReq = new TSExecuteBatchStatementReq(sessionId, batchSQLList); TSStatus execResp = callWithRetryAndReconnect( @@ -509,6 +512,7 @@ public ResultSet executeQuery(String sql, long timeoutInMS) throws SQLException private ResultSet executeQuerySQL(String sql, long timeoutInMS) throws TException, SQLException { isCancelled = false; + warningChain = null; TSExecuteStatementReq execReq = new TSExecuteStatementReq(sessionId, sql, stmtId); int rows = fetchSize; if (maxRows != 0 && fetchSize > maxRows) { @@ -600,6 +604,7 @@ public int executeUpdate(String arg0, String[] arg1) throws SQLException { private int executeUpdateSQL(final String sql) throws TException, IoTDBSQLException, SQLException { + warningChain = null; final TSExecuteStatementReq execReq = new TSExecuteStatementReq(sessionId, sql, stmtId); final TSExecuteStatementResp execResp = callWithRetryAndReconnect( @@ -619,9 +624,15 @@ private int executeUpdateSQL(final String sql) } catch (final StatementExecutionException e) { throw new IoTDBSQLException(e.getMessage(), execResp.getStatus()); } + setExecutionWarning(execResp.getStatus()); return 0; } + private void setExecutionWarning(TSStatus status) { + final String message = status.getMessage(); + warningChain = message == null || message.isEmpty() ? null : new SQLWarning(message); + } + @Override public Connection getConnection() { return connection; 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 5cfd5342e05e6..18ce5b705efb5 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; @@ -33,6 +36,7 @@ import org.mockito.MockitoAnnotations; import java.sql.SQLException; +import java.sql.SQLWarning; import java.time.ZoneId; import static org.junit.Assert.assertEquals; @@ -106,4 +110,23 @@ public void setTimeoutTest() throws SQLException { statement.setQueryTimeout(100); Assert.assertEquals(100, statement.getQueryTimeout()); } + + @SuppressWarnings("resource") + @Test + public void executionStatusMessageIsExposedAsWarningAndClearedOnNextExecution() throws Exception { + TSExecuteStatementResp response = new TSExecuteStatementResp(); + response.setStatus( + new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()).setMessage("submission result")); + when(client.executeStatementV2(any())).thenReturn(response); + + IoTDBStatement statement = new IoTDBStatement(connection, client, sessionId, zoneID, 0, 1L); + Assert.assertFalse(statement.execute("MIGRATE REGION 1 FROM 2 TO 3")); + SQLWarning warning = statement.getWarnings(); + Assert.assertNotNull(warning); + Assert.assertEquals("submission result", warning.getMessage()); + + response.setStatus(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())); + Assert.assertFalse(statement.execute("FLUSH")); + 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 03ec800325182..12a381bc73e1a 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,19 @@ 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_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360 = + "Total regions: %d, successfully submitted: %d, failed to submit: %d\n"; + public static final String MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29 = + "Region %d: Successfully submitted\n"; + public static final String MESSAGE_REGION_ARG_ARG_01229B27 = + "Region %d: %s\n"; + 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 f145e04052ec8..07712b3a3dc95 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,19 @@ 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_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360 = + "Region 总数:%d,成功提交:%d,提交失败:%d\n"; + public static final String MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29 = + "Region %d:提交成功\n"; + public static final String MESSAGE_REGION_ARG_ARG_01229B27 = + "Region %d:%s\n"; + 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 dfbb440d9d6aa..9b80cecd960e0 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 @@ -963,7 +963,7 @@ private TSStatus checkExtendRegion( 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; } @@ -1157,26 +1157,26 @@ public TSStatus migrateRegion(TMigrateRegionReq migrateRegionReq) { 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(migrateRegionReq.getFromId()); + final TDataNodeLocation destDataNode = + getRegisteredDataNodeLocationOrNull(migrateRegionReq.getToId()); + if (originalDataNode == null) { return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode()) .setMessage( String.format( - "Source DataNode %s does not exist in the cluster", + ManagerMessages + .MESSAGE_SOURCE_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_2255633C, migrateRegionReq.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", + ManagerMessages + .MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF, migrateRegionReq.getToId())); } - final TDataNodeLocation originalDataNode = originalDataNodeConfiguration.getLocation(); - final TDataNodeLocation destDataNode = destDataNodeConfiguration.getLocation(); final RegionMaintainHandler handler = env.getRegionMaintainHandler(); TSStatus resp = new TSStatus(); @@ -1189,15 +1189,16 @@ public TSStatus migrateRegion(TMigrateRegionReq migrateRegionReq) { migrateOneRegion( migrateRegionReq, theRegionId, originalDataNode, destDataNode, handler); if (subStatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) { - messageBuilder.append("region ").append(theRegionId).append(": Successfully submitted\n"); + messageBuilder.append( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29, theRegionId)); success++; } else { - messageBuilder - .append("region ") - .append(theRegionId) - .append(": ") - .append(subStatus.getMessage()) - .append('\n'); + messageBuilder.append( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_ARG_01229B27, + theRegionId, + subStatus.getMessage())); } resp.addToSubStatus(subStatus); } @@ -1205,10 +1206,13 @@ public TSStatus migrateRegion(TMigrateRegionReq migrateRegionReq) { messageBuilder.insert( 0, String.format( - "Total regions: %d, successfully submitted: %d, failed to submit: %d\n", - total, success, total - success)); + ManagerMessages + .MESSAGE_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360, + total, + success, + total - success)); resp.setCode( - total == success + total > 0 && total == success ? TSStatusCode.SUCCESS_STATUS.getStatusCode() : TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode()); resp.setMessage(messageBuilder.toString()); @@ -1230,7 +1234,9 @@ private TSStatus migrateOneRegion( } else { LOGGER.error(ManagerMessages.GET_REGION_GROUP_ID_FAIL); return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode()) - .setMessage(ManagerMessages.GET_REGION_GROUP_ID_FAIL); + .setMessage( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, theRegionId)); } // select coordinator for adding peer @@ -1288,7 +1294,9 @@ private TSStatus migrateOneRegion( * 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) { @@ -1301,11 +1309,17 @@ public TSStatus reconstructRegion(TReconstructRegionReq req) { return new TSStatus(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode()) .setMessage( String.format( - "Target DataNode %s does not exist in the cluster", req.getDataNodeId())); + ManagerMessages + .MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF, + req.getDataNodeId())); } try (AutoCloseableLock ignoredLock = AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) { List procedures = new ArrayList<>(); + TSStatus resp = new TSStatus(); + StringBuilder messageBuilder = new StringBuilder(); + int total = 0; + int success = 0; Set seenRegionIds = new HashSet<>(); for (int x : req.getRegionIds()) { if (!seenRegionIds.add(x)) { @@ -1316,6 +1330,7 @@ public TSStatus reconstructRegion(TReconstructRegionReq req) { req.getDataNodeId()); continue; } + total++; Optional regionIdOptional = configManager.getPartitionManager().findTConsensusGroupIdByRegionId(x); if (!regionIdOptional.isPresent()) { @@ -1324,6 +1339,14 @@ public TSStatus reconstructRegion(TReconstructRegionReq req) { .LOG_SKIP_NON_EXISTENT_REGION_ID_ARG_IN_RECONSTRUCTREGION_REQUEST_TO_DATANODE_ARG_7F76D789, x, req.getDataNodeId()); + TSStatus subStatus = + new TSStatus(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode()) + .setMessage( + String.format(ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, x)); + resp.addToSubStatus(subStatus); + messageBuilder.append( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_ARG_01229B27, x, subStatus.getMessage())); continue; } TConsensusGroupId regionId = regionIdOptional.get(); @@ -1341,8 +1364,12 @@ public TSStatus reconstructRegion(TReconstructRegionReq req) { return status; } procedures.add(new ReconstructRegionProcedure(regionId, targetDataNode, coordinator)); + resp.addToSubStatus(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())); + messageBuilder.append( + String.format(ManagerMessages.MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29, x)); + success++; } - // all checks pass, submit all procedures + // Submit all procedures that passed validation after the complete request has been checked. procedures.forEach( reconstructRegionProcedure -> { this.executor.submitProcedure(reconstructRegionProcedure); @@ -1350,8 +1377,22 @@ public TSStatus reconstructRegion(TReconstructRegionReq req) { ManagerMessages.RECONSTRUCTREGION_SUBMIT_RECONSTRUCTREGIONPROCEDURE_SUCCESSFULLY, reconstructRegionProcedure); }); + + messageBuilder.insert( + 0, + String.format( + ManagerMessages + .MESSAGE_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360, + total, + success, + total - success)); + resp.setCode( + total > 0 && total == success + ? TSStatusCode.SUCCESS_STATUS.getStatusCode() + : TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode()); + resp.setMessage(messageBuilder.toString()); + return resp; } - return RpcUtils.SUCCESS_STATUS; } public TSStatus extendRegions(TExtendRegionReq req) { @@ -1377,15 +1418,14 @@ private TSStatus processExtendOrRemoveRegions( total++; TSStatus subStatus = regionAction.apply(regionId, req); if (subStatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) { - messageBuilder.append("region ").append(regionId).append(": Successfully submitted\n"); + messageBuilder.append( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29, regionId)); success++; } else { - messageBuilder - .append("region ") - .append(regionId) - .append(": ") - .append(subStatus.getMessage()) - .append('\n'); + messageBuilder.append( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_ARG_01229B27, regionId, subStatus.getMessage())); } resp.addToSubStatus(subStatus); } @@ -1393,11 +1433,16 @@ private TSStatus processExtendOrRemoveRegions( messageBuilder.insert( 0, String.format( - "Total regions: %d, successfully submitted: %d, failed to submit: %d\n", - total, success, total - success)); + ManagerMessages + .MESSAGE_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360, + total, + success, + total - success)); resp.setCode( - total == success ? TSStatusCode.SUCCESS_STATUS.getStatusCode() : errorCode.getStatusCode()); + total > 0 && total == success + ? TSStatusCode.SUCCESS_STATUS.getStatusCode() + : errorCode.getStatusCode()); resp.setMessage(messageBuilder.toString()); return resp; } @@ -1413,7 +1458,9 @@ private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) { } else { LOGGER.error(ManagerMessages.GET_REGION_GROUP_ID_FAIL); return new TSStatus(TSStatusCode.EXTEND_REGION_ERROR.getStatusCode()) - .setMessage(ManagerMessages.GET_REGION_GROUP_ID_FAIL); + .setMessage( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, theRegionId)); } // find target dn @@ -1425,7 +1472,9 @@ private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) { return new TSStatus(TSStatusCode.EXTEND_REGION_ERROR.getStatusCode()) .setMessage( String.format( - "Target DataNode %s does not exist in the cluster", req.getDataNodeId())); + ManagerMessages + .MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF, + req.getDataNodeId())); } // select coordinator for adding peer RegionMaintainHandler handler = env.getRegionMaintainHandler(); @@ -1466,12 +1515,14 @@ private TSStatus removeOneRegion(int theRegionId, TRemoveRegionReq req) { } else { LOGGER.error(ManagerMessages.GET_REGION_GROUP_ID_FAIL); return new TSStatus(TSStatusCode.REMOVE_REGION_PEER_ERROR.getStatusCode()) - .setMessage(ManagerMessages.GET_REGION_GROUP_ID_FAIL); + .setMessage( + String.format( + ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, theRegionId)); } // find target dn final TDataNodeLocation targetDataNode = - configManager.getNodeManager().getRegisteredDataNode(req.getDataNodeId()).getLocation(); + getRegisteredDataNodeLocationOrNull(req.getDataNodeId()); // select coordinator for removing peer RegionMaintainHandler handler = env.getRegionMaintainHandler(); 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 index 4e2b58daa631e..4514033f7fc7a 100644 --- 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 @@ -139,8 +139,11 @@ 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()); + TSStatus status = manager.reconstructRegion(request); + assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(), status.getCode()); + assertTrue(status.getMessage().contains("Total regions: 3")); + assertTrue(status.getMessage().contains("successfully submitted: 2")); + assertTrue(status.getMessage().contains("failed to submit: 1")); ArgumentCaptor captor = ArgumentCaptor.forClass(ReconstructRegionProcedure.class); @@ -153,11 +156,24 @@ public void testDuplicateAndNonExistentRegionIdsAreSkippedInInputOrder() { } @Test - public void testRequestWithNoUsableRegionIdsSucceedsWithoutSubmittingProcedure() { + public void testRequestWithNoUsableRegionIdsFailsWithoutSubmittingProcedure() { TReconstructRegionReq request = new TReconstructRegionReq(Arrays.asList(99, 99), 7, Model.TREE); - assertEquals( - TSStatusCode.SUCCESS_STATUS.getStatusCode(), manager.reconstructRegion(request).getCode()); + TSStatus status = manager.reconstructRegion(request); + assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(), status.getCode()); + assertTrue(status.getMessage().contains("Total regions: 1")); + assertTrue(status.getMessage().contains("failed to submit: 1")); + verify(executor, times(0)).submitProcedure(any()); + } + + @Test + public void testRequestWithNoRegionIdsFailsWithoutSubmittingProcedure() { + TReconstructRegionReq request = + new TReconstructRegionReq(Collections.emptyList(), 7, Model.TREE); + + TSStatus status = manager.reconstructRegion(request); + assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(), status.getCode()); + assertTrue(status.getMessage().contains("Total regions: 0")); verify(executor, times(0)).submitProcedure(any()); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java index d7dd2abcb1d38..a8481557eaf4a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java @@ -3909,7 +3909,7 @@ public SettableFuture migrateRegion(final MigrateRegionTask mi future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); + future.set(new ConfigTaskResult(status)); } } catch (Exception e) { future.setException(e); @@ -4097,7 +4097,7 @@ public SettableFuture reconstructRegion( future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); + future.set(new ConfigTaskResult(status)); } } catch (Exception e) { future.setException(e); @@ -4120,7 +4120,7 @@ public SettableFuture extendRegion(ExtendRegionTask extendRegi future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); + future.set(new ConfigTaskResult(status)); } } catch (Exception e) { future.setException(e); @@ -4143,7 +4143,7 @@ public SettableFuture removeRegion(RemoveRegionTask removeRegi future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); + future.set(new ConfigTaskResult(status)); } } catch (Exception e) { future.setException(e); From 835d16dc663fdfb2434e29b9e5dcd9bf694767d0 Mon Sep 17 00:00:00 2001 From: Yongzao <532741407@qq.com> Date: Wed, 23 Sep 2026 14:37:55 +0800 Subject: [PATCH 2/3] Validate complete region operation requests before submitting procedures --- .../org/apache/iotdb/cli/AbstractCli.java | 6 +- .../org/apache/iotdb/cli/AbstractCliTest.java | 44 ++ .../org/apache/iotdb/jdbc/IoTDBStatement.java | 11 - .../apache/iotdb/jdbc/IoTDBStatementTest.java | 48 +- .../confignode/i18n/ManagerMessages.java | 15 +- .../confignode/i18n/ManagerMessages.java | 15 +- .../confignode/manager/ProcedureManager.java | 501 ++++++------------ ...ProcedureManagerReconstructRegionTest.java | 206 ------- .../ProcedureManagerRegionOperationTest.java | 401 ++++++++++++++ .../executor/ClusterConfigTaskExecutor.java | 8 +- .../execution/ConfigExecutionTest.java | 27 + 11 files changed, 684 insertions(+), 598 deletions(-) delete mode 100644 iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerReconstructRegionTest.java create mode 100644 iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerRegionOperationTest.java diff --git a/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java b/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java index 409d8c488bca3..aad1074860fbf 100644 --- a/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java +++ b/iotdb-client/cli/src/main/java/org/apache/iotdb/cli/AbstractCli.java @@ -46,7 +46,6 @@ import java.sql.ResultSet; import java.sql.ResultSetMetaData; import java.sql.SQLException; -import java.sql.SQLWarning; import java.sql.Statement; import java.time.ZoneId; import java.util.ArrayList; @@ -705,10 +704,7 @@ private static int executeQuery(CliContext ctx, IoTDBConnection connection, Stri } } } else { - // Config statements may return a detailed submission result through JDBC warnings. - SQLWarning warning = statement.getWarnings(); - ctx.getPrinter() - .println("Msg: " + (warning == null ? SUCCESS_MESSAGE : warning.getMessage())); + ctx.getPrinter().println("Msg: " + SUCCESS_MESSAGE); } } catch (Exception e) { ctx.getPrinter().println("Msg: " + e); 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 01644d58412cc..deba69d50be19 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 @@ 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.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.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 void setUp() throws Exception { 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/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java b/iotdb-client/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java index 2d10ff3a4ba03..2af6e970a9255 100644 --- a/iotdb-client/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java +++ b/iotdb-client/jdbc/src/main/java/org/apache/iotdb/jdbc/IoTDBStatement.java @@ -356,7 +356,6 @@ private T callWithRetryAndReconnect(TFunction rpc, Function */ private boolean executeSQL(String sql) throws TException, SQLException { isCancelled = false; - warningChain = null; TSExecuteStatementReq execReq = new TSExecuteStatementReq(sessionId, sql, stmtId); int rows = fetchSize; if (maxRows != 0 && fetchSize > maxRows) { @@ -383,7 +382,6 @@ private boolean executeSQL(String sql) throws TException, SQLException { } catch (StatementExecutionException e) { throw new IoTDBSQLException(e.getMessage(), execResp.getStatus()); } - setExecutionWarning(execResp.getStatus()); if (execResp.isSetDatabase()) { connection.changeDefaultDatabase(execResp.getDatabase()); @@ -441,7 +439,6 @@ public int[] executeBatch() throws SQLException { private int[] executeBatchSQL() throws TException, BatchUpdateException, SQLException { isCancelled = false; - warningChain = null; TSExecuteBatchStatementReq execReq = new TSExecuteBatchStatementReq(sessionId, batchSQLList); TSStatus execResp = callWithRetryAndReconnect( @@ -512,7 +509,6 @@ public ResultSet executeQuery(String sql, long timeoutInMS) throws SQLException private ResultSet executeQuerySQL(String sql, long timeoutInMS) throws TException, SQLException { isCancelled = false; - warningChain = null; TSExecuteStatementReq execReq = new TSExecuteStatementReq(sessionId, sql, stmtId); int rows = fetchSize; if (maxRows != 0 && fetchSize > maxRows) { @@ -604,7 +600,6 @@ public int executeUpdate(String arg0, String[] arg1) throws SQLException { private int executeUpdateSQL(final String sql) throws TException, IoTDBSQLException, SQLException { - warningChain = null; final TSExecuteStatementReq execReq = new TSExecuteStatementReq(sessionId, sql, stmtId); final TSExecuteStatementResp execResp = callWithRetryAndReconnect( @@ -624,15 +619,9 @@ private int executeUpdateSQL(final String sql) } catch (final StatementExecutionException e) { throw new IoTDBSQLException(e.getMessage(), execResp.getStatus()); } - setExecutionWarning(execResp.getStatus()); return 0; } - private void setExecutionWarning(TSStatus status) { - final String message = status.getMessage(); - warningChain = message == null || message.isEmpty() ? null : new SQLWarning(message); - } - @Override public Connection getConnection() { return connection; 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 18ce5b705efb5..2f7d993e37d7d 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 @@ -36,7 +36,6 @@ import org.mockito.MockitoAnnotations; import java.sql.SQLException; -import java.sql.SQLWarning; import java.time.ZoneId; import static org.junit.Assert.assertEquals; @@ -113,20 +112,37 @@ public void setTimeoutTest() throws SQLException { @SuppressWarnings("resource") @Test - public void executionStatusMessageIsExposedAsWarningAndClearedOnNextExecution() throws Exception { - TSExecuteStatementResp response = new TSExecuteStatementResp(); - response.setStatus( - new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()).setMessage("submission result")); - when(client.executeStatementV2(any())).thenReturn(response); - - IoTDBStatement statement = new IoTDBStatement(connection, client, sessionId, zoneID, 0, 1L); - Assert.assertFalse(statement.execute("MIGRATE REGION 1 FROM 2 TO 3")); - SQLWarning warning = statement.getWarnings(); - Assert.assertNotNull(warning); - Assert.assertEquals("submission result", warning.getMessage()); - - response.setStatus(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())); - Assert.assertFalse(statement.execute("FLUSH")); - Assert.assertNull(statement.getWarnings()); + 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 12a381bc73e1a..6a36ac7130231 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,13 +295,14 @@ 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_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360 = - "Total regions: %d, successfully submitted: %d, failed to submit: %d\n"; - public static final String MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29 = - "Region %d: Successfully submitted\n"; - public static final String MESSAGE_REGION_ARG_ARG_01229B27 = - "Region %d: %s\n"; + 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 = 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 07712b3a3dc95..3c35171dd06f9 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,13 +293,14 @@ 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_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360 = - "Region 总数:%d,成功提交:%d,提交失败:%d\n"; - public static final String MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29 = - "Region %d:提交成功\n"; - public static final String MESSAGE_REGION_ARG_ARG_01229B27 = - "Region %d:%s\n"; + 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 = 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 9b80cecd960e0..e6efa881da1e9 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.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.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.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; @@ -973,7 +970,7 @@ private TSStatus checkExtendRegion( private TSStatus checkRemoveRegion( TRemoveRegionReq req, TConsensusGroupId regionId, - @Nullable TDataNodeLocation targetDataNode, + TDataNodeLocation targetDataNode, TDataNodeLocation coordinator) { String failMessage = regionOperationCommonCheck( @@ -991,12 +988,11 @@ private TSStatus checkRemoveRegion( .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,22 +1148,19 @@ private String checkRegionOperationModelCorrectness(TConsensusGroupId regionId, // 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 TDataNodeLocation originalDataNode = - getRegisteredDataNodeLocationOrNull(migrateRegionReq.getFromId()); - final TDataNodeLocation destDataNode = - getRegisteredDataNodeLocationOrNull(migrateRegionReq.getToId()); + getRegisteredDataNodeLocationOrNull(req.getFromId()); + final TDataNodeLocation destDataNode = getRegisteredDataNodeLocationOrNull(req.getToId()); if (originalDataNode == null) { return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode()) .setMessage( String.format( ManagerMessages .MESSAGE_SOURCE_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_2255633C, - migrateRegionReq.getFromId())); + req.getFromId())); } if (destDataNode == null) { return new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode()) @@ -1175,117 +1168,49 @@ public TSStatus migrateRegion(TMigrateRegionReq migrateRegionReq) { String.format( ManagerMessages .MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF, - migrateRegionReq.getToId())); + 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 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( - String.format( - ManagerMessages.MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29, theRegionId)); - success++; - } else { - messageBuilder.append( - String.format( - ManagerMessages.MESSAGE_REGION_ARG_ARG_01229B27, - theRegionId, - subStatus.getMessage())); + List 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 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( - ManagerMessages - .MESSAGE_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360, - total, - success, - total - success)); - resp.setCode( - total > 0 && total == success - ? TSStatusCode.SUCCESS_STATUS.getStatusCode() - : TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode()); - resp.setMessage(messageBuilder.toString()); - return resp; + return submitRegionOperationProcedures(procedures); } } - private TSStatus migrateOneRegion( - TMigrateRegionReq migrateRegionReq, - int theRegionId, - TDataNodeLocation originalDataNode, - TDataNodeLocation destDataNode, - RegionMaintainHandler handler) { - TConsensusGroupId regionGroupId; - Optional 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( - String.format( - ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, theRegionId)); - } - - // 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; - } - - // 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()); - } - /** * Resolve the location of a registered DataNode, or return {@code null} if the given id does not * belong to any registered DataNode (e.g. it is a ConfigNode id or simply does not exist). {@link @@ -1300,56 +1225,28 @@ private TDataNodeLocation getRegisteredDataNodeLocationOrNull(int dataNodeId) { } 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( - ManagerMessages - .MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF, - 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 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 procedures = new ArrayList<>(); - TSStatus resp = new TSStatus(); - StringBuilder messageBuilder = new StringBuilder(); - int total = 0; - int success = 0; - Set 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; - } - total++; - Optional 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()); - TSStatus subStatus = - new TSStatus(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode()) - .setMessage( - String.format(ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, x)); - resp.addToSubStatus(subStatus); - messageBuilder.append( - String.format( - ManagerMessages.MESSAGE_REGION_ARG_ARG_01229B27, x, subStatus.getMessage())); - continue; - } - TConsensusGroupId regionId = regionIdOptional.get(); + for (TConsensusGroupId regionId : regionIds) { final TDataNodeLocation coordinator = handler .filterDataNodeWithOtherRegionReplica( @@ -1359,116 +1256,22 @@ public TSStatus reconstructRegion(TReconstructRegionReq req) { 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)); - resp.addToSubStatus(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())); - messageBuilder.append( - String.format(ManagerMessages.MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29, x)); - success++; - } - // Submit all procedures that passed validation after the complete request has been checked. - procedures.forEach( - reconstructRegionProcedure -> { - this.executor.submitProcedure(reconstructRegionProcedure); - LOGGER.info( - ManagerMessages.RECONSTRUCTREGION_SUBMIT_RECONSTRUCTREGIONPROCEDURE_SUCCESSFULLY, - reconstructRegionProcedure); - }); - - messageBuilder.insert( - 0, - String.format( - ManagerMessages - .MESSAGE_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360, - total, - success, - total - success)); - resp.setCode( - total > 0 && total == success - ? TSStatusCode.SUCCESS_STATUS.getStatusCode() - : TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode()); - resp.setMessage(messageBuilder.toString()); - return resp; - } - } - - 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 TSStatus processExtendOrRemoveRegions( - Iterable regionIds, - R req, - BiFunction 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( - String.format( - ManagerMessages.MESSAGE_REGION_ARG_SUCCESSFULLY_SUBMITTED_BB0F2E29, regionId)); - success++; - } else { - messageBuilder.append( - String.format( - ManagerMessages.MESSAGE_REGION_ARG_ARG_01229B27, regionId, subStatus.getMessage())); } - resp.addToSubStatus(subStatus); + return submitRegionOperationProcedures(procedures); } - - messageBuilder.insert( - 0, - String.format( - ManagerMessages - .MESSAGE_TOTAL_REGIONS_ARG_SUCCESSFULLY_SUBMITTED_ARG_FAILED_TO_SUBMIT_ARG_2F69D360, - total, - success, - total - success)); - - resp.setCode( - total > 0 && total == success - ? TSStatusCode.SUCCESS_STATUS.getStatusCode() - : errorCode.getStatusCode()); - resp.setMessage(messageBuilder.toString()); - return resp; } - private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) { + public TSStatus extendRegions(TExtendRegionReq req) { try (AutoCloseableLock ignoredLock = AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) { - TConsensusGroupId regionId; - Optional 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( - String.format( - ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, theRegionId)); - } - - // 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( @@ -1476,105 +1279,119 @@ private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) { .MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF, 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); + + List 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 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 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( String.format( - ManagerMessages.MESSAGE_REGION_ARG_DOES_NOT_EXIST_3C8400C9, theRegionId)); + ManagerMessages + .MESSAGE_TARGET_DATANODE_ARG_DOES_NOT_EXIST_IN_THE_CLUSTER_679D59AF, + req.getDataNodeId())); } - // find target dn - final TDataNodeLocation targetDataNode = - getRegisteredDataNodeLocationOrNull(req.getDataNodeId()); - - // 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 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 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 requestedRegionIds, List regionIds, TSStatusCode errorCode) { + if (requestedRegionIds == null || requestedRegionIds.isEmpty()) { + return new TSStatus(errorCode.getStatusCode()) + .setMessage(ManagerMessages.MESSAGE_REGION_IDS_MUST_NOT_BE_EMPTY_B42DAAFD); + } + Set 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 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> 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 4514033f7fc7a..0000000000000 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerReconstructRegionTest.java +++ /dev/null @@ -1,206 +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 executor; - private NodeManager nodeManager; - private PartitionManager partitionManager; - private final ConcurrentHashMap> 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 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); - - TSStatus status = manager.reconstructRegion(request); - assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(), status.getCode()); - assertTrue(status.getMessage().contains("Total regions: 3")); - assertTrue(status.getMessage().contains("successfully submitted: 2")); - assertTrue(status.getMessage().contains("failed to submit: 1")); - - ArgumentCaptor 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 testRequestWithNoUsableRegionIdsFailsWithoutSubmittingProcedure() { - TReconstructRegionReq request = new TReconstructRegionReq(Arrays.asList(99, 99), 7, Model.TREE); - - TSStatus status = manager.reconstructRegion(request); - assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(), status.getCode()); - assertTrue(status.getMessage().contains("Total regions: 1")); - assertTrue(status.getMessage().contains("failed to submit: 1")); - verify(executor, times(0)).submitProcedure(any()); - } - - @Test - public void testRequestWithNoRegionIdsFailsWithoutSubmittingProcedure() { - TReconstructRegionReq request = - new TReconstructRegionReq(Collections.emptyList(), 7, Model.TREE); - - TSStatus status = manager.reconstructRegion(request); - assertEquals(TSStatusCode.RECONSTRUCT_REGION_ERROR.getStatusCode(), status.getCode()); - assertTrue(status.getMessage().contains("Total regions: 0")); - 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 0000000000000..d0fb6efbb0ad2 --- /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 parameters() { + List 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> procedures = + new ConcurrentHashMap<>(); + + private ProcedureManager manager; + private ProcedureExecutor 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 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 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 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 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 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/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java index a8481557eaf4a..d7dd2abcb1d38 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java @@ -3909,7 +3909,7 @@ public SettableFuture migrateRegion(final MigrateRegionTask mi future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(status)); + future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); } } catch (Exception e) { future.setException(e); @@ -4097,7 +4097,7 @@ public SettableFuture reconstructRegion( future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(status)); + future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); } } catch (Exception e) { future.setException(e); @@ -4120,7 +4120,7 @@ public SettableFuture extendRegion(ExtendRegionTask extendRegi future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(status)); + future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); } } catch (Exception e) { future.setException(e); @@ -4143,7 +4143,7 @@ public SettableFuture removeRegion(RemoveRegionTask removeRegi future.setException(new IoTDBException(status)); return future; } else { - future.set(new ConfigTaskResult(status)); + future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS)); } } catch (Exception e) { future.setException(e); 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 e2522b2730dea..3d11defae7b87 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 void normalConfigTaskTest() { 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 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 = From 52034de39e06c54fc793e6f96acb180ad2fc0275 Mon Sep 17 00:00:00 2001 From: Yongzao <532741407@qq.com> Date: Wed, 23 Sep 2026 15:57:11 +0800 Subject: [PATCH 3/3] Reject invalid multi-region submissions in ITs --- .../IoTDBMigrateMultiRegionForIoTV1IT.java | 9 +- ...BRegionGroupExpandAndShrinkForIoTV1IT.java | 123 ++++++++---------- 2 files changed, 56 insertions(+), 76 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 da361b32a7cf0..b915f5f83f0cf 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.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 void multiRegionMigrateTest() throws Exception { 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 8bb1671b2d38e..c4902cbc93142 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.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 void extendRegionToInvalidDataNodeTest() throws Exception { 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 void multiRegionNormalTest() throws Exception { } } - /** 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 void multiRegionExpandPartialExistTest() throws Exception { Map> regionMap = getAllRegionMap(statement); Set allDataNodeId = getAllDataNodes(statement); - List allRegions = new ArrayList<>(regionMap.keySet()); - List selectedRegions = allRegions.subList(0, Math.min(3, allRegions.size())); + Assert.assertEquals(2, regionMap.size()); + List selectedRegions = new ArrayList<>(regionMap.keySet()); int targetDataNode = findDataNodeNotContainsAnyRegion(allDataNodeId, regionMap, selectedRegions); - // first expand some regions individually - List 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 void multiRegionShrinkPartialNotExistTest() throws Exception { .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 void multiRegionShrinkPartialNotExistTest() throws Exception { Map> regionMap = getAllRegionMap(statement); Set allDataNodeId = getAllDataNodes(statement); - List allRegions = new ArrayList<>(regionMap.keySet()); - List selectedRegions = allRegions.subList(0, Math.min(3, allRegions.size())); + Assert.assertEquals(2, regionMap.size()); + List selectedRegions = new ArrayList<>(regionMap.keySet()); int targetDataNode = findDataNodeNotContainsAnyRegion(allDataNodeId, regionMap, selectedRegions); @@ -404,33 +397,34 @@ public void multiRegionShrinkPartialNotExistTest() throws Exception { // first expand all regions to target DataNode multiRegionGroupExpand(statement, client, selectedRegions, targetDataNode); - // then shrink some regions individually - List 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> 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 @@ private void executeMultiRegionOperation( 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,