From a17f55eba3959653b205174d9b0cf22c466ae5d2 Mon Sep 17 00:00:00 2001 From: yuluo-yx Date: Sun, 9 Aug 2026 12:23:09 +0800 Subject: [PATCH] fix(controller): detect inactive raft masters --- .../impl/manager/RaftReplicasInfoManager.java | 2 +- .../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 046ced90c0a..d92ab556672 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 ControllerResult checkNotActiveBroker(Check List 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 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 b47f072c2c0..6062ee39e25 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.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 void testCheckNotActiveBrokerBrokerLiveTableNotEmptyIdentifiesNotActiveBr 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 replicaInfoTable = (Map) + 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 syncStateSetInfoTable = (Map) + FieldUtils.readField(raftReplicasInfoManager, "syncStateSetInfoTable", true); + syncStateSetInfoTable.put(brokerName, syncStateInfo); + + Map 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 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 brokerLiveTable = new HashMap<>();