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 dc49e1a2 feat: manage topics, groups, messages and DLQ against a live
cluster (#798)
dc49e1a2 is described below
commit dc49e1a24f24053ca546bc4aa82b504b10ef2548
Author: lizhimins <[email protected]>
AuthorDate: Mon Aug 3 14:32:48 2026 +0800
feat: manage topics, groups, messages and DLQ against a live cluster (#798)
---
.../studio/instance/dlq/DLQProviderStub.java | 6 +-
.../instance/message/MessageProviderStub.java | 6 +-
.../instance/topic/CloudMetadataProvider.java | 5 +-
.../studio/instance/topic/NameSrvAdminClient.java | 5 +-
.../studio/queryhistory/QueryHistoryService.java | 68 +++
.../studio/rocketmq/RocketMQDLQProvider.java | 291 ++++++++++++
.../studio/rocketmq/RocketMQMessageProvider.java | 454 ++++++++++++++++++
.../studio/rocketmq/RocketMQMetadataProvider.java | 514 +++++++++++++++++++++
8 files changed, 1341 insertions(+), 8 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
index 34558993..9a301a3b 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
@@ -18,11 +18,13 @@ package org.apache.rocketmq.studio.instance.dlq;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import lombok.extern.slf4j.Slf4j;
-import org.springframework.stereotype.Component;
import java.util.List;
-@Component
+/**
+ * Fallback {@link DLQProvider} used only by tests; not registered as a Spring
bean. The live
+ * implementation is {@code RocketMQDLQProvider}.
+ */
@Slf4j
public class DLQProviderStub implements DLQProvider {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
index f99882ab..5c2f8f0c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProviderStub.java
@@ -18,11 +18,13 @@ package org.apache.rocketmq.studio.instance.message;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import lombok.extern.slf4j.Slf4j;
-import org.springframework.stereotype.Component;
import java.util.List;
-@Component
+/**
+ * Fallback {@link MessageProvider} used only by tests; not registered as a
Spring bean. The live
+ * implementation is {@code RocketMQMessageProvider}.
+ */
@Slf4j
public class MessageProviderStub implements MessageProvider {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/CloudMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/CloudMetadataProvider.java
index c9a58504..a0053447 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/CloudMetadataProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/CloudMetadataProvider.java
@@ -21,12 +21,13 @@ import
org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
import lombok.extern.slf4j.Slf4j;
-import org.springframework.stereotype.Component;
import java.util.List;
+/**
+ * Stub metadata provider. Disabled in favor of RocketMQMetadataProvider.
+ */
@Slf4j
-@Component
public class CloudMetadataProvider implements MetadataProvider {
@Override
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/NameSrvAdminClient.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/NameSrvAdminClient.java
index 15f9df57..60f05a34 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/NameSrvAdminClient.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/NameSrvAdminClient.java
@@ -19,10 +19,11 @@ package org.apache.rocketmq.studio.instance.topic;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import lombok.extern.slf4j.Slf4j;
-import org.springframework.stereotype.Component;
+/**
+ * Stub admin client. Disabled in favor of RocketMQAdminClientImpl.
+ */
@Slf4j
-@Component
public class NameSrvAdminClient implements AdminClient {
@Override
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryService.java
b/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryService.java
new file mode 100644
index 00000000..18e512aa
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/queryhistory/QueryHistoryService.java
@@ -0,0 +1,68 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.queryhistory;
+
+import lombok.extern.slf4j.Slf4j;
+import org.apache.rocketmq.studio.persistence.entity.RmqMessageQuery;
+import org.apache.rocketmq.studio.persistence.entity.RmqTraceQuery;
+import org.apache.rocketmq.studio.persistence.mapper.RmqMessageQueryMapper;
+import org.apache.rocketmq.studio.persistence.mapper.RmqTraceQueryMapper;
+import org.springframework.stereotype.Service;
+
+import java.time.LocalDateTime;
+
+@Slf4j
+@Service
+public class QueryHistoryService {
+
+ private final RmqMessageQueryMapper messageQueryMapper;
+ private final RmqTraceQueryMapper traceQueryMapper;
+
+ public QueryHistoryService(RmqMessageQueryMapper messageQueryMapper,
+ RmqTraceQueryMapper traceQueryMapper) {
+ this.messageQueryMapper = messageQueryMapper;
+ this.traceQueryMapper = traceQueryMapper;
+ }
+
+ public void recordMessageQuery(String queryType, String topic, String
msgId,
+ String tag, String key, Long startTime,
+ Long endTime, int resultCount) {
+ RmqMessageQuery query = new RmqMessageQuery();
+ query.setQueryType(queryType);
+ query.setTopic(topic);
+ query.setMsgId(msgId);
+ query.setTag(tag);
+ query.setMessageKey(key);
+ query.setStartTime(startTime);
+ query.setEndTime(endTime);
+ query.setResultCount(resultCount);
+ query.setQueriedAt(LocalDateTime.now());
+ messageQueryMapper.insert(query);
+ log.debug("Message query recorded: type={} topic={}", queryType,
topic);
+ }
+
+ public void recordTraceQuery(String msgId, String topic, int nodeCount,
int consumerCount) {
+ RmqTraceQuery query = new RmqTraceQuery();
+ query.setMsgId(msgId);
+ query.setTopic(topic);
+ query.setNodeCount(nodeCount);
+ query.setConsumerCount(consumerCount);
+ query.setQueriedAt(LocalDateTime.now());
+ traceQueryMapper.insert(query);
+ log.debug("Trace query recorded: msgId={} topic={}", msgId, topic);
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
new file mode 100644
index 00000000..e9d1a1d9
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
@@ -0,0 +1,291 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.rocketmq;
+
+import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
+import org.apache.rocketmq.client.consumer.PullResult;
+import org.apache.rocketmq.client.consumer.PullStatus;
+import org.apache.rocketmq.client.producer.DefaultMQProducer;
+import org.apache.rocketmq.client.producer.SendResult;
+import org.apache.rocketmq.common.MixAll;
+import org.apache.rocketmq.common.message.Message;
+import org.apache.rocketmq.common.message.MessageConst;
+import org.apache.rocketmq.common.message.MessageExt;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
+import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
+import org.apache.rocketmq.remoting.protocol.body.TopicList;
+import org.apache.rocketmq.studio.instance.dlq.DLQGroupVO;
+import org.apache.rocketmq.studio.instance.dlq.DLQProvider;
+import org.apache.rocketmq.studio.ops.audit.AuditService;
+import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.ObjectProvider;
+import org.springframework.context.annotation.Primary;
+import org.springframework.stereotype.Service;
+import org.springframework.util.StringUtils;
+
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * Real {@link DLQProvider} backed by the RocketMQ admin API. Lists
dead-letter groups by scanning
+ * {@code %DLQ%} topics and resends dead-letter messages back to a target
topic.
+ */
+@Service
+@Primary
+public class RocketMQDLQProvider implements DLQProvider {
+
+ private static final Logger log =
LoggerFactory.getLogger(RocketMQDLQProvider.class);
+
+ private static final long ONE_HOUR_MILLIS = 3600_000L;
+ private static final int RESEND_HARD_CAP = 5000;
+
+ private final ObjectProvider<DefaultMQAdminExt> adminExtProvider;
+ private final AuditService auditService;
+ private final RocketMQProperties properties;
+
+ public RocketMQDLQProvider(ObjectProvider<DefaultMQAdminExt>
adminExtProvider,
+ AuditService auditService,
+ RocketMQProperties properties) {
+ this.adminExtProvider = adminExtProvider;
+ this.auditService = auditService;
+ this.properties = properties;
+ }
+
+ @Override
+ public List<DLQGroupVO> listDLQGroups(String clusterId) {
+ DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
+ if (adminExt == null) {
+ log.warn("DefaultMQAdminExt is not configured, returning empty DLQ
group list");
+ return Collections.emptyList();
+ }
+
+ Set<String> topics;
+ try {
+ TopicList topicList = adminExt.fetchAllTopicList();
+ topics = topicList == null ? Collections.emptySet() :
topicList.getTopicList();
+ } catch (Exception e) {
+ log.warn("Failed to fetch topic list for DLQ scan: {}",
e.getMessage());
+ return Collections.emptyList();
+ }
+
+ List<DLQGroupVO> groups = new ArrayList<>();
+ for (String topic : topics) {
+ if (topic == null ||
!topic.startsWith(MixAll.DLQ_GROUP_TOPIC_PREFIX)) {
+ continue;
+ }
+ String groupName =
topic.substring(MixAll.DLQ_GROUP_TOPIC_PREFIX.length());
+ if (!StringUtils.hasText(groupName)) {
+ continue;
+ }
+ groups.add(buildDLQGroup(adminExt, groupName, topic));
+ }
+ return groups;
+ }
+
+ private DLQGroupVO buildDLQGroup(DefaultMQAdminExt adminExt, String
groupName, String dlqTopic) {
+ long messageCount = 0L;
+ LocalDateTime lastEnqueueTime = null;
+ try {
+ TopicStatsTable statsTable = adminExt.examineTopicStats(dlqTopic);
+ if (statsTable != null && statsTable.getOffsetTable() != null) {
+ long latestUpdate = 0L;
+ for (Map.Entry<MessageQueue, TopicOffset> entry :
statsTable.getOffsetTable().entrySet()) {
+ TopicOffset offset = entry.getValue();
+ if (offset == null) {
+ continue;
+ }
+ messageCount += Math.max(0L, offset.getMaxOffset() -
offset.getMinOffset());
+ if (offset.getLastUpdateTimestamp() > latestUpdate) {
+ latestUpdate = offset.getLastUpdateTimestamp();
+ }
+ }
+ if (latestUpdate > 0L) {
+ lastEnqueueTime = LocalDateTime.ofInstant(
+ Instant.ofEpochMilli(latestUpdate),
ZoneId.systemDefault());
+ }
+ }
+ } catch (Exception e) {
+ log.warn("Failed to examine stats for DLQ topic {}: {}", dlqTopic,
e.getMessage());
+ }
+
+ return DLQGroupVO.builder()
+ .groupName(groupName)
+ .dlqTopic(dlqTopic)
+ .messageCount(messageCount)
+ .lastEnqueueTime(lastEnqueueTime)
+ .retryCount(0)
+ .status(messageCount > 0 ? "ACTIVE" : "EMPTY")
+ .build();
+ }
+
+ @Override
+ public void resendMessages(String groupName, Long startTime, Long endTime,
String targetTopic) {
+ DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
+
+ long end = endTime != null ? endTime : System.currentTimeMillis();
+ long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
+
+ List<MessageExt> deadLetters = collectDeadLetters(dlqTopic, begin,
end);
+ int resent = 0;
+ int failed = 0;
+ if (!deadLetters.isEmpty()) {
+ DefaultMQProducer producer = newProducer(groupName);
+ try {
+ producer.start();
+ for (MessageExt deadLetter : deadLetters) {
+ if (resendOne(producer, deadLetter, targetTopic)) {
+ resent++;
+ } else {
+ failed++;
+ }
+ }
+ } catch (Exception e) {
+ log.warn("Failed to start resend producer for group {}: {}",
groupName, e.getMessage());
+ failed += deadLetters.size() - resent;
+ } finally {
+ producer.shutdown();
+ }
+ }
+
+ String detail = String.format("group=%s, dlqTopic=%s, targetTopic=%s,
matched=%d, resent=%d, failed=%d",
+ groupName, dlqTopic, StringUtils.hasText(targetTopic) ?
targetTopic : "<original>",
+ deadLetters.size(), resent, failed);
+ recordAudit(groupName, detail, failed == 0 ? "SUCCESS" : "PARTIAL");
+ log.info("DLQ resend completed: {}", detail);
+ }
+
+ private List<MessageExt> collectDeadLetters(String dlqTopic, long begin,
long end) {
+ DefaultMQPullConsumer consumer = newPullConsumer();
+ List<MessageExt> result = new ArrayList<>();
+ try {
+ consumer.start();
+ Set<MessageQueue> queues =
consumer.fetchSubscribeMessageQueues(dlqTopic);
+ outer:
+ for (MessageQueue queue : queues) {
+ long minOffset = consumer.searchOffset(queue, begin);
+ long maxOffset = consumer.searchOffset(queue, end);
+ for (long offset = minOffset; offset <= maxOffset; ) {
+ if (result.size() >= RESEND_HARD_CAP) {
+ break outer;
+ }
+ PullResult pullResult = consumer.pull(queue, "*", offset,
32);
+ offset = pullResult.getNextBeginOffset();
+ if (pullResult.getPullStatus() != PullStatus.FOUND
+ || pullResult.getMsgFoundList() == null) {
+ break;
+ }
+ for (MessageExt messageExt : pullResult.getMsgFoundList())
{
+ if (messageExt.getStoreTimestamp() >= begin
+ && messageExt.getStoreTimestamp() <= end) {
+ result.add(messageExt);
+ if (result.size() >= RESEND_HARD_CAP) {
+ break outer;
+ }
+ }
+ }
+ }
+ }
+ } catch (Exception e) {
+ log.warn("Failed to collect dead letters from {}: {}", dlqTopic,
e.getMessage());
+ } finally {
+ consumer.shutdown();
+ }
+ return result;
+ }
+
+ private boolean resendOne(DefaultMQProducer producer, MessageExt
deadLetter, String targetTopic) {
+ String destination = resolveTargetTopic(deadLetter, targetTopic);
+ if (!StringUtils.hasText(destination)) {
+ log.warn("Skip resend of msgId={}: no target topic resolvable",
deadLetter.getMsgId());
+ return false;
+ }
+ try {
+ Message message = new Message(destination, deadLetter.getBody());
+ if (StringUtils.hasText(deadLetter.getTags())) {
+ message.setTags(deadLetter.getTags());
+ }
+ if (StringUtils.hasText(deadLetter.getKeys())) {
+ message.setKeys(deadLetter.getKeys());
+ }
+ message.putUserProperty("DLQ_ORIGIN_MESSAGE_ID",
deadLetter.getMsgId());
+ message.putUserProperty("DLQ_ORIGIN_TOPIC", deadLetter.getTopic());
+ SendResult sendResult = producer.send(message);
+ log.debug("Resent dead letter msgId={} to topic={}, sendStatus={}",
+ deadLetter.getMsgId(), destination,
sendResult.getSendStatus());
+ return true;
+ } catch (Exception e) {
+ log.warn("Failed to resend dead letter msgId={} to topic={}: {}",
+ deadLetter.getMsgId(), destination, e.getMessage());
+ return false;
+ }
+ }
+
+ private String resolveTargetTopic(MessageExt deadLetter, String
targetTopic) {
+ if (StringUtils.hasText(targetTopic)) {
+ return targetTopic;
+ }
+ Map<String, String> properties = deadLetter.getProperties();
+ if (properties != null) {
+ String origin =
properties.get(MessageConst.PROPERTY_DLQ_ORIGIN_TOPIC);
+ if (StringUtils.hasText(origin)) {
+ return origin;
+ }
+ String retryTopic =
properties.get(MessageConst.PROPERTY_RETRY_TOPIC);
+ if (StringUtils.hasText(retryTopic)) {
+ return retryTopic;
+ }
+ }
+ return null;
+ }
+
+ private DefaultMQPullConsumer newPullConsumer() {
+ DefaultMQPullConsumer consumer = new
DefaultMQPullConsumer("studio-dlq-query-group");
+ consumer.setInstanceName("studio-dlq-query-" +
System.currentTimeMillis());
+ if (StringUtils.hasText(properties.getNamesrvAddr())) {
+ consumer.setNamesrvAddr(properties.getNamesrvAddr());
+ }
+ return consumer;
+ }
+
+ private DefaultMQProducer newProducer(String groupName) {
+ DefaultMQProducer producer = new
DefaultMQProducer("studio-dlq-resend-" + groupName);
+ producer.setInstanceName("studio-dlq-resend-" +
System.currentTimeMillis());
+ producer.setRetryTimesWhenSendFailed(2);
+ if (StringUtils.hasText(properties.getNamesrvAddr())) {
+ producer.setNamesrvAddr(properties.getNamesrvAddr());
+ }
+ return producer;
+ }
+
+ private void recordAudit(String groupName, String detail, String result) {
+ try {
+ auditService.record("RESEND_DLQ", groupName, detail, result);
+ } catch (Exception e) {
+ log.warn("Failed to record DLQ resend audit: {}", e.getMessage());
+ }
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
new file mode 100644
index 00000000..27c73672
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
@@ -0,0 +1,454 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.rocketmq;
+
+import org.apache.rocketmq.client.QueryResult;
+import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
+import org.apache.rocketmq.client.consumer.PullResult;
+import org.apache.rocketmq.client.consumer.PullStatus;
+import org.apache.rocketmq.common.message.MessageDecoder;
+import org.apache.rocketmq.common.message.MessageExt;
+import org.apache.rocketmq.common.message.MessageId;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
+import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
+import org.apache.rocketmq.studio.instance.message.MessageProvider;
+import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
+import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
+import org.apache.rocketmq.studio.queryhistory.QueryHistoryService;
+import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
+import org.apache.rocketmq.tools.admin.api.MessageTrack;
+import org.apache.rocketmq.tools.admin.api.TrackType;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.ObjectProvider;
+import org.springframework.context.annotation.Primary;
+import org.springframework.stereotype.Service;
+import org.springframework.util.StringUtils;
+
+import java.net.InetSocketAddress;
+import java.net.SocketAddress;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Set;
+
+/**
+ * Real {@link MessageProvider} backed by the RocketMQ admin API. Supports
message lookup by
+ * message id, by business key and by topic + time range, as well as message
trace retrieval.
+ * Falls back to empty results when adminExt is not configured or a query
fails.
+ */
+@Service
+@Primary
+public class RocketMQMessageProvider implements MessageProvider {
+
+ private static final Logger log =
LoggerFactory.getLogger(RocketMQMessageProvider.class);
+
+ private static final String TRACE_TOPIC = "RMQ_SYS_TRACE_TOPIC";
+ private static final char FIELD_SEPARATOR = '\u0001';
+ private static final int KEY_QUERY_MAX = 64;
+ private static final int TRACE_QUERY_MAX = 64;
+ private static final int DEFAULT_TOPIC_LIMIT = 200;
+ private static final int TOPIC_QUERY_HARD_CAP = 2000;
+ private static final long ONE_HOUR_MILLIS = 3600_000L;
+ private static final long ONE_DAY_MILLIS = 24 * ONE_HOUR_MILLIS;
+
+ private final ObjectProvider<DefaultMQAdminExt> adminExtProvider;
+ private final QueryHistoryService queryHistoryService;
+ private final RocketMQProperties properties;
+
+ public RocketMQMessageProvider(ObjectProvider<DefaultMQAdminExt>
adminExtProvider,
+ QueryHistoryService queryHistoryService,
+ RocketMQProperties properties) {
+ this.adminExtProvider = adminExtProvider;
+ this.queryHistoryService = queryHistoryService;
+ this.properties = properties;
+ }
+
+ @Override
+ public List<MessageRecordVO> queryMessages(String topic, String msgId,
String tag, String key,
+ Long startTime, Long endTime) {
+ DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
+ if (adminExt == null) {
+ log.warn("DefaultMQAdminExt is not configured, returning empty
message list");
+ return Collections.emptyList();
+ }
+
+ long end = endTime != null ? endTime : System.currentTimeMillis();
+ long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
+
+ List<MessageRecordVO> result;
+ String queryType;
+ if (StringUtils.hasText(msgId)) {
+ queryType = "MSG_ID";
+ result = queryByMsgId(adminExt, topic, msgId);
+ } else if (StringUtils.hasText(topic) && StringUtils.hasText(key)) {
+ queryType = "KEY";
+ result = queryByKey(adminExt, topic, key, tag, begin, end);
+ } else if (StringUtils.hasText(topic)) {
+ queryType = "TOPIC";
+ result = queryByTopic(topic, tag, begin, end, DEFAULT_TOPIC_LIMIT);
+ } else {
+ log.warn("queryMessages requires at least one of msgId/topic,
returning empty list");
+ return Collections.emptyList();
+ }
+
+ recordMessageQuery(queryType, topic, msgId, tag, key, startTime,
endTime, result.size());
+ return result;
+ }
+
+ private List<MessageRecordVO> queryByMsgId(DefaultMQAdminExt adminExt,
String topic, String msgId) {
+ MessageExt messageExt = null;
+ if (StringUtils.hasText(topic)) {
+ try {
+ messageExt = adminExt.viewMessage(topic, msgId);
+ } catch (Exception e) {
+ log.warn("viewMessage(topic={}, msgId={}) failed: {}", topic,
msgId, e.getMessage());
+ }
+ }
+ if (messageExt == null) {
+ messageExt = viewMessageByOffsetId(adminExt, msgId);
+ }
+ if (messageExt == null) {
+ return Collections.emptyList();
+ }
+ return Collections.singletonList(toRecordVO(messageExt));
+ }
+
+ /**
+ * Locate a message purely by its offset msgId by decoding the broker
address embedded in the
+ * id and querying that broker directly. Used when no topic hint is
available.
+ */
+ private MessageExt viewMessageByOffsetId(DefaultMQAdminExt adminExt,
String msgId) {
+ try {
+ MessageId messageId = MessageDecoder.decodeMessageId(msgId);
+ SocketAddress address = messageId.getAddress();
+ if (!(address instanceof InetSocketAddress)) {
+ return null;
+ }
+ InetSocketAddress inet = (InetSocketAddress) address;
+ String brokerAddr = inet.getAddress().getHostAddress() + ":" +
inet.getPort();
+ return adminExt.getDefaultMQAdminExtImpl()
+ .getMqClientInstance()
+ .getMQClientAPIImpl()
+ .viewMessage(brokerAddr, msgId, 3000L, 3000L);
+ } catch (Exception e) {
+ log.warn("viewMessage by decoded offset id failed for msgId={}:
{}", msgId, e.getMessage());
+ return null;
+ }
+ }
+
+ private List<MessageRecordVO> queryByKey(DefaultMQAdminExt adminExt,
String topic, String key,
+ String tag, long begin, long end)
{
+ try {
+ QueryResult queryResult = adminExt.queryMessage(topic, key,
KEY_QUERY_MAX, begin, end);
+ if (queryResult == null || queryResult.getMessageList() == null) {
+ return Collections.emptyList();
+ }
+ List<MessageRecordVO> result = new ArrayList<>();
+ for (MessageExt messageExt : queryResult.getMessageList()) {
+ if (matchesTag(messageExt, tag)) {
+ result.add(toRecordVO(messageExt));
+ }
+ }
+ return result;
+ } catch (Exception e) {
+ log.warn("queryMessage(topic={}, key={}) failed: {}", topic, key,
e.getMessage());
+ return Collections.emptyList();
+ }
+ }
+
+ /**
+ * Scan a topic within a time range using a short-lived pull consumer,
mirroring the approach
+ * used by the RocketMQ dashboard for time-range topic queries.
+ */
+ private List<MessageRecordVO> queryByTopic(String topic, String tag, long
begin, long end, int limit) {
+ DefaultMQPullConsumer consumer = newPullConsumer("studio-msg-query");
+ List<MessageRecordVO> result = new ArrayList<>();
+ try {
+ consumer.start();
+ Set<MessageQueue> queues =
consumer.fetchSubscribeMessageQueues(topic);
+ outer:
+ for (MessageQueue queue : queues) {
+ long minOffset = consumer.searchOffset(queue, begin);
+ long maxOffset = consumer.searchOffset(queue, end);
+ for (long offset = minOffset; offset <= maxOffset; ) {
+ if (result.size() >= Math.min(limit,
TOPIC_QUERY_HARD_CAP)) {
+ break outer;
+ }
+ PullResult pullResult = consumer.pull(queue, "*", offset,
32);
+ offset = pullResult.getNextBeginOffset();
+ if (pullResult.getPullStatus() != PullStatus.FOUND
+ || pullResult.getMsgFoundList() == null) {
+ break;
+ }
+ for (MessageExt messageExt : pullResult.getMsgFoundList())
{
+ if (messageExt.getStoreTimestamp() < begin
+ || messageExt.getStoreTimestamp() > end) {
+ continue;
+ }
+ if (!matchesTag(messageExt, tag)) {
+ continue;
+ }
+ result.add(toRecordVO(messageExt));
+ if (result.size() >= Math.min(limit,
TOPIC_QUERY_HARD_CAP)) {
+ break outer;
+ }
+ }
+ }
+ }
+ } catch (Exception e) {
+ log.warn("queryByTopic(topic={}) failed: {}", topic,
e.getMessage());
+ } finally {
+ consumer.shutdown();
+ }
+ return result;
+ }
+
+ @Override
+ public TraceRecordVO getMessageTrace(String msgId) {
+ DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
+ if (adminExt == null) {
+ log.warn("DefaultMQAdminExt is not configured, returning empty
trace");
+ return emptyTrace();
+ }
+
+ long now = System.currentTimeMillis();
+ long begin = now - ONE_HOUR_MILLIS;
+ long end = now + 60_000L;
+
+ List<TraceNodeVO> nodes = new ArrayList<>();
+ List<ConsumerStatusVO> consumerStatus = new ArrayList<>();
+
+ try {
+ QueryResult traceResult = adminExt.queryMessage(TRACE_TOPIC,
msgId, TRACE_QUERY_MAX, begin, end);
+ if (traceResult != null && traceResult.getMessageList() != null) {
+ for (MessageExt traceMessage : traceResult.getMessageList()) {
+ parseTraceBody(traceMessage.getBody(), msgId, nodes,
consumerStatus);
+ }
+ }
+ } catch (Exception e) {
+ log.warn("Trace query for msgId={} failed: {}", msgId,
e.getMessage());
+ }
+
+ recordTraceQuery(msgId, null, nodes.size(), consumerStatus.size());
+ return TraceRecordVO.builder()
+ .nodes(nodes)
+ .consumerStatus(consumerStatus)
+ .build();
+ }
+
+ /**
+ * Parse a trace message body. Each line is one trace context whose fields
are separated by
+ * the SOH character ({@code \u0001}); the first field is the trace type.
+ */
+ private void parseTraceBody(byte[] body, String targetMsgId,
List<TraceNodeVO> nodes,
+ List<ConsumerStatusVO> consumerStatus) {
+ if (body == null || body.length == 0) {
+ return;
+ }
+ String data = new String(body, StandardCharsets.UTF_8);
+ for (String line : data.split("\n")) {
+ if (!StringUtils.hasText(line)) {
+ continue;
+ }
+ String[] fields = line.split(String.valueOf(FIELD_SEPARATOR), -1);
+ if (fields.length == 0) {
+ continue;
+ }
+ try {
+ switch (fields[0].trim()) {
+ case "Pub":
+ nodes.add(buildProduceNode(fields));
+ break;
+ case "SubAfter":
+ nodes.add(buildConsumeNode(fields));
+ consumerStatus.add(buildConsumerStatus(fields));
+ break;
+ case "EndTransaction":
+ nodes.add(buildTransactionNode(fields));
+ break;
+ default:
+ // SubBefore and unknown types are not surfaced as
timeline nodes.
+ break;
+ }
+ } catch (Exception e) {
+ log.debug("Skipping unparseable trace line: {}",
e.getMessage());
+ }
+ }
+ }
+
+ // Pub layout: type, time, region, group, topic, msgId, tags, keys,
storeHost, clientHost,
+ // retryTimes, msgType, status, costTime
+ private TraceNodeVO buildProduceNode(String[] f) {
+ return TraceNodeVO.builder()
+ .title("produce")
+ .timestamp(parseLong(field(f, 1)))
+ .status(parseBoolean(field(f, 12)) ? "finish" : "failed")
+ .costTime(parseLong(field(f, 13)))
+ .description("producer=" + field(f, 3) + ", clientHost=" +
field(f, 9)
+ + ", storeHost=" + field(f, 8))
+ .build();
+ }
+
+ // SubAfter layout: type, time, region, group, traceId, msgId, retryTimes,
keys, costTime,
+ // status, timestamp
+ private TraceNodeVO buildConsumeNode(String[] f) {
+ return TraceNodeVO.builder()
+ .title("consume")
+ .timestamp(parseLong(field(f, 1)))
+ .status(parseBoolean(field(f, 9)) ? "finish" : "failed")
+ .costTime(parseLong(field(f, 8)))
+ .description("group=" + field(f, 3) + ", retryTimes=" +
field(f, 6))
+ .build();
+ }
+
+ private ConsumerStatusVO buildConsumerStatus(String[] f) {
+ return ConsumerStatusVO.builder()
+ .group(field(f, 3))
+ .deliveryStatus(parseBoolean(field(f, 9)) ?
DeliveryStatus.success : DeliveryStatus.failed)
+ .consumeTime(parseLong(field(f, 1)))
+ .retryCount((int) parseLong(field(f, 6)))
+ .build();
+ }
+
+ // EndTransaction layout: type, time, region, group, topic, msgId, tags,
keys, storeHost,
+ // clientHost, retryTimes, msgType, status,
transactionId, txState
+ private TraceNodeVO buildTransactionNode(String[] f) {
+ return TraceNodeVO.builder()
+ .title("endTransaction")
+ .timestamp(parseLong(field(f, 1)))
+ .status("finish")
+ .costTime(0L)
+ .description("group=" + field(f, 3) + ", transactionState=" +
field(f, 14))
+ .build();
+ }
+
+ private List<ConsumerStatusVO> fallbackConsumerStatus(DefaultMQAdminExt
adminExt, MessageExt message) {
+ List<ConsumerStatusVO> result = new ArrayList<>();
+ try {
+ List<MessageTrack> tracks = adminExt.messageTrackDetail(message);
+ if (tracks == null) {
+ return result;
+ }
+ for (MessageTrack track : tracks) {
+ result.add(ConsumerStatusVO.builder()
+ .group(track.getConsumerGroup())
+ .deliveryStatus(mapTrackType(track.getTrackType()))
+ .consumeTime(0L)
+ .retryCount(0)
+ .build());
+ }
+ } catch (Exception e) {
+ log.warn("messageTrackDetail fallback failed for msgId={}: {}",
message.getMsgId(), e.getMessage());
+ }
+ return result;
+ }
+
+ private DeliveryStatus mapTrackType(TrackType trackType) {
+ if (trackType == null) {
+ return DeliveryStatus.pending;
+ }
+ switch (trackType) {
+ case CONSUMED:
+ case CONSUME_BROADCASTING:
+ case CONSUMED_BUT_FILTERED:
+ return DeliveryStatus.success;
+ case NOT_CONSUME_YET:
+ case PULL:
+ case NOT_ONLINE:
+ return DeliveryStatus.pending;
+ default:
+ return DeliveryStatus.failed;
+ }
+ }
+
+ private MessageRecordVO toRecordVO(MessageExt messageExt) {
+ byte[] body = messageExt.getBody();
+ return MessageRecordVO.builder()
+ .msgId(messageExt.getMsgId())
+ .topic(messageExt.getTopic())
+ .tag(messageExt.getTags())
+ .key(messageExt.getKeys())
+ .body(body == null ? null : new String(body,
StandardCharsets.UTF_8))
+ .storeTime(messageExt.getStoreTimestamp())
+ .bornHost(String.valueOf(messageExt.getBornHost()))
+ .storeHost(String.valueOf(messageExt.getStoreHost()))
+ .properties(messageExt.getProperties())
+ .size(messageExt.getStoreSize())
+ .build();
+ }
+
+ private boolean matchesTag(MessageExt messageExt, String tag) {
+ if (!StringUtils.hasText(tag) || "*".equals(tag)) {
+ return true;
+ }
+ return tag.equals(messageExt.getTags());
+ }
+
+ private DefaultMQPullConsumer newPullConsumer(String groupPrefix) {
+ DefaultMQPullConsumer consumer = new DefaultMQPullConsumer(groupPrefix
+ "-group");
+ consumer.setInstanceName(groupPrefix + "-" +
System.currentTimeMillis());
+ if (StringUtils.hasText(properties.getNamesrvAddr())) {
+ consumer.setNamesrvAddr(properties.getNamesrvAddr());
+ }
+ return consumer;
+ }
+
+ private void recordMessageQuery(String queryType, String topic, String
msgId, String tag, String key,
+ Long startTime, Long endTime, int
resultCount) {
+ try {
+ queryHistoryService.recordMessageQuery(queryType, topic, msgId,
tag, key, startTime, endTime,
+ resultCount);
+ } catch (Exception e) {
+ log.warn("Failed to record message query history: {}",
e.getMessage());
+ }
+ }
+
+ private void recordTraceQuery(String msgId, String topic, int nodeCount,
int consumerCount) {
+ try {
+ queryHistoryService.recordTraceQuery(msgId, topic, nodeCount,
consumerCount);
+ } catch (Exception e) {
+ log.warn("Failed to record trace query history: {}",
e.getMessage());
+ }
+ }
+
+ private static TraceRecordVO emptyTrace() {
+ return TraceRecordVO.builder()
+ .nodes(Collections.emptyList())
+ .consumerStatus(Collections.emptyList())
+ .build();
+ }
+
+ private static String field(String[] fields, int index) {
+ return index < fields.length ? fields[index] : "";
+ }
+
+ private static long parseLong(String value) {
+ try {
+ return Long.parseLong(value.trim());
+ } catch (Exception e) {
+ return 0L;
+ }
+ }
+
+ private static boolean parseBoolean(String value) {
+ return "true".equalsIgnoreCase(value == null ? "" : value.trim());
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMetadataProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMetadataProvider.java
new file mode 100644
index 00000000..f8bcbf84
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMetadataProvider.java
@@ -0,0 +1,514 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.rocketmq;
+
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.remoting.protocol.admin.ConsumeStats;
+import org.apache.rocketmq.remoting.protocol.admin.OffsetWrapper;
+import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
+import org.apache.rocketmq.remoting.protocol.body.ConsumerConnection;
+import org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
+import org.apache.rocketmq.remoting.protocol.heartbeat.SubscriptionData;
+import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import org.apache.rocketmq.remoting.protocol.route.QueueData;
+import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
+import
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
+import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
+import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
+import org.apache.rocketmq.studio.common.domain.enums.TopicType;
+import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
+import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
+import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
+import org.apache.rocketmq.studio.instance.topic.BrokerRouteVO;
+import org.apache.rocketmq.studio.instance.topic.MetadataProvider;
+import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
+import org.apache.rocketmq.studio.instance.topic.TopicVO;
+import org.apache.rocketmq.studio.persistence.entity.RmqGroup;
+import org.apache.rocketmq.studio.persistence.entity.RmqTopic;
+import org.apache.rocketmq.studio.persistence.mapper.RmqGroupMapper;
+import org.apache.rocketmq.studio.persistence.mapper.RmqTopicMapper;
+import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.context.annotation.Primary;
+import org.springframework.stereotype.Service;
+import org.springframework.util.StringUtils;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+/**
+ * Real MetadataProvider implementation.
+ *
+ * <p>Topic and consumer group listings are served from the studio metadata
database, which is the
+ * source of record: creation writes there and reads never fall back to the
broker. Route,
+ * consumer, progress and subscription queries stay live through
DefaultMQAdminExt, so a record
+ * without a broker route surfaces as an empty route list instead of being
hidden.
+ */
+@Service
+@Primary
+public class RocketMQMetadataProvider implements MetadataProvider {
+
+ private static final Logger log =
LoggerFactory.getLogger(RocketMQMetadataProvider.class);
+
+ private static final Set<String> SYSTEM_TOPIC_PREFIXES = Set.of(
+ "RMQ_SYS_", "SCHEDULE_TOPIC_", "%RETRY%", "%DLQ%", "CID_"
+ );
+
+ private static final Set<String> SYSTEM_TOPICS = Set.of(
+ "TBW102", "SELF_TEST_TOPIC", "DefaultCluster",
"OFFSET_MOVED_EVENT",
+ "broker", "SCHEDULE_TOPIC_XXXX", "RMQ_SYS_TRANS_HALF_TOPIC",
+ "RMQ_SYS_TRACE_TOPIC", "RMQ_SYS_TRANS_OP_HALF_TOPIC"
+ );
+
+ private final DefaultMQAdminExt adminExt;
+ private final RmqTopicMapper topicMapper;
+ private final RmqGroupMapper groupMapper;
+
+ @Autowired
+ public RocketMQMetadataProvider(@Autowired(required = false)
DefaultMQAdminExt adminExt,
+ RmqTopicMapper topicMapper,
+ RmqGroupMapper groupMapper) {
+ this.adminExt = adminExt;
+ this.topicMapper = topicMapper;
+ this.groupMapper = groupMapper;
+ }
+
+ @Override
+ public List<TopicVO> listTopics(String clusterId, String type, String
search) {
+ LambdaQueryWrapper<RmqTopic> query = new LambdaQueryWrapper<RmqTopic>()
+ .eq(StringUtils.hasText(clusterId), RmqTopic::getClusterId,
clusterId)
+ .eq(StringUtils.hasText(type), RmqTopic::getTopicType, type)
+ .like(StringUtils.hasText(search), RmqTopic::getName, search)
+ .orderByAsc(RmqTopic::getName);
+
+ List<TopicVO> result = new ArrayList<>();
+ for (RmqTopic entity : topicMapper.selectList(query)) {
+ if (isSystemTopic(entity.getName(), Collections.emptySet())) {
+ continue;
+ }
+ result.add(toTopicVO(entity));
+ }
+ return result;
+ }
+
+ private TopicVO toTopicVO(RmqTopic entity) {
+ TopicVO vo = new TopicVO();
+ vo.setId(entity.getName());
+ vo.setName(entity.getName());
+ vo.setClusterId(entity.getClusterId());
+ vo.setType(parseTopicType(entity.getTopicType()));
+ vo.setReadQueues(entity.getReadQueueNums() == null ? 0 :
entity.getReadQueueNums());
+ vo.setWriteQueues(entity.getWriteQueueNums() == null ? 0 :
entity.getWriteQueueNums());
+ vo.setPerm(parseTopicPerm(entity.getPerm()));
+ vo.setRemark(entity.getRemark());
+ vo.setCreatedAt(entity.getCreatedAt());
+ vo.setUpdatedAt(entity.getUpdatedAt());
+ return vo;
+ }
+
+ private TopicType parseTopicType(String type) {
+ if (!StringUtils.hasText(type)) {
+ return TopicType.NORMAL;
+ }
+ try {
+ return TopicType.valueOf(type);
+ } catch (IllegalArgumentException ex) {
+ log.debug("Unknown topic type {}, falling back to NORMAL", type);
+ return TopicType.NORMAL;
+ }
+ }
+
+ private TopicPerm parseTopicPerm(Integer perm) {
+ if (perm == null) {
+ return TopicPerm.RW;
+ }
+ return switch (perm) {
+ case 2 -> TopicPerm.WO;
+ case 4 -> TopicPerm.RO;
+ default -> TopicPerm.RW;
+ };
+ }
+
+ @Override
+ public List<ConsumerGroupVO> listConsumerGroups(String clusterId, String
search) {
+ LambdaQueryWrapper<RmqGroup> query = new LambdaQueryWrapper<RmqGroup>()
+ .eq(StringUtils.hasText(clusterId), RmqGroup::getClusterId,
clusterId)
+ .like(StringUtils.hasText(search), RmqGroup::getName, search)
+ .orderByAsc(RmqGroup::getName);
+
+ List<ConsumerGroupVO> result = new ArrayList<>();
+ for (RmqGroup entity : groupMapper.selectList(query)) {
+ ConsumerGroupVO vo = new ConsumerGroupVO();
+ vo.setId(entity.getName());
+ vo.setName(entity.getName());
+ vo.setClusterId(entity.getClusterId());
+ vo.setConsumeType(parseConsumeType(entity.getMessageModel()));
+ vo.setRetryMaxTimes(entity.getMaxRetry() == null ? 0 :
entity.getMaxRetry());
+ vo.setCreatedAt(entity.getCreatedAt());
+ vo.setUpdatedAt(entity.getUpdatedAt());
+
+ if (adminExt != null) {
+ enrichGroupWithConnectionInfo(vo, entity.getName());
+ }
+ result.add(vo);
+ }
+ return result;
+ }
+
+ private ConsumeType parseConsumeType(String messageModel) {
+ if (!StringUtils.hasText(messageModel)) {
+ return ConsumeType.CLUSTERING;
+ }
+ try {
+ return ConsumeType.valueOf(messageModel);
+ } catch (IllegalArgumentException ex) {
+ log.debug("Unknown message model {}, falling back to CLUSTERING",
messageModel);
+ return ConsumeType.CLUSTERING;
+ }
+ }
+
+ @Override
+ public List<BrokerRouteVO> getTopicRoutes(String name) {
+ if (adminExt == null) {
+ return Collections.emptyList();
+ }
+
+ try {
+ TopicRouteData routeData = adminExt.examineTopicRouteInfo(name);
+ if (routeData == null) {
+ return Collections.emptyList();
+ }
+
+ Map<String, BrokerData> brokerDataMap = new HashMap<>();
+ if (routeData.getBrokerDatas() != null) {
+ for (BrokerData bd : routeData.getBrokerDatas()) {
+ brokerDataMap.put(bd.getBrokerName(), bd);
+ }
+ }
+
+ List<BrokerRouteVO> routes = new ArrayList<>();
+ if (routeData.getQueueDatas() != null) {
+ for (QueueData qd : routeData.getQueueDatas()) {
+ BrokerData bd = brokerDataMap.get(qd.getBrokerName());
+ String brokerAddr = "";
+ if (bd != null && bd.getBrokerAddrs() != null &&
!bd.getBrokerAddrs().isEmpty()) {
+ brokerAddr =
bd.getBrokerAddrs().values().iterator().next();
+ }
+
+ routes.add(BrokerRouteVO.builder()
+ .brokerName(qd.getBrokerName())
+ .brokerAddr(brokerAddr)
+ .writeQueues(qd.getWriteQueueNums())
+ .readQueues(qd.getReadQueueNums())
+ .perm(mapPerm(qd.getPerm()))
+ .build());
+ }
+ }
+ return routes;
+ } catch (Exception e) {
+ log.warn("Failed to get routes for topic {}: {}", name,
e.getMessage());
+ return Collections.emptyList();
+ }
+ }
+
+ @Override
+ public List<TopicConsumerVO> getTopicConsumers(String name) {
+ if (adminExt == null) {
+ return Collections.emptyList();
+ }
+
+ try {
+ Set<String> allGroups = collectAllConsumerGroups();
+ List<TopicConsumerVO> consumers = new ArrayList<>();
+
+ for (String group : allGroups) {
+ try {
+ ConsumeStats stats = adminExt.examineConsumeStats(group,
name);
+ if (stats == null || stats.getOffsetTable() == null ||
stats.getOffsetTable().isEmpty()) {
+ continue;
+ }
+
+ long diffTotal = 0;
+ for (Map.Entry<MessageQueue, OffsetWrapper> entry :
stats.getOffsetTable().entrySet()) {
+ OffsetWrapper ow = entry.getValue();
+ diffTotal += Math.max(0, ow.getBrokerOffset() -
ow.getConsumerOffset());
+ }
+
+ ConsumeType consumeType = ConsumeType.CLUSTERING;
+ String messageModel = "CLUSTERING";
+ try {
+ ConsumerConnection conn =
adminExt.examineConsumerConnectionInfo(group);
+ if (conn != null && conn.getConsumeType() != null) {
+ messageModel = conn.getConsumeType().name();
+ if (conn.getMessageModel() != null) {
+ messageModel = conn.getMessageModel().name();
+ }
+ }
+ } catch (Exception ignored) {
+ // group may be offline
+ }
+
+ consumers.add(TopicConsumerVO.builder()
+ .group(group)
+ .consumeType(consumeType)
+ .messageModel(messageModel)
+ .consumeTps(stats.getConsumeTps())
+ .diffTotal(diffTotal)
+ .build());
+ } catch (Exception ignored) {
+ // This group does not subscribe to this topic
+ }
+ }
+ return consumers;
+ } catch (Exception e) {
+ log.warn("Failed to get consumers for topic {}: {}", name,
e.getMessage());
+ return Collections.emptyList();
+ }
+ }
+
+ @Override
+ public List<QueueProgressVO> getGroupProgress(String name) {
+ if (adminExt == null) {
+ return Collections.emptyList();
+ }
+
+ try {
+ ConsumeStats stats = adminExt.examineConsumeStats(name);
+ if (stats == null || stats.getOffsetTable() == null) {
+ return Collections.emptyList();
+ }
+
+ List<QueueProgressVO> progress = new ArrayList<>();
+ for (Map.Entry<MessageQueue, OffsetWrapper> entry :
stats.getOffsetTable().entrySet()) {
+ MessageQueue mq = entry.getKey();
+ OffsetWrapper ow = entry.getValue();
+ long diff = Math.max(0, ow.getBrokerOffset() -
ow.getConsumerOffset());
+
+ progress.add(QueueProgressVO.builder()
+ .broker(mq.getBrokerName())
+ .queueId(mq.getQueueId())
+ .brokerOffset(ow.getBrokerOffset())
+ .consumerOffset(ow.getConsumerOffset())
+ .diffTotal(diff)
+ .build());
+ }
+
+ progress.sort((a, b) -> {
+ int cmp = a.getBroker().compareToIgnoreCase(b.getBroker());
+ return cmp != 0 ? cmp : Integer.compare(a.getQueueId(),
b.getQueueId());
+ });
+ return progress;
+ } catch (Exception e) {
+ log.warn("Failed to get progress for group {}: {}", name,
e.getMessage());
+ return Collections.emptyList();
+ }
+ }
+
+ @Override
+ public List<SubscriptionEntryVO> getGroupSubscriptions(String name) {
+ if (adminExt == null) {
+ return Collections.emptyList();
+ }
+
+ try {
+ ConsumerConnection conn =
adminExt.examineConsumerConnectionInfo(name);
+ if (conn == null || conn.getSubscriptionTable() == null) {
+ return Collections.emptyList();
+ }
+
+ List<SubscriptionEntryVO> subscriptions = new ArrayList<>();
+ for (Map.Entry<String, SubscriptionData> entry :
conn.getSubscriptionTable().entrySet()) {
+ SubscriptionData sd = entry.getValue();
+ subscriptions.add(SubscriptionEntryVO.builder()
+ .topic(sd.getTopic())
+ .expression(sd.getSubString())
+ .type(sd.getExpressionType())
+
.filterMode("CLASS_FILTER".equals(sd.getExpressionType()) ? "SQL" : "TAG")
+ .build());
+ }
+ return subscriptions;
+ } catch (Exception e) {
+ log.warn("Failed to get subscriptions for group {}: {}", name,
e.getMessage());
+ return Collections.emptyList();
+ }
+ }
+
+ // ── Helper methods ──────────────────────────────────────────────────
+
+ private TopicVO buildTopicVO(String topicName) {
+ try {
+ TopicRouteData routeData =
adminExt.examineTopicRouteInfo(topicName);
+ if (routeData == null || routeData.getQueueDatas() == null ||
routeData.getQueueDatas().isEmpty()) {
+ // Topic exists in nameserver but has no route data
+ TopicVO vo = new TopicVO();
+ vo.setId(topicName);
+ vo.setName(topicName);
+ return vo;
+ }
+
+ QueueData firstQd = routeData.getQueueDatas().get(0);
+ TopicVO vo = new TopicVO();
+ vo.setId(topicName);
+ vo.setName(topicName);
+ vo.setWriteQueues(firstQd.getWriteQueueNums());
+ vo.setReadQueues(firstQd.getReadQueueNums());
+ vo.setPerm(mapPerm(firstQd.getPerm()));
+ vo.setType(inferTopicType(topicName));
+ return vo;
+ } catch (Exception e) {
+ log.debug("Failed to get route for topic {}: {}", topicName,
e.getMessage());
+ TopicVO vo = new TopicVO();
+ vo.setId(topicName);
+ vo.setName(topicName);
+ return vo;
+ }
+ }
+
+ private Set<String> collectAllConsumerGroups() {
+ Set<String> allGroups = new HashSet<>();
+ try {
+ ClusterInfo clusterInfo = adminExt.examineBrokerClusterInfo();
+ if (clusterInfo == null || clusterInfo.getBrokerAddrTable() ==
null) {
+ return allGroups;
+ }
+
+ Set<String> processedAddrs = new HashSet<>();
+ for (BrokerData brokerData :
clusterInfo.getBrokerAddrTable().values()) {
+ if (brokerData.getBrokerAddrs() == null) {
+ continue;
+ }
+ // Use master address (brokerId = 0) preferentially
+ String masterAddr = brokerData.getBrokerAddrs().get(0L);
+ if (masterAddr == null &&
!brokerData.getBrokerAddrs().isEmpty()) {
+ masterAddr =
brokerData.getBrokerAddrs().values().iterator().next();
+ }
+ if (masterAddr == null || processedAddrs.contains(masterAddr))
{
+ continue;
+ }
+ processedAddrs.add(masterAddr);
+
+ try {
+ SubscriptionGroupWrapper wrapper =
adminExt.getAllSubscriptionGroup(masterAddr, 5000);
+ if (wrapper != null && wrapper.getSubscriptionGroupTable()
!= null) {
+ java.util.concurrent.ConcurrentMap<String,
SubscriptionGroupConfig> table = wrapper.getSubscriptionGroupTable();
+ allGroups.addAll(table.keySet());
+ }
+ } catch (Exception e) {
+ log.debug("Failed to get subscription groups from broker
{}: {}", masterAddr, e.getMessage());
+ }
+ }
+ } catch (Exception e) {
+ log.warn("Failed to collect consumer groups: {}", e.getMessage());
+ }
+ return allGroups;
+ }
+
+ private void enrichGroupWithConnectionInfo(ConsumerGroupVO vo, String
groupName) {
+ try {
+ ConsumerConnection conn =
adminExt.examineConsumerConnectionInfo(groupName);
+ if (conn != null) {
+ if (conn.getConnectionSet() != null) {
+ vo.setOnlineInstances(conn.getConnectionSet().size());
+ }
+ if (conn.getMessageModel() != null) {
+ vo.setSubscriptionMode(
+
"BROADCASTING".equals(conn.getMessageModel().name())
+ ?
org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode.Pop
+ :
org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode.Push
+ );
+ }
+ if (conn.getSubscriptionTable() != null) {
+ vo.setSubscribedTopics(new
ArrayList<>(conn.getSubscriptionTable().keySet()));
+ }
+ }
+ } catch (Exception ignored) {
+ // Group may be offline, that's fine
+ }
+
+ // Try to get lag info
+ try {
+ ConsumeStats stats = adminExt.examineConsumeStats(groupName);
+ if (stats != null && stats.getOffsetTable() != null) {
+ long totalLag = 0;
+ for (OffsetWrapper ow : stats.getOffsetTable().values()) {
+ totalLag += Math.max(0, ow.getBrokerOffset() -
ow.getConsumerOffset());
+ }
+ vo.setTotalLag(totalLag);
+ }
+ } catch (Exception ignored) {
+ // No stats available
+ }
+ }
+
+ private Set<String> getBrokerNames() {
+ Set<String> names = new HashSet<>();
+ try {
+ ClusterInfo clusterInfo = adminExt.examineBrokerClusterInfo();
+ if (clusterInfo != null && clusterInfo.getBrokerAddrTable() !=
null) {
+ names.addAll(clusterInfo.getBrokerAddrTable().keySet());
+ }
+ } catch (Exception ignored) {
+ }
+ return names;
+ }
+
+ private boolean isSystemTopic(String topicName, Set<String> brokerNames) {
+ if (SYSTEM_TOPICS.contains(topicName)) {
+ return true;
+ }
+ for (String prefix : SYSTEM_TOPIC_PREFIXES) {
+ if (topicName.startsWith(prefix)) {
+ return true;
+ }
+ }
+ // Skip topics that match broker names
+ return brokerNames.contains(topicName);
+ }
+
+ private TopicPerm mapPerm(int perm) {
+ // RocketMQ perm: 6=RW, 4=R, 2=W
+ if (perm == 6) {
+ return TopicPerm.RW;
+ } else if (perm == 4) {
+ return TopicPerm.RO;
+ } else if (perm == 2) {
+ return TopicPerm.WO;
+ }
+ return TopicPerm.RW;
+ }
+
+ private org.apache.rocketmq.studio.common.domain.enums.TopicType
inferTopicType(String topicName) {
+ if (topicName.contains("TRANS") || topicName.contains("trans")) {
+ return
org.apache.rocketmq.studio.common.domain.enums.TopicType.TRANSACTION;
+ }
+ if (topicName.contains("DELAY") || topicName.contains("delay") ||
topicName.contains("SCHEDULE")) {
+ return
org.apache.rocketmq.studio.common.domain.enums.TopicType.DELAY;
+ }
+ if (topicName.contains("FIFO") || topicName.contains("fifo") ||
topicName.contains("ORDER")) {
+ return
org.apache.rocketmq.studio.common.domain.enums.TopicType.FIFO;
+ }
+ return org.apache.rocketmq.studio.common.domain.enums.TopicType.NORMAL;
+ }
+}