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

tanxinyu pushed a commit to branch jira3512
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit d979d9d8c0a7c5935ded6746c5a8ddc794b2a2cd
Author: LebronAl <[email protected]>
AuthorDate: Fri Jun 17 13:45:32 2022 +0800

    finish
---
 .../org/apache/iotdb/consensus/IConsensus.java     |  2 ++
 .../multileader/MultiLeaderConsensus.java          |  5 +++++
 .../iotdb/consensus/ratis/RatisConsensus.java      | 12 +++++++++++
 .../consensus/standalone/StandAloneConsensus.java  |  6 ++++++
 .../iotdb/consensus/ratis/RatisConsensusTest.java  |  4 ----
 .../service/thrift/impl/InternalServiceImpl.java   | 25 +++++++++++++++++++++-
 thrift-commons/src/main/thrift/common.thrift       |  5 +++--
 7 files changed, 52 insertions(+), 7 deletions(-)

diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java 
b/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java
index f7c195eefc..60982cc49b 100644
--- a/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java
+++ b/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java
@@ -63,4 +63,6 @@ public interface IConsensus {
   boolean isLeader(ConsensusGroupId groupId);
 
   Peer getLeader(ConsensusGroupId groupId);
+
+  List<ConsensusGroupId> getAllConsensusGroupIds();
 }
diff --git 
a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java
 
b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java
index fbd9b280f3..b9e762a1db 100644
--- 
a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java
+++ 
b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java
@@ -236,6 +236,11 @@ public class MultiLeaderConsensus implements IConsensus {
     return new Peer(groupId, thisNode);
   }
 
+  @Override
+  public List<ConsensusGroupId> getAllConsensusGroupIds() {
+    return new ArrayList<>(stateMachineMap.keySet());
+  }
+
   public MultiLeaderServerImpl getImpl(ConsensusGroupId groupId) {
     return stateMachineMap.get(groupId);
   }
diff --git 
a/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java 
b/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java
index 8687dba132..33a523eb06 100644
--- 
a/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java
+++ 
b/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java
@@ -530,6 +530,18 @@ class RatisConsensus implements IConsensus {
     return new Peer(groupId, leaderEndpoint);
   }
 
+  @Override
+  public List<ConsensusGroupId> getAllConsensusGroupIds() {
+    List<ConsensusGroupId> ids = new ArrayList<>();
+    server
+        .getGroupIds()
+        .forEach(
+            groupId -> {
+              ids.add(Utils.fromRaftGroupIdToConsensusGroupId(groupId));
+            });
+    return ids;
+  }
+
   @Override
   public ConsensusGenericResponse triggerSnapshot(ConsensusGroupId groupId) {
     RaftGroupId raftGroupId = Utils.fromConsensusGroupIdToRaftGroupId(groupId);
diff --git 
a/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java
 
b/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java
index 1aa568416b..d87e27837f 100644
--- 
a/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java
+++ 
b/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java
@@ -44,6 +44,7 @@ import java.io.IOException;
 import java.nio.file.DirectoryStream;
 import java.nio.file.Files;
 import java.nio.file.Path;
+import java.util.ArrayList;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
@@ -218,6 +219,11 @@ class StandAloneConsensus implements IConsensus {
     return new Peer(groupId, thisNode);
   }
 
+  @Override
+  public List<ConsensusGroupId> getAllConsensusGroupIds() {
+    return new ArrayList<>(stateMachineMap.keySet());
+  }
+
   private String buildPeerDir(ConsensusGroupId groupId) {
     return storageDir + File.separator + groupId.getType().getValue() + "_" + 
groupId.getId();
   }
diff --git 
a/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
 
b/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
index 456a27fdde..b44d96460d 100644
--- 
a/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
+++ 
b/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java
@@ -119,10 +119,6 @@ public class RatisConsensusTest {
     }
   }
 
-  private int getLeaderOrdinal() {
-    return servers.get(0).getLeader(gid).getEndpoint().port - 6000;
-  }
-
   @Test
   public void basicConsensus3Copy() throws Exception {
     servers.get(0).addConsensusGroup(group.getGroupId(), group.getPeers());
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java
index 17efadbf68..557f3f3fa9 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java
@@ -95,7 +95,9 @@ import org.slf4j.LoggerFactory;
 
 import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.HashMap;
 import java.util.List;
+import java.util.Map;
 import java.util.Random;
 import java.util.stream.Collectors;
 
@@ -305,7 +307,7 @@ public class InternalServiceImpl implements 
InternalService.Iface {
 
   @Override
   public THeartbeatResp getHeartBeat(THeartbeatReq req) throws TException {
-    THeartbeatResp resp = new THeartbeatResp(req.getHeartbeatTimestamp());
+    THeartbeatResp resp = new THeartbeatResp(req.getHeartbeatTimestamp(), 
getJudgedLeaders());
     Random whetherToGetMetric = new Random();
     if 
(MetricConfigDescriptor.getInstance().getMetricConfig().getEnableMetric()
         && whetherToGetMetric.nextDouble() < loadBalanceThreshold) {
@@ -327,6 +329,27 @@ public class InternalServiceImpl implements 
InternalService.Iface {
     return resp;
   }
 
+  private Map<TConsensusGroupId, Boolean> getJudgedLeaders() {
+    Map<TConsensusGroupId, Boolean> result = new HashMap<>();
+    DataRegionConsensusImpl.getInstance()
+        .getAllConsensusGroupIds()
+        .forEach(
+            groupId -> {
+              result.put(
+                  groupId.convertToTConsensusGroupId(),
+                  DataRegionConsensusImpl.getInstance().isLeader(groupId));
+            });
+    SchemaRegionConsensusImpl.getInstance()
+        .getAllConsensusGroupIds()
+        .forEach(
+            groupId -> {
+              result.put(
+                  groupId.convertToTConsensusGroupId(),
+                  SchemaRegionConsensusImpl.getInstance().isLeader(groupId));
+            });
+    return result;
+  }
+
   private long getMemory(String gaugeName) {
     long result = 0;
     try {
diff --git a/thrift-commons/src/main/thrift/common.thrift 
b/thrift-commons/src/main/thrift/common.thrift
index e314a6327a..c381d5fb2d 100644
--- a/thrift-commons/src/main/thrift/common.thrift
+++ b/thrift-commons/src/main/thrift/common.thrift
@@ -83,8 +83,9 @@ struct THeartbeatReq {
 
 struct THeartbeatResp {
   1: required i64 heartbeatTimestamp
-  2: optional i16 cpu
-  3: optional i16 memory
+  2: required map<TConsensusGroupId, bool> judgedLeaders
+  3: optional i16 cpu
+  4: optional i16 memory
 }
 
 struct TDataNodeInfo {

Reply via email to