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

RongtongJin pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git


The following commit(s) were added to refs/heads/develop by this push:
     new 7eee0fc366 fix(controller): detect inactive raft masters (#10885)
7eee0fc366 is described below

commit 7eee0fc366e869cbfe92b5c763f1a7daafc91af9
Author: shown <[email protected]>
AuthorDate: Thu Aug 13 15:46:06 2026 +0800

    fix(controller): detect inactive raft masters (#10885)
---
 .../impl/manager/RaftReplicasInfoManager.java      |  2 +-
 .../impl/manager/RaftReplicasInfoManagerTest.java  | 36 ++++++++++++++++++++++
 2 files changed, 37 insertions(+), 1 deletion(-)

diff --git 
a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManager.java
 
b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManager.java
index 046ced90c0..d92ab55667 100644
--- 
a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManager.java
+++ 
b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManager.java
@@ -126,7 +126,7 @@ public class RaftReplicasInfoManager extends 
ReplicasInfoManager {
         List<String> needReElectBrokerNames = scanNeedReelectBrokerSets(new 
BrokerValidPredicate() {
             @Override
             public boolean check(String clusterName, String brokerName, Long 
brokerId) {
-                return !isBrokerActive(clusterName, brokerName, brokerId, 
checkTime);
+                return isBrokerActive(clusterName, brokerName, brokerId, 
checkTime);
             }
         });
         Set<String> alreadyReportedBrokerName = 
notActiveBrokerIdentityInfoList.stream()
diff --git 
a/controller/src/test/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManagerTest.java
 
b/controller/src/test/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManagerTest.java
index b47f072c2c..6062ee39e2 100644
--- 
a/controller/src/test/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManagerTest.java
+++ 
b/controller/src/test/java/org/apache/rocketmq/controller/impl/manager/RaftReplicasInfoManagerTest.java
@@ -35,8 +35,11 @@ import org.junit.runner.RunWith;
 import org.mockito.Mock;
 import org.mockito.junit.MockitoJUnitRunner;
 
+import java.nio.charset.StandardCharsets;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.HashMap;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 
@@ -142,6 +145,39 @@ public class RaftReplicasInfoManagerTest {
         assertNotNull(result.getBody());
     }
 
+    @Test
+    @SuppressWarnings("unchecked")
+    public void testCheckNotActiveBrokerReportsInactiveMasterWithActiveSlave() 
throws IllegalAccessException {
+        String clusterName = "cluster";
+        String brokerName = "broker";
+        BrokerReplicaInfo brokerReplicaInfo = new 
BrokerReplicaInfo(clusterName, brokerName);
+        brokerReplicaInfo.addBroker(1L, "127.0.0.1:10911", "master");
+        brokerReplicaInfo.addBroker(2L, "127.0.0.1:10912", "slave");
+        Map<String, BrokerReplicaInfo> replicaInfoTable = (Map<String, 
BrokerReplicaInfo>)
+            FieldUtils.readField(raftReplicasInfoManager, "replicaInfoTable", 
true);
+        replicaInfoTable.put(brokerName, brokerReplicaInfo);
+
+        SyncStateInfo syncStateInfo = new SyncStateInfo(clusterName, 
brokerName);
+        syncStateInfo.updateMasterInfo(1L);
+        syncStateInfo.updateSyncStateSetInfo(new HashSet<>(Arrays.asList(1L, 
2L)));
+        Map<String, SyncStateInfo> syncStateSetInfoTable = (Map<String, 
SyncStateInfo>)
+            FieldUtils.readField(raftReplicasInfoManager, 
"syncStateSetInfoTable", true);
+        syncStateSetInfoTable.put(brokerName, syncStateInfo);
+
+        Map<BrokerIdentityInfo, BrokerLiveInfo> brokerLiveTable = new 
HashMap<>();
+        brokerLiveTable.put(new BrokerIdentityInfo(clusterName, brokerName, 
2L),
+            new BrokerLiveInfo(brokerName, "127.0.0.1:10912", 2L, 
System.currentTimeMillis(),
+                60_000L, null, 1, 100L, 1));
+        FieldUtils.writeDeclaredField(raftReplicasInfoManager, 
"brokerLiveTable", brokerLiveTable, true);
+
+        ControllerResult<CheckNotActiveBrokerResponse> result = 
raftReplicasInfoManager
+            .checkNotActiveBroker(new CheckNotActiveBrokerRequest());
+
+        assertEquals(ResponseCode.SUCCESS, result.getResponseCode());
+        assertTrue(new String(result.getBody(), StandardCharsets.UTF_8)
+            .contains("\"brokerName\":\"broker\""));
+    }
+
     @Test
     public void testCheckNotActiveBrokerSerializeErrorSetsErrorRemark() throws 
IllegalAccessException {
         Map<BrokerIdentityInfo, BrokerLiveInfo> brokerLiveTable = new 
HashMap<>();

Reply via email to