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<>();