This is an automated email from the ASF dual-hosted git repository. haonan pushed a commit to branch ssl_between_nodes in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 4503a14721cd755c90c3b41f9cb60f648a826161 Author: HTHou <[email protected]> AuthorDate: Tue Jun 17 11:34:03 2025 +0800 Finishing consensus part --- .../manager/consensus/ConsensusManager.java | 5 ++ .../iotdb/consensus/config/IoTConsensusConfig.java | 76 +++++++++++++++++++++- .../consensus/config/PipeConsensusConfig.java | 76 +++++++++++++++++++++- .../apache/iotdb/consensus/config/RatisConfig.java | 25 +++++++ .../iotdb/consensus/ratis/RatisConsensus.java | 4 +- .../db/consensus/DataRegionConsensusImpl.java | 18 +++++ .../db/consensus/SchemaRegionConsensusImpl.java | 11 +++- 7 files changed, 207 insertions(+), 8 deletions(-) diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java index 0189e33ceb6..deda5c0a566 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java @@ -169,6 +169,11 @@ public class ConsensusManager { CONF.getConfigNodeRatisGrpcFlowControlWindow())) .setLeaderOutstandingAppendsMax( CONF.getConfigNodeRatisGrpcLeaderOutstandingAppendsMax()) + .setEnableSSL(COMMON_CONF.isEnableSSL()) + .setSslKeyStorePath(COMMON_CONF.getKeyStorePath()) + .setSslKeyStorePassword(COMMON_CONF.getKeyStorePwd()) + .setSslTrustStorePath(COMMON_CONF.getTrustStorePath()) + .setSslTrustStorePassword(COMMON_CONF.getTrustStorePwd()) .build()) .setRpc( RatisConfig.Rpc.newBuilder() diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/IoTConsensusConfig.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/IoTConsensusConfig.java index 0621aee23ef..32c4664b60d 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/IoTConsensusConfig.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/IoTConsensusConfig.java @@ -84,6 +84,12 @@ public class IoTConsensusConfig { private final int thriftMaxFrameSize; private final int maxClientNumForEachNode; + private final boolean isEnableSSL; + private final String sslTrustStorePath; + private final String sslTrustStorePassword; + private final String sslKeyStorePath; + private final String sslKeyStorePassword; + private RPC( int rpcSelectorThreadNum, int rpcMinConcurrentClientNum, @@ -94,7 +100,12 @@ public class IoTConsensusConfig { int connectionTimeoutInMs, boolean printLogWhenThriftClientEncounterException, int thriftMaxFrameSize, - int maxClientNumForEachNode) { + int maxClientNumForEachNode, + boolean isEnableSSL, + String sslTrustStorePath, + String sslTrustStorePassword, + String sslKeyStorePath, + String sslKeyStorePassword) { this.rpcSelectorThreadNum = rpcSelectorThreadNum; this.rpcMinConcurrentClientNum = rpcMinConcurrentClientNum; this.rpcMaxConcurrentClientNum = rpcMaxConcurrentClientNum; @@ -105,6 +116,11 @@ public class IoTConsensusConfig { this.printLogWhenThriftClientEncounterException = printLogWhenThriftClientEncounterException; this.thriftMaxFrameSize = thriftMaxFrameSize; this.maxClientNumForEachNode = maxClientNumForEachNode; + this.isEnableSSL = isEnableSSL; + this.sslTrustStorePath = sslTrustStorePath; + this.sslTrustStorePassword = sslTrustStorePassword; + this.sslKeyStorePath = sslKeyStorePath; + this.sslKeyStorePassword = sslKeyStorePassword; } public int getRpcSelectorThreadNum() { @@ -147,6 +163,26 @@ public class IoTConsensusConfig { return maxClientNumForEachNode; } + public boolean isEnableSSL() { + return isEnableSSL; + } + + public String getSslTrustStorePath() { + return sslTrustStorePath; + } + + public String getSslTrustStorePassword() { + return sslTrustStorePassword; + } + + public String getSslKeyStorePath() { + return sslKeyStorePath; + } + + public String getSslKeyStorePassword() { + return sslKeyStorePassword; + } + public static RPC.Builder newBuilder() { return new RPC.Builder(); } @@ -165,6 +201,12 @@ public class IoTConsensusConfig { private int thriftMaxFrameSize = 536870912; private int maxClientNumForEachNode = DefaultProperty.MAX_CLIENT_NUM_FOR_EACH_NODE; + private boolean isEnableSSL = false; + private String sslTrustStorePath = ""; + private String sslTrustStorePassword = ""; + private String sslKeyStorePath = ""; + private String sslKeyStorePassword = ""; + public RPC.Builder setRpcSelectorThreadNum(int rpcSelectorThreadNum) { this.rpcSelectorThreadNum = rpcSelectorThreadNum; return this; @@ -218,6 +260,31 @@ public class IoTConsensusConfig { return this; } + public Builder setEnableSSL(boolean isEnableSSL) { + this.isEnableSSL = isEnableSSL; + return this; + } + + public Builder setSslTrustStorePath(String sslTrustStorePath) { + this.sslTrustStorePath = sslTrustStorePath; + return this; + } + + public Builder setSslTrustStorePassword(String sslTrustStorePassword) { + this.sslTrustStorePassword = sslTrustStorePassword; + return this; + } + + public Builder setSslKeyStorePath(String sslKeyStorePath) { + this.sslKeyStorePath = sslKeyStorePath; + return this; + } + + public Builder setSslKeyStorePassword(String sslKeyStorePassword) { + this.sslKeyStorePassword = sslKeyStorePassword; + return this; + } + public RPC build() { return new RPC( rpcSelectorThreadNum, @@ -229,7 +296,12 @@ public class IoTConsensusConfig { connectionTimeoutInMs, printLogWhenThriftClientEncounterException, thriftMaxFrameSize, - maxClientNumForEachNode); + maxClientNumForEachNode, + isEnableSSL, + sslTrustStorePath, + sslTrustStorePassword, + sslKeyStorePath, + sslKeyStorePassword); } } } diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/PipeConsensusConfig.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/PipeConsensusConfig.java index f0366cf0087..da06c60a624 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/PipeConsensusConfig.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/PipeConsensusConfig.java @@ -87,17 +87,33 @@ public class PipeConsensusConfig { private final int connectionTimeoutInMs; private final int thriftMaxFrameSize; + private boolean isEnableSSL = false; + private String sslTrustStorePath = ""; + private String sslTrustStorePassword = ""; + private String sslKeyStorePath = ""; + private String sslKeyStorePassword = ""; + public RPC( int rpcMaxConcurrentClientNum, int thriftServerAwaitTimeForStopService, boolean isRpcThriftCompressionEnabled, int connectionTimeoutInMs, - int thriftMaxFrameSize) { + int thriftMaxFrameSize, + boolean isEnableSSL, + String sslTrustStorePath, + String sslTrustStorePassword, + String sslKeyStorePath, + String sslKeyStorePassword) { this.rpcMaxConcurrentClientNum = rpcMaxConcurrentClientNum; this.thriftServerAwaitTimeForStopService = thriftServerAwaitTimeForStopService; this.isRpcThriftCompressionEnabled = isRpcThriftCompressionEnabled; this.connectionTimeoutInMs = connectionTimeoutInMs; this.thriftMaxFrameSize = thriftMaxFrameSize; + this.isEnableSSL = isEnableSSL; + this.sslTrustStorePath = sslTrustStorePath; + this.sslTrustStorePassword = sslTrustStorePassword; + this.sslKeyStorePath = sslKeyStorePath; + this.sslKeyStorePassword = sslKeyStorePassword; } public int getRpcMaxConcurrentClientNum() { @@ -120,6 +136,26 @@ public class PipeConsensusConfig { return thriftMaxFrameSize; } + public boolean isEnableSSL() { + return isEnableSSL; + } + + public String getSslTrustStorePath() { + return sslTrustStorePath; + } + + public String getSslTrustStorePassword() { + return sslTrustStorePassword; + } + + public String getSslKeyStorePath() { + return sslKeyStorePath; + } + + public String getSslKeyStorePassword() { + return sslKeyStorePassword; + } + public static RPC.Builder newBuilder() { return new RPC.Builder(); } @@ -131,6 +167,12 @@ public class PipeConsensusConfig { private int connectionTimeoutInMs = (int) TimeUnit.SECONDS.toMillis(60); private int thriftMaxFrameSize = 536870912; + private boolean isEnableSSL = false; + private String sslTrustStorePath = ""; + private String sslTrustStorePassword = ""; + private String sslKeyStorePath = ""; + private String sslKeyStorePassword = ""; + public RPC.Builder setRpcMaxConcurrentClientNum(int rpcMaxConcurrentClientNum) { this.rpcMaxConcurrentClientNum = rpcMaxConcurrentClientNum; return this; @@ -157,13 +199,43 @@ public class PipeConsensusConfig { return this; } + public Builder setEnableSSL(boolean isEnableSSL) { + this.isEnableSSL = isEnableSSL; + return this; + } + + public Builder setSslTrustStorePath(String sslTrustStorePath) { + this.sslTrustStorePath = sslTrustStorePath; + return this; + } + + public Builder setSslTrustStorePassword(String sslTrustStorePassword) { + this.sslTrustStorePassword = sslTrustStorePassword; + return this; + } + + public Builder setSslKeyStorePath(String sslKeyStorePath) { + this.sslKeyStorePath = sslKeyStorePath; + return this; + } + + public Builder setSslKeyStorePassword(String sslKeyStorePassword) { + this.sslKeyStorePassword = sslKeyStorePassword; + return this; + } + public RPC build() { return new RPC( rpcMaxConcurrentClientNum, thriftServerAwaitTimeForStopService, isRpcThriftCompressionEnabled, connectionTimeoutInMs, - thriftMaxFrameSize); + thriftMaxFrameSize, + isEnableSSL, + sslTrustStorePath, + sslTrustStorePassword, + sslKeyStorePath, + sslKeyStorePassword); } } } diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java index f5d8aa48fec..e3781338bae 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java @@ -829,6 +829,31 @@ public class RatisConfig { this.leaderOutstandingAppendsMax = leaderOutstandingAppendsMax; return this; } + + public Grpc.Builder setEnableSSL(boolean isEnableSSL) { + this.isEnableSSL = isEnableSSL; + return this; + } + + public Grpc.Builder setSslTrustStorePath(String sslTrustStorePath) { + this.sslTrustStorePath = sslTrustStorePath; + return this; + } + + public Grpc.Builder setSslTrustStorePassword(String sslTrustStorePassword) { + this.sslTrustStorePassword = sslTrustStorePassword; + return this; + } + + public Grpc.Builder setSslKeyStorePath(String sslKeyStorePath) { + this.sslKeyStorePath = sslKeyStorePath; + return this; + } + + public Grpc.Builder setSslKeyStorePassword(String sslKeyStorePassword) { + this.sslKeyStorePassword = sslKeyStorePassword; + return this; + } } } diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java index 96779812c0a..a4ca82f75e3 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java @@ -264,7 +264,7 @@ class RatisConsensus implements IConsensus { try { diskGuardian.stop(); } catch (InterruptedException e) { - logger.warn("{}: interrupted when shutting down add Executor with exception {}", this, e); + logger.warn("{}: interrupted when shutting down add Executor with exception ", this, e); Thread.currentThread().interrupt(); } finally { clientManager.close(); @@ -815,7 +815,7 @@ class RatisConsensus implements IConsensus { try { leaderId = server.get().getDivision(raftGroupId).getInfo().getLeaderId(); } catch (IOException e) { - logger.warn("fetch division info for group " + groupId + " failed due to: ", e); + logger.warn("fetch division info for group {} failed due to: ", groupId, e); return null; } if (leaderId == null) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java index 4785f03c483..47bb6033b5c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java @@ -21,6 +21,8 @@ package org.apache.iotdb.db.consensus; import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType; import org.apache.iotdb.common.rpc.thrift.TEndPoint; +import org.apache.iotdb.commons.conf.CommonConfig; +import org.apache.iotdb.commons.conf.CommonDescriptor; import org.apache.iotdb.commons.consensus.ConsensusGroupId; import org.apache.iotdb.commons.consensus.DataRegionId; import org.apache.iotdb.commons.memory.IMemoryBlock; @@ -78,6 +80,7 @@ public class DataRegionConsensusImpl { private static class DataRegionConsensusImplHolder { private static final IoTDBConfig CONF = IoTDBDescriptor.getInstance().getConfig(); + private static final CommonConfig COMMON_CONF = CommonDescriptor.getInstance().getConfig(); private static final DataNodeMemoryConfig MEMORY_CONFIG = IoTDBDescriptor.getInstance().getMemoryConfig(); @@ -135,6 +138,11 @@ public class DataRegionConsensusImpl { CONF.getThriftServerAwaitTimeForStopService()) .setThriftMaxFrameSize(CONF.getThriftMaxFrameSize()) .setMaxClientNumForEachNode(CONF.getMaxClientNumForEachNode()) + .setEnableSSL(COMMON_CONF.isEnableSSL()) + .setSslKeyStorePath(COMMON_CONF.getKeyStorePath()) + .setSslKeyStorePassword(COMMON_CONF.getKeyStorePwd()) + .setSslTrustStorePath(COMMON_CONF.getTrustStorePath()) + .setSslTrustStorePassword(COMMON_CONF.getTrustStorePwd()) .build()) .setReplication( IoTConsensusConfig.Replication.newBuilder() @@ -159,6 +167,11 @@ public class DataRegionConsensusImpl { .setThriftServerAwaitTimeForStopService( CONF.getThriftServerAwaitTimeForStopService()) .setThriftMaxFrameSize(CONF.getThriftMaxFrameSize()) + .setEnableSSL(COMMON_CONF.isEnableSSL()) + .setSslKeyStorePath(COMMON_CONF.getKeyStorePath()) + .setSslKeyStorePassword(COMMON_CONF.getKeyStorePwd()) + .setSslTrustStorePath(COMMON_CONF.getTrustStorePath()) + .setSslTrustStorePassword(COMMON_CONF.getTrustStorePwd()) .build()) .setPipe( PipeConsensusConfig.Pipe.newBuilder() @@ -205,6 +218,11 @@ public class DataRegionConsensusImpl { CONF.getDataRatisConsensusGrpcFlowControlWindow())) .setLeaderOutstandingAppendsMax( CONF.getDataRatisConsensusGrpcLeaderOutstandingAppendsMax()) + .setEnableSSL(COMMON_CONF.isEnableSSL()) + .setSslKeyStorePath(COMMON_CONF.getKeyStorePath()) + .setSslKeyStorePassword(COMMON_CONF.getKeyStorePwd()) + .setSslTrustStorePath(COMMON_CONF.getTrustStorePath()) + .setSslTrustStorePassword(COMMON_CONF.getTrustStorePwd()) .build()) .setRpc( RatisConfig.Rpc.newBuilder() diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java index b3e79fa1eec..5195e2c4ccd 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java @@ -21,6 +21,8 @@ package org.apache.iotdb.db.consensus; import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType; import org.apache.iotdb.common.rpc.thrift.TEndPoint; +import org.apache.iotdb.commons.conf.CommonConfig; +import org.apache.iotdb.commons.conf.CommonDescriptor; import org.apache.iotdb.commons.consensus.SchemaRegionId; import org.apache.iotdb.consensus.ConsensusFactory; import org.apache.iotdb.consensus.IConsensus; @@ -56,7 +58,8 @@ public class SchemaRegionConsensusImpl { private static class SchemaRegionConsensusImplHolder { - private static IoTDBConfig CONF; + private static final IoTDBConfig CONF = IoTDBDescriptor.getInstance().getConfig(); + private static final CommonConfig COMMON_CONF = CommonDescriptor.getInstance().getConfig(); private static IConsensus INSTANCE; // Make sure both statics are initialized. @@ -65,7 +68,6 @@ public class SchemaRegionConsensusImpl { } private static void reinitializeStatics() { - CONF = IoTDBDescriptor.getInstance().getConfig(); INSTANCE = ConsensusFactory.getConsensusImpl( CONF.getSchemaRegionConsensusProtocolClass(), @@ -102,6 +104,11 @@ public class SchemaRegionConsensusImpl { .setLeaderOutstandingAppendsMax( CONF .getSchemaRatisConsensusGrpcLeaderOutstandingAppendsMax()) + .setEnableSSL(COMMON_CONF.isEnableSSL()) + .setSslKeyStorePath(COMMON_CONF.getKeyStorePath()) + .setSslKeyStorePassword(COMMON_CONF.getKeyStorePwd()) + .setSslTrustStorePath(COMMON_CONF.getTrustStorePath()) + .setSslTrustStorePassword(COMMON_CONF.getTrustStorePwd()) .build()) .setRpc( RatisConfig.Rpc.newBuilder()
