This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git
The following commit(s) were added to refs/heads/rocketmq-studio by this push:
new c0a27842 fix: distinguish runtime diagnostic failures from empty data
(#1143)
c0a27842 is described below
commit c0a27842a6bc1a4e85b15f094f2f9c49f011a7f1
Author: aias00 <[email protected]>
AuthorDate: Fri Aug 7 00:33:07 2026 -0700
fix: distinguish runtime diagnostic failures from empty data (#1143)
* fix: distinguish failed consumer connection scans
* fix: distinguish failed runtime diagnostics
* fix: expose client scan discovery failures
---
.../provider/apache/RocketMQClientProvider.java | 15 ++++-
.../provider/apache/RocketMQMetadataProvider.java | 5 +-
.../apache/RocketMQClientProviderTest.java | 67 ++++++++++++++++++++++
.../apache/RocketMQMetadataProviderTest.java | 38 ++++++++++++
4 files changed, 122 insertions(+), 3 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
index 574f0d3f..19a43ee0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProvider.java
@@ -212,12 +212,14 @@ public class RocketMQClientProvider implements
ClientProvider {
private List<ClientConnectionVO> findConsumerConnections(MQAdminExt
adminExt, String clusterId) {
List<ClientConnectionVO> result = new ArrayList<>();
Set<String> groups = collectSubscriptionGroups(adminExt);
+ int successfulGroupQueries = 0;
for (String group : groups) {
if (isSystemGroup(group)) {
continue;
}
try {
ConsumerConnection consumerConnection =
adminExt.examineConsumerConnectionInfo(group);
+ successfulGroupQueries++;
if (consumerConnection == null ||
consumerConnection.getConnectionSet() == null) {
continue;
}
@@ -231,6 +233,9 @@ public class RocketMQClientProvider implements
ClientProvider {
log.warn("Failed to examine consumer connection for group={},
skipping", group, e);
}
}
+ if (!groups.isEmpty() && successfulGroupQueries == 0) {
+ throw new BusinessException(502, "Failed to query consumer
connections from all groups");
+ }
return result;
}
@@ -241,11 +246,14 @@ public class RocketMQClientProvider implements
ClientProvider {
clusterInfo = adminExt.examineBrokerClusterInfo();
} catch (Exception e) {
log.warn("Failed to fetch cluster info for consumer connection
scan", e);
- return groups;
+ throw new BusinessException(502, "Failed to discover brokers for
consumer connections: "
+ + rootMessage(e));
}
if (clusterInfo == null || clusterInfo.getBrokerAddrTable() == null) {
return groups;
}
+ int attemptedBrokerQueries = 0;
+ int successfulBrokerQueries = 0;
for (BrokerData brokerData :
clusterInfo.getBrokerAddrTable().values()) {
if (brokerData == null) {
continue;
@@ -255,8 +263,10 @@ public class RocketMQClientProvider implements
ClientProvider {
continue;
}
try {
+ attemptedBrokerQueries++;
SubscriptionGroupWrapper wrapper =
adminExt.getAllSubscriptionGroup(brokerAddr,
SUBSCRIPTION_GROUP_TIMEOUT_MILLIS);
+ successfulBrokerQueries++;
if (wrapper != null && wrapper.getSubscriptionGroupTable() !=
null) {
groups.addAll(wrapper.getSubscriptionGroupTable().keySet());
}
@@ -264,6 +274,9 @@ public class RocketMQClientProvider implements
ClientProvider {
log.warn("Failed to fetch subscription groups from broker={},
skipping", brokerAddr, e);
}
}
+ if (attemptedBrokerQueries > 0 && successfulBrokerQueries == 0) {
+ throw new BusinessException(502, "Failed to query subscription
groups from all brokers");
+ }
return groups;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
index 06d43089..f63a1128 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
@@ -29,6 +29,7 @@ import org.apache.rocketmq.remoting.protocol.route.QueueData;
import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.tools.admin.MQAdminExt;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
@@ -243,7 +244,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
return routes;
} catch (Exception e) {
log.warn("Failed to get routes for topic {}: {}", name,
e.getMessage());
- return Collections.emptyList();
+ throw new BusinessException(502, "Failed to get routes for topic "
+ name + ": " + e.getMessage());
}
}
@@ -376,7 +377,7 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
return progress;
} catch (Exception e) {
log.warn("Failed to get progress for group {}: {}", name,
e.getMessage());
- return Collections.emptyList();
+ throw new BusinessException(502, "Failed to get progress for group
" + name + ": " + e.getMessage());
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
index 25becb0e..8993dfa1 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClientProviderTest.java
@@ -49,6 +49,7 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.verify;
@@ -222,6 +223,16 @@ class RocketMQClientProviderTest {
verify(adminExt).getAllSubscriptionGroup("127.0.0.1:10911", 5000L);
}
+ @Test
+ void consumerScanFailsWhenBrokerDiscoveryFails() throws Exception {
+ when(adminExt.examineBrokerClusterInfo()).thenThrow(new
IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() -> provider.findConnections("instance-a",
"cluster-a", "Consumer"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to discover brokers for consumer
connections: broker unavailable")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
@Test
void consumerScanSkipsNullConnectionEntries() throws Exception {
ClusterInfo clusterInfo = new ClusterInfo();
@@ -246,6 +257,62 @@ class RocketMQClientProviderTest {
assertThat(connections.get(0).getClientId()).isEqualTo("consumer-client");
}
+ @Test
+ void consumerScanFailsWhenEveryBrokerGroupQueryFails() throws Exception {
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo(
+ "127.0.0.1:10911", "127.0.0.2:10911"));
+ when(adminExt.getAllSubscriptionGroup(anyString(),
anyLong())).thenThrow(
+ new IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() -> provider.findConnections("instance-a",
"cluster-a", "Consumer"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to query subscription groups from all
brokers")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
+ @Test
+ void consumerScanFailsWhenEveryGroupConnectionQueryFails() throws
Exception {
+ SubscriptionGroupWrapper wrapper = subscriptionGroups("group-a",
"group-b");
+
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo("127.0.0.1:10911"));
+ when(adminExt.getAllSubscriptionGroup("127.0.0.1:10911",
5000L)).thenReturn(wrapper);
+ when(adminExt.examineConsumerConnectionInfo(anyString()))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() -> provider.findConnections("instance-a",
"cluster-a", "Consumer"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to query consumer connections from all
groups")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
+ @Test
+ void consumerScanReturnsPartialResultsWhenOneGroupQueryFails() throws
Exception {
+ SubscriptionGroupWrapper wrapper = subscriptionGroups("group-a",
"group-b");
+ ConsumerConnection consumerConnection = new ConsumerConnection();
+ consumerConnection.setConnectionSet(new
HashSet<>(List.of(connection("consumer-client", "10.0.0.2:1000"))));
+
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo("127.0.0.1:10911"));
+ when(adminExt.getAllSubscriptionGroup("127.0.0.1:10911",
5000L)).thenReturn(wrapper);
+ when(adminExt.examineConsumerConnectionInfo("group-a"))
+ .thenThrow(new IllegalStateException("broker unavailable"));
+
when(adminExt.examineConsumerConnectionInfo("group-b")).thenReturn(consumerConnection);
+
+ List<ClientConnectionVO> connections =
provider.findConnections("instance-a", "cluster-a", "Consumer");
+
+ assertThat(connections).singleElement().satisfies(connection -> {
+ assertThat(connection.getClientId()).isEqualTo("consumer-client");
+ assertThat(connection.getGroupOrTopic()).isEqualTo("group-b");
+ });
+ }
+
+ private static SubscriptionGroupWrapper subscriptionGroups(String...
names) {
+ SubscriptionGroupWrapper wrapper = new SubscriptionGroupWrapper();
+ ConcurrentHashMap<String, SubscriptionGroupConfig> groups = new
ConcurrentHashMap<>();
+ for (String name : names) {
+ groups.put(name, new SubscriptionGroupConfig());
+ }
+ wrapper.setSubscriptionGroupTable(groups);
+ return wrapper;
+ }
+
private static Connection connection(String clientId, String clientAddr) {
Connection connection = new Connection();
connection.setClientId(clientId);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
index 90d4607c..19459427 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
@@ -18,6 +18,9 @@ package org.apache.rocketmq.studio.provider.apache;
import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
+import org.apache.rocketmq.tools.admin.MQAdminExt;
import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
@@ -35,7 +38,10 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.lenient;
+import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
@@ -129,4 +135,36 @@ class RocketMQMetadataProviderTest {
.containsExactlyElementsOf(subscriptions);
verify(runtimeAdminClientResolver,
org.mockito.Mockito.times(2)).execute(eq("instance-a"), any());
}
+
+ @Test
+ void getTopicRoutesSurfacesAdminFailure() throws Exception {
+ DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ when(admin.examineTopicRouteInfo("TopicA")).thenThrow(new
IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() -> newLiveProvider(admin).getTopicRoutes(null,
"TopicA"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to get routes for topic TopicA: broker
unavailable")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
+ @Test
+ void getGroupProgressSurfacesAdminFailure() throws Exception {
+ DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ when(admin.examineConsumeStats("group-a")).thenThrow(new
IllegalStateException("broker unavailable"));
+
+ assertThatThrownBy(() -> newLiveProvider(admin).getGroupProgress(null,
"group-a"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Failed to get progress for group group-a: broker
unavailable")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
+ }
+
+ private RocketMQMetadataProvider newLiveProvider(MQAdminExt admin) throws
Exception {
+ MqAdminExtFactory factory = mock(MqAdminExtFactory.class);
+ RocketMQProperties liveProperties = new RocketMQProperties();
+ liveProperties.setNamesrvAddr("10.0.0.1:9876");
+ lenient().when(factory.execute(anyString(), any(),
any())).thenAnswer(invocation ->
+
invocation.<MqAdminExtFactory.AdminAction<Object>>getArgument(2).apply(admin));
+ return new RocketMQMetadataProvider(factory, liveProperties,
topicMapper, groupMapper,
+ runtimeAdminClientResolver);
+ }
}