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 9f9a0784 feat(tencent): add consumer group management, message query
and trace (#1641)
9f9a0784 is described below
commit 9f9a07843ed1949cb8ecc1cd596154bd0e4353d5
Author: Lei Zhiyuan <[email protected]>
AuthorDate: Tue Aug 11 21:11:33 2026 +0800
feat(tencent): add consumer group management, message query and trace
(#1641)
* feat(tencent): support consumer group management on Tencent Cloud
RocketMQ 5.x
Implement the consumer group surface for the Tencent Cloud provider (Trocket
v20230308 OpenAPI):
- listConsumerGroups: page through DescribeConsumerGroupList with name
filter;
enrich createdAt and the real consume model per group via
DescribeConsumerGroup
(off the cheap count path to avoid N+1 on instance listings)
- createConsumerGroup / deleteConsumerGroup: Create/DeleteConsumerGroup with
retry times and orderly delivery defaults
- getGroupProgress / getGroupSubscriptions: DescribeTopicListByGroup to map
per-topic consumer lag and subscription entries
- resetOffset: ResetConsumerGroupOffset by timestamp (topic required)
- countGroups: reuse the cheap list path
Also set instances/subscribedTopics to empty lists so the consumer group
detail
page no longer throws on null (blank modal), and bump
tencentcloud-sdk-java-trocket from 3.1.1292 to 3.1.1498 which adds
DescribeConsumerGroupResponse#getConsumeModel to map ConsumeModel ->
ConsumeType.
Unit tests cover list (mapping, consume model, timestamps), create, delete,
progress/subscriptions, and reset offset.
* feat(tencent): support message query and trace on Tencent Cloud RocketMQ
5.x
Implement the message surface for the Tencent Cloud provider (Trocket
v20230308 OpenAPI), backed by the message-query APIs:
- queryMessages: query by msgId via DescribeMessage (full detail with
body/properties), or by topic/time/key/tag via DescribeMessageList
(offset/limit paging), mapping MessageItem -> MessageRecordVO
- getMessageTrace: DescribeMessageTrace, parsing the per-stage JSON
payload (produce/persist/consume) into TraceNodeVO nodes and
ConsumerStatusVO entries (consume logs map to delivery status + retry
count)
Unit tests cover msgId detail query, topic/key list query, and trace
stage mapping.
* fix(tencent): set empty TaskRequestId when starting DescribeMessageList
query
DescribeMessageList is an async task-based API: the first call must pass
TaskRequestId as an empty string to start the query, and later pages reuse
the id returned by the previous response. The prior implementation omitted
it, causing 'The request is missing the required parameter TaskRequestId'
from the OpenAPI when querying messages by topic/time/key/tag.
* fix(tencent): make ProduceTime parsing more robust and log parse failures
DescribeMessageList returns ProduceTime as 'yyyy-MM-dd HH:mm:ss' (no
milliseconds). The existing single formatter already handled that, but the
page showed storeTime as epoch 0 (1970), so ProduceTime may be null,
empty, or in an unexpected format at runtime. Broaden the accepted
formatters and add warn logging to capture the actual raw ProduceTime when
parsing fails, so the mismatch can be diagnosed against the live API.
* fix(tencent): surface tag and key for msgId query from DescribeMessage
properties
DescribeMessage returns the tag and key only inside Properties as TAGS and
KEYS, not as top-level fields. The msgId query path therefore left
MessageRecordVO.tag/key null, so the page showed neither. Extract TAGS/KEYS
from the parsed properties and map them onto the record.
* fix(tencent): pass topic through to message trace API
DescribeMessageTrace requires a non-empty Topic, but the trace endpoint and
frontend only carried msgId + instanceId, so Tencent rejected the empty
topic ('parameter value [Topic] = [] is invalid'). Thread the message topic
through the whole trace path:
- MessageController.getMessageTrace accepts an optional topic param
- MessageService / InstanceProvider / MessageProvider.getMessageTrace gain a
topic argument (Apache/Aliyun providers ignore it as before)
- TencentInstanceProvider uses the topic in DescribeMessageTraceRequest
- Frontend passes record.topic when opening the trace tab
Unit tests updated for the new signatures.
* fix(tencent): page message list with a stable TaskRequestId per logical
query
DescribeMessageList is an async task-based API. Each logical query is
identified by a TaskRequestId: the first call starts it, and paging through
the same query reuses that id (returned in the response), while a new query
uses a fresh id. The prior implementation passed an empty string and only
fetched the first page, so multi-page results were truncated and the task
id semantics were wrong.
Generate a random TaskRequestId per query and page through the whole result
set (advancing Offset until TotalCount is reached or MAX_PAGES), reusing the
same task id across pages. The frontend message table is not
server-paginated
(pagination=false), so the server aggregates all pages.
Add a unit test asserting offset advances while the task id stays stable
across pages.
* fix(tencent): terminate message list paging on short page, not just
TotalCount
Align the Tencent message-list paging termination with the Aliyun provider:
stop when the last page returns fewer rows than MESSAGE_LIMIT, using
TotalCount only as a secondary safety net. Some Tencent query tasks may not
populate TotalCount reliably, and relying on it could truncate multi-page
results. Add a unit test that drives a full page then a short page and
asserts offset advances while TaskRequestId stays stable.
* refactor(tencent): address review readability and document topic
requirement
- Break the message-list 'last page' predicate into named booleans
(shortPage / allCollected) so the || / && precedence is explicit without
relying on operator order.
- Document in MessageController that the optional topic param is required
by the Tencent provider (DescribeMessageTrace) but not by Apache/Aliyun.
---
server/pom.xml | 2 +-
.../studio/instance/message/MessageController.java | 7 +-
.../studio/instance/message/MessageProvider.java | 2 +-
.../instance/message/MessageProviderStub.java | 2 +-
.../studio/instance/message/MessageService.java | 8 +-
.../rocketmq/studio/provider/InstanceProvider.java | 2 +-
.../provider/alibaba/AliyunInstanceProvider.java | 2 +-
.../provider/apache/ApacheInstanceProvider.java | 4 +-
.../provider/apache/RocketMQMessageProvider.java | 2 +-
.../provider/tencent/TencentInstanceProvider.java | 579 ++++++++++++++++++++-
.../instance/message/MessageControllerTest.java | 7 +-
.../instance/message/MessageProviderStubTest.java | 2 +-
.../instance/message/MessageServiceTest.java | 2 +-
.../alibaba/AliyunInstanceProviderTest.java | 2 +-
.../apache/RocketMQMessageProviderTest.java | 6 +-
.../apache/RocketMQMetadataProviderTest.java | 2 -
.../tencent/TencentInstanceProviderTest.java | 261 ++++++++++
web/src/api/message.ts | 4 +-
web/src/pages/instance/message.tsx | 2 +-
web/src/services/messageService.ts | 3 +-
20 files changed, 858 insertions(+), 43 deletions(-)
diff --git a/server/pom.xml b/server/pom.xml
index c142baf1..c9e1cc12 100644
--- a/server/pom.xml
+++ b/server/pom.xml
@@ -108,7 +108,7 @@
<dependency>
<groupId>com.tencentcloudapi</groupId>
<artifactId>tencentcloud-sdk-java-trocket</artifactId>
- <version>3.1.1292</version>
+ <version>3.1.1498</version>
</dependency>
<!-- rocketmq-common also brings Kotlin-based okio-jvm 3.x; keep its
runtime available. -->
<dependency>
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
index 3b2ddbfb..3fb6da3c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java
@@ -46,7 +46,10 @@ public class MessageController {
}
@GetMapping("/{msgId}/trace")
- public Result<TraceRecordVO> getMessageTrace(@PathVariable String msgId,
@RequestParam String instanceId) {
- return Result.ok(messageService.getMessageTrace(instanceId, msgId));
+ public Result<TraceRecordVO> getMessageTrace(@PathVariable String msgId,
@RequestParam String instanceId,
+ // Optional for the
Apache/Aliyun providers; the Tencent
+ // provider requires a
non-empty topic (DescribeMessageTrace).
+ @RequestParam(required =
false) String topic) {
+ return Result.ok(messageService.getMessageTrace(instanceId, msgId,
topic));
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
index fe79a31e..c4cc365a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageProvider.java
@@ -23,5 +23,5 @@ public interface MessageProvider {
List<MessageRecordVO> queryMessages(String instanceId, String topic,
String msgId, String tag, String key, Long startTime,
Long endTime);
- TraceRecordVO getMessageTrace(String instanceId, String msgId);
+ TraceRecordVO getMessageTrace(String instanceId, String msgId, String
topic);
}
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 932d7117..7b24291b 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
@@ -36,7 +36,7 @@ public class MessageProviderStub implements MessageProvider {
}
@Override
- public TraceRecordVO getMessageTrace(String instanceId, String msgId) {
+ public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
log.warn("MessageProviderStub.getMessageTrace called but no real
message provider is configured");
throw unsupported();
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
index 5e5a7da5..e4602bbe 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageService.java
@@ -44,14 +44,14 @@ public class MessageService {
.orElseGet(() -> messageProvider.queryMessages(instanceId,
topic, msgId, tag, key, startTime, endTime));
}
- public TraceRecordVO getMessageTrace(String instanceId, String msgId) {
+ public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
if (!StringUtils.hasText(msgId)) {
throw new BusinessException(400, "msgId is required");
}
- log.info("Getting message trace: msgId={}", msgId);
+ log.info("Getting message trace: msgId={}, topic={}", msgId, topic);
return providerRegistry.byInstanceId(instanceId)
- .map(provider -> provider.getMessageTrace(instanceId, msgId))
- .orElseGet(() -> messageProvider.getMessageTrace(instanceId,
msgId));
+ .map(provider -> provider.getMessageTrace(instanceId, msgId,
topic))
+ .orElseGet(() -> messageProvider.getMessageTrace(instanceId,
msgId, topic));
}
private void validateTopicQueryWindow(String topic, String msgId, String
key, Long startTime, Long endTime) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
index 4582a4a6..9cdecb69 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/InstanceProvider.java
@@ -79,5 +79,5 @@ public interface InstanceProvider {
List<MessageRecordVO> queryMessages(String instanceId, String topic,
String msgId,
String tag, String key, Long
startTime, Long endTime);
- TraceRecordVO getMessageTrace(String instanceId, String msgId);
+ TraceRecordVO getMessageTrace(String instanceId, String msgId, String
topic);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
index 447249c0..4e449078 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProvider.java
@@ -430,7 +430,7 @@ public class AliyunInstanceProvider implements
InstanceProvider {
}
@Override
- public TraceRecordVO getMessageTrace(String instanceId, String msgId) {
+ public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
Context ctx = resolve(instanceId);
GetTraceRequest request = GetTraceRequest.builder()
.instanceId(ctx.cloudInstanceId())
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
index 3bf81839..9e82244d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProvider.java
@@ -131,8 +131,8 @@ public class ApacheInstanceProvider implements
InstanceProvider {
}
@Override
- public TraceRecordVO getMessageTrace(String instanceId, String msgId) {
- return messageProvider.getMessageTrace(instanceId, msgId);
+ public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
+ return messageProvider.getMessageTrace(instanceId, msgId, topic);
}
private boolean matchesInstance(String topicInstanceId, String instanceId)
{
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index 01433857..511d66cd 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -264,7 +264,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
}
@Override
- public TraceRecordVO getMessageTrace(String instanceId, String msgId) {
+ public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
return runtimeAdminClientResolver.execute(instanceId,
adminExt -> getMessageTrace(instanceId, (DefaultMQAdminExt)
adminExt, msgId));
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
index 5ab7482d..aba0521f 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java
@@ -16,16 +16,35 @@
*/
package org.apache.rocketmq.studio.provider.tencent;
+import com.tencentcloudapi.trocket.v20230308.models.ConsumeGroupItem;
+import com.tencentcloudapi.trocket.v20230308.models.CreateConsumerGroupRequest;
import com.tencentcloudapi.trocket.v20230308.models.CreateTopicRequest;
+import com.tencentcloudapi.trocket.v20230308.models.DeleteConsumerGroupRequest;
import com.tencentcloudapi.trocket.v20230308.models.DeleteTopicRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeConsumerGroupRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeConsumerGroupResponse;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeConsumerGroupListRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeConsumerGroupListResponse;
+import com.tencentcloudapi.trocket.v20230308.models.DescribeMessageListRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeMessageListResponse;
+import com.tencentcloudapi.trocket.v20230308.models.DescribeMessageRequest;
+import com.tencentcloudapi.trocket.v20230308.models.DescribeMessageResponse;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeMessageTraceRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeMessageTraceResponse;
+import com.tencentcloudapi.trocket.v20230308.models.MessageItem;
+import com.tencentcloudapi.trocket.v20230308.models.MessageTraceItem;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListByGroupRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListByGroupResponse;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListRequest;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListResponse;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicRequest;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicResponse;
import com.tencentcloudapi.trocket.v20230308.models.ModifyTopicRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.ResetConsumerGroupOffsetRequest;
import com.tencentcloudapi.trocket.v20230308.models.SubscriptionData;
import com.tencentcloudapi.trocket.v20230308.models.TopicItem;
import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
+import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
@@ -35,7 +54,9 @@ import org.apache.rocketmq.studio.instance.InstanceVO;
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.message.ConsumerStatusVO;
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.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
@@ -43,22 +64,32 @@ import org.apache.rocketmq.studio.provider.InstanceProvider;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneId;
+import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
+import java.util.Map;
+import java.util.UUID;
/**
* Tencent Cloud TDMQ RocketMQ 5.x topic operations backed by Trocket
v20230308 OpenAPI.
*
- * <p>Topic management is supported independently from the other
instance-scoped operations. The
- * remaining operations intentionally retain the provider's
unsupported-operation behavior until
- * their corresponding Tencent Cloud APIs are mapped to Studio's common
models.</p>
+ * <p>In addition to topic/consumer-group management, message query is
supported via
+ * DescribeMessageList (list), DescribeMessage (detail) and
DescribeMessageTrace (trace).</p>
*/
@RequiredArgsConstructor
+@Slf4j
@Component
public class TencentInstanceProvider implements InstanceProvider {
@@ -68,7 +99,21 @@ public class TencentInstanceProvider implements
InstanceProvider {
static final int DEFAULT_QUEUE_NUM = 8;
static final int MIN_QUEUE_NUM = 3;
static final int MAX_QUEUE_NUM = 16;
- private static final String NOT_IMPLEMENTED = "Tencent Cloud operation is
not implemented yet";
+ static final int DEFAULT_MAX_RETRY_TIMES = 16;
+ static final int MESSAGE_LIMIT = 100;
+ private static final DateTimeFormatter TENCENT_TIME_FORMATTER =
DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss[,SSS][,SS]");
+ private static final DateTimeFormatter[] TENCENT_TIME_FORMATTERS = {
+ TENCENT_TIME_FORMATTER,
+ DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"),
+ DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS"),
+ DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH:mm:ss.SSSZ")
+ };
+ private static final long ONE_HOUR_MILLIS = 60L * 60L * 1000L;
+ private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
+ private static final String STAGE_PRODUCE = "produce";
+ private static final String STAGE_PERSIST = "persist";
+ private static final String STAGE_CONSUME = "consume";
private final TencentClientFactory clientFactory;
private final InstanceRepository instanceRepository;
@@ -97,7 +142,7 @@ public class TencentInstanceProvider implements
InstanceProvider {
@Override
public int countGroups(String instanceId) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ return listConsumerGroups(instanceId, null, false).size();
}
@Override
@@ -257,43 +302,450 @@ public class TencentInstanceProvider implements
InstanceProvider {
@Override
public List<ConsumerGroupVO> listConsumerGroups(String instanceId, String
search) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ return listConsumerGroups(instanceId, search, true);
+ }
+
+ private List<ConsumerGroupVO> listConsumerGroups(String instanceId, String
search, boolean enrichTimes) {
+ Context context = resolve(instanceId);
+ List<ConsumerGroupVO> groups = new ArrayList<>();
+ for (int page = 0; page < MAX_PAGES; page++) {
+ DescribeConsumerGroupListRequest request = new
DescribeConsumerGroupListRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setOffset((long) page * PAGE_SIZE);
+ request.setLimit((long) PAGE_SIZE);
+ DescribeConsumerGroupListResponse response =
clientFactory.call(context.credentialId(), context.regionId(),
+ client -> client.DescribeConsumerGroupList(request));
+ ConsumeGroupItem[] data = response == null ? null :
response.getData();
+ if (data == null || data.length == 0) {
+ break;
+ }
+ for (ConsumeGroupItem item : data) {
+ if (item == null) {
+ continue;
+ }
+ ConsumerGroupVO group = toConsumerGroup(item, instanceId);
+ if (matchesSearch(search, group.getName())) {
+ if (enrichTimes) {
+ enrichConsumerGroupDetail(context, group);
+ }
+ groups.add(group);
+ }
+ }
+ if (data.length < PAGE_SIZE) {
+ break;
+ }
+ }
+ return groups;
+ }
+
+ /**
+ * DescribeConsumerGroupList exposes limited fields, so resolve the
creation timestamp and the
+ * real consume model from DescribeConsumerGroup. Kept off the cheap count
path to avoid N+1
+ * calls for instance listings.
+ */
+ private void enrichConsumerGroupDetail(Context context, ConsumerGroupVO
group) {
+ try {
+ DescribeConsumerGroupRequest request = new
DescribeConsumerGroupRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setConsumerGroup(group.getName());
+ DescribeConsumerGroupResponse response =
clientFactory.call(context.credentialId(), context.regionId(),
+ client -> client.DescribeConsumerGroup(request));
+ if (response == null) {
+ return;
+ }
+ group.setCreatedAt(toLocalDateTime(response.getCreatedTime()));
+ if (StringUtils.hasText(response.getConsumeModel())) {
+
group.setConsumeType(toConsumeType(response.getConsumeModel()));
+ }
+ } catch (BusinessException ignored) {
+ // A single group detail lookup failure should not fail the whole
list.
+ }
}
@Override
public ConsumerGroupVO createConsumerGroup(String instanceId,
ConsumerGroupVO group) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ Context context = resolve(instanceId);
+ validateConsumerGroup(group);
+ CreateConsumerGroupRequest request = new CreateConsumerGroupRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setConsumerGroup(group.getName());
+ request.setMaxRetryTimes((long) retryMaxTimes(group));
+ request.setConsumeEnable(true);
+ request.setConsumeMessageOrderly(isOrderly(group));
+ clientFactory.call(context.credentialId(), context.regionId(), client
-> client.CreateConsumerGroup(request));
+ group.setInstanceId(instanceId);
+ group.setRetryMaxTimes(retryMaxTimes(group));
+ group.setSubscribedTopics(java.util.List.of());
+ group.setCreatedAt(LocalDateTime.now());
+ group.setUpdatedAt(LocalDateTime.now());
+ return group;
}
@Override
public void deleteConsumerGroup(String instanceId, String groupName) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ Context context = resolve(instanceId);
+ requireGroupName(groupName);
+ DeleteConsumerGroupRequest request = new DeleteConsumerGroupRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setConsumerGroup(groupName);
+ clientFactory.call(context.credentialId(), context.regionId(), client
-> client.DeleteConsumerGroup(request));
}
@Override
public List<QueueProgressVO> getGroupProgress(String instanceId, String
groupName) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ Context context = resolve(instanceId);
+ requireGroupName(groupName);
+ List<SubscriptionData> subscriptions =
listTopicSubscriptionsByGroup(context, groupName);
+ List<QueueProgressVO> rows = new ArrayList<>();
+ for (SubscriptionData subscription : subscriptions) {
+ if (subscription == null) {
+ continue;
+ }
+ rows.add(QueueProgressVO.builder()
+ .broker("topic:" + subscription.getTopic())
+ .queueId(0)
+ .brokerOffset(0L)
+ .consumerOffset(0L)
+ .diffTotal(subscription.getConsumerLag() == null ? 0L :
subscription.getConsumerLag())
+ .build());
+ }
+ return rows;
}
@Override
public List<SubscriptionEntryVO> getGroupSubscriptions(String instanceId,
String groupName) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ Context context = resolve(instanceId);
+ requireGroupName(groupName);
+ List<SubscriptionData> subscriptions =
listTopicSubscriptionsByGroup(context, groupName);
+ List<SubscriptionEntryVO> entries = new ArrayList<>();
+ for (SubscriptionData subscription : subscriptions) {
+ if (subscription == null) {
+ continue;
+ }
+ entries.add(toSubscriptionEntry(subscription));
+ }
+ return entries;
}
@Override
public void resetOffset(String instanceId, String groupName, long
timestamp, String topic) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ Context context = resolve(instanceId);
+ requireGroupName(groupName);
+ if (!StringUtils.hasText(topic)) {
+ throw new BusinessException(400, "Topic is required to reset
consumer group offset");
+ }
+ ResetConsumerGroupOffsetRequest request = new
ResetConsumerGroupOffsetRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setConsumerGroup(groupName);
+ request.setTopic(topic);
+ if (timestamp > 0L) {
+ request.setResetTimestamp(timestamp);
+ } else {
+ request.setResetTimestamp(System.currentTimeMillis());
+ }
+ clientFactory.call(context.credentialId(), context.regionId(), client
-> client.ResetConsumerGroupOffset(request));
}
@Override
public List<MessageRecordVO> queryMessages(String instanceId, String
topic, String msgId,
String tag, String key, Long
startTime, Long endTime) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ Context context = resolve(instanceId);
+ requireTopic(topic);
+ // Querying by message ID returns the full detail (body, properties
and tracks) via
+ // DescribeMessage, mirroring the msgId path of the base provider.
+ if (StringUtils.hasText(msgId)) {
+ MessageRecordVO record = toRecordVO(describeMessage(context,
topic, msgId));
+ return record == null ? Collections.emptyList() :
Collections.singletonList(record);
+ }
+
+ long end = endTime != null ? endTime : System.currentTimeMillis();
+ long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
+ if (begin >= end) {
+ throw new BusinessException(400, "Message query start time must be
before end time");
+ }
+
+ // DescribeMessageList is an async, task-based query: each logical
query is identified by a
+ // TaskRequestId, and paging through that query reuses the same id
(the response returns it
+ // for the next page). A fresh random id starts a brand-new query. The
frontend message table
+ // is not server-paginated (pagination=false), so page through the
whole result set here.
+ String taskRequestId = UUID.randomUUID().toString();
+ List<MessageRecordVO> result = new ArrayList<>();
+ for (int page = 0; page < MAX_PAGES; page++) {
+ DescribeMessageListRequest request = new
DescribeMessageListRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setTopic(topic);
+ request.setStartTime(begin);
+ request.setEndTime(end);
+ request.setTaskRequestId(taskRequestId);
+ if (StringUtils.hasText(key)) {
+ request.setMsgKey(key);
+ }
+ if (StringUtils.hasText(tag)) {
+ request.setTag(tag);
+ }
+ request.setOffset((long) result.size());
+ request.setLimit((long) MESSAGE_LIMIT);
+ DescribeMessageListResponse response =
clientFactory.call(context.credentialId(), context.regionId(),
+ client -> client.DescribeMessageList(request));
+ MessageItem[] data = response == null ? null : response.getData();
+ long total = response == null ? 0L : (response.getTotalCount() ==
null ? 0L : response.getTotalCount());
+ if (data != null) {
+ for (MessageItem item : data) {
+ if (item != null) {
+ result.add(toRecordVO(item, topic));
+ }
+ }
+ }
+ // Stop on the last page (returned fewer rows than requested) or
once all results have
+ // been collected. Like the Aliyun provider, the short-page check
is the primary signal
+ // so we do not rely on TotalCount, which may not be populated for
every query.
+ int returned = data == null ? 0 : data.length;
+ // Stop on the last page (returned fewer rows than requested) or
once all results have
+ // been collected. Like the Aliyun provider, the short-page check
is the primary signal
+ // so we do not rely on TotalCount, which may not be populated for
every query.
+ if (isLastPage(returned, total, result.size())) {
+ break;
+ }
+ }
+ return result;
}
@Override
- public TraceRecordVO getMessageTrace(String instanceId, String msgId) {
- throw new UnsupportedOperationException(NOT_IMPLEMENTED);
+ public TraceRecordVO getMessageTrace(String instanceId, String msgId,
String topic) {
+ Context context = resolve(instanceId);
+ requireMsgId(msgId);
+ requireTopic(topic);
+
+ DescribeMessageTraceRequest request = new
DescribeMessageTraceRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setTopic(topic);
+ request.setMsgId(msgId);
+ DescribeMessageTraceResponse response =
clientFactory.call(context.credentialId(), context.regionId(),
+ client -> client.DescribeMessageTrace(request));
+
+ List<TraceNodeVO> nodes = new ArrayList<>();
+ List<ConsumerStatusVO> consumerStatus = new ArrayList<>();
+ MessageTraceItem[] items = response == null ? null :
response.getData();
+ if (items != null) {
+ for (MessageTraceItem item : items) {
+ buildTraceStage(item, nodes, consumerStatus);
+ }
+ }
+ return TraceRecordVO.builder()
+ .nodes(nodes)
+ .consumerStatus(consumerStatus)
+ .build();
+ }
+
+ private DescribeMessageResponse describeMessage(Context context, String
topic, String msgId) {
+ DescribeMessageRequest request = new DescribeMessageRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setTopic(topic);
+ request.setMsgId(msgId);
+ return clientFactory.call(context.credentialId(), context.regionId(),
+ client -> client.DescribeMessage(request));
+ }
+
+ private static void buildTraceStage(MessageTraceItem item,
List<TraceNodeVO> nodes,
+ List<ConsumerStatusVO> consumerStatus)
{
+ String stage = item == null ? null : item.getStage();
+ String data = item == null ? null : item.getData();
+ if (!StringUtils.hasText(stage) || !StringUtils.hasText(data)) {
+ return;
+ }
+ try {
+ JsonNode root = OBJECT_MAPPER.readTree(data);
+ switch (stage) {
+ case STAGE_PRODUCE:
+ nodes.add(buildProduceNode(root));
+ break;
+ case STAGE_PERSIST:
+ nodes.add(buildPersistNode(root));
+ break;
+ case STAGE_CONSUME:
+ buildConsumeNodes(root, nodes, consumerStatus);
+ break;
+ default:
+ break;
+ }
+ } catch (Exception e) {
+ // Unparseable trace payloads are skipped rather than failing the
whole trace.
+ }
+ }
+
+ private static TraceNodeVO buildProduceNode(JsonNode root) {
+ return TraceNodeVO.builder()
+ .title(STAGE_PRODUCE)
+
.timestamp(parseTraceTime(root.path("ProduceTime").asText(null)))
+ .status(toTraceStatus(root.path("Status").asInt(0)))
+ .costTime(root.path("Duration").asLong(0L))
+ .description("producer=" +
root.path("ProducerAddr").asText(""))
+ .build();
+ }
+
+ private static TraceNodeVO buildPersistNode(JsonNode root) {
+ return TraceNodeVO.builder()
+ .title(STAGE_PERSIST)
+
.timestamp(parseTraceTime(root.path("PersistTime").asText(null)))
+ .status(toTraceStatus(root.path("Status").asInt(0)))
+ .costTime(0L)
+ .description("store persist")
+ .build();
+ }
+
+ private static void buildConsumeNodes(JsonNode root, List<TraceNodeVO>
nodes,
+ List<ConsumerStatusVO>
consumerStatus) {
+ JsonNode logs = root.path("RocketMqConsumeLogs");
+ if (!logs.isArray()) {
+ return;
+ }
+ for (JsonNode log : logs) {
+ String group = log.path("ConsumerGroup").asText("");
+ long pushTime = parseTraceTime(log.path("PushTime").asText(null));
+ int status = log.path("Status").asInt(0);
+ int retryTimes = log.path("RetryTimes").asInt(0);
+ nodes.add(TraceNodeVO.builder()
+ .title(STAGE_CONSUME)
+ .timestamp(pushTime)
+ .status(toConsumeTraceStatus(status))
+ .costTime(0L)
+ .description("group=" + group + ", consumer=" +
log.path("ConsumerAddr").asText(""))
+ .build());
+ consumerStatus.add(ConsumerStatusVO.builder()
+ .group(group)
+ .deliveryStatus(toDeliveryStatus(status))
+ .consumeTime(pushTime)
+ .retryCount(retryTimes)
+ .build());
+ }
+ }
+
+ private static MessageRecordVO toRecordVO(DescribeMessageResponse
response) {
+ if (response == null) {
+ return null;
+ }
+ Map<String, String> properties =
parseProperties(response.getProperties());
+ // DescribeMessage carries the tag and key inside Properties (TAGS /
KEYS), not as
+ // top-level fields, so surface them onto the record for the message
list/detail page.
+ return MessageRecordVO.builder()
+ .msgId(defaultIfBlank(response.getMessageId(), ""))
+ .topic(defaultIfBlank(response.getShowTopicName(), ""))
+ .tag(properties.get("TAGS"))
+ .key(properties.get("KEYS"))
+ .body(response.getBody())
+ .bodyEncoding("UTF-8")
+ .bodyTruncated(false)
+ .storeTime(parseTraceTime(response.getProduceTime()))
+ .bornHost(response.getProducerAddr())
+ .properties(properties)
+ .propertiesTruncated(false)
+ .size(0)
+ .build();
+ }
+
+ private static MessageRecordVO toRecordVO(MessageItem item, String topic) {
+ if (item == null) {
+ return null;
+ }
+ return MessageRecordVO.builder()
+ .msgId(defaultIfBlank(item.getMsgId(), ""))
+ .topic(topic)
+ .tag(item.getTags())
+ .key(item.getKeys())
+ .body(null)
+ .bodyEncoding(null)
+ .bodyTruncated(false)
+ .storeTime(parseTraceTime(item.getProduceTime()))
+ .bornHost(item.getProducerAddr())
+ .properties(Collections.emptyMap())
+ .propertiesTruncated(false)
+ .size(0)
+ .build();
+ }
+
+ private static Map<String, String> parseProperties(String raw) {
+ if (!StringUtils.hasText(raw)) {
+ return Collections.emptyMap();
+ }
+ try {
+ JsonNode root = OBJECT_MAPPER.readTree(raw);
+ if (root == null || !root.isObject()) {
+ return Collections.emptyMap();
+ }
+ Map<String, String> properties = new LinkedHashMap<>();
+ root.fields().forEachRemaining(entry ->
properties.put(entry.getKey(), entry.getValue().asText("")));
+ return properties;
+ } catch (Exception e) {
+ return Collections.emptyMap();
+ }
+ }
+
+ private static boolean isLastPage(int returned, long total, int collected)
{
+ // The last page is one that returned no rows or fewer rows than
requested. When the API
+ // reports a TotalCount, also stop once all expected results have been
collected.
+ boolean shortPage = returned == 0 || returned < MESSAGE_LIMIT;
+ boolean allCollected = total > 0 && collected >= total;
+ return shortPage || allCollected;
+ }
+
+ private static long parseTraceTime(String value) {
+ if (!StringUtils.hasText(value)) {
+ log.warn("Tencent message query: empty ProduceTime, storeTime=0");
+ return 0L;
+ }
+ String trimmed = value.trim();
+ for (DateTimeFormatter formatter : TENCENT_TIME_FORMATTERS) {
+ try {
+ return LocalDateTime.parse(trimmed, formatter)
+
.atZone(ZoneId.systemDefault()).toInstant().toEpochMilli();
+ } catch (Exception ignored) {
+ // try next format
+ }
+ }
+ try {
+ return Long.parseLong(trimmed);
+ } catch (NumberFormatException ignored) {
+ log.warn("Tencent message query: unparseable ProduceTime={},
storeTime=0", value);
+ }
+ return 0L;
+ }
+
+ private static String toTraceStatus(int status) {
+ return status == 0 ? "finish" : "failed";
+ }
+
+ private static String toConsumeTraceStatus(int status) {
+ // Tencent consume log Status uses the RocketMQ convention where 2
means consumed.
+ return status == 2 ? "finish" : "failed";
+ }
+
+ private static DeliveryStatus toDeliveryStatus(int status) {
+ // Tencent consume log status: 0/1 in-flight, 2 consumed, others
failed.
+ switch (status) {
+ case 2:
+ return DeliveryStatus.success;
+ case 0:
+ case 1:
+ return DeliveryStatus.pending;
+ default:
+ return DeliveryStatus.failed;
+ }
+ }
+
+ private static String defaultIfBlank(String value, String fallback) {
+ return value == null || value.isBlank() ? fallback : value;
+ }
+
+ private static void requireTopic(String topic) {
+ if (!StringUtils.hasText(topic)) {
+ throw new BusinessException(400, "topic is required for Tencent
Cloud message query");
+ }
+ }
+
+ private static void requireMsgId(String msgId) {
+ if (!StringUtils.hasText(msgId)) {
+ throw new BusinessException(400, "msgId is required for message
trace query");
+ }
}
private Context resolve(String instanceId) {
@@ -336,6 +788,98 @@ public class TencentInstanceProvider implements
InstanceProvider {
.build();
}
+ private static ConsumerGroupVO toConsumerGroup(ConsumeGroupItem item,
String instanceId) {
+ ConsumerGroupVO group = new ConsumerGroupVO();
+ group.setName(item.getConsumerGroup());
+ group.setInstanceId(instanceId);
+ group.setClusterId(item.getClusterIdV4());
+ group.setNamespace(item.getNamespaceV4());
+ group.setConsumeType(toConsumeType(item.getConsumeMessageOrderly()));
+ group.setDeliveryOrderType(item.getConsumeMessageOrderly() == null ||
!item.getConsumeMessageOrderly()
+ ? "Concurrently" : "Orderly");
+ group.setRetryMaxTimes(toInt(item.getMaxRetryTimes()));
+ group.setSubscribedTopics(java.util.List.of());
+ group.setInstances(java.util.List.of());
+ return group;
+ }
+
+ private List<SubscriptionData> listTopicSubscriptionsByGroup(Context
context, String groupName) {
+ List<SubscriptionData> all = new ArrayList<>();
+ for (int page = 0; page < MAX_PAGES; page++) {
+ DescribeTopicListByGroupRequest request = new
DescribeTopicListByGroupRequest();
+ request.setInstanceId(context.cloudInstanceId());
+ request.setConsumerGroup(groupName);
+ request.setOffset((long) page * PAGE_SIZE);
+ request.setLimit((long) PAGE_SIZE);
+ DescribeTopicListByGroupResponse response =
clientFactory.call(context.credentialId(), context.regionId(),
+ client -> client.DescribeTopicListByGroup(request));
+ SubscriptionData[] data = response == null ? null :
response.getData();
+ if (data == null || data.length == 0) {
+ break;
+ }
+ all.addAll(Arrays.asList(data));
+ if (data.length < PAGE_SIZE) {
+ break;
+ }
+ }
+ return all;
+ }
+
+ private static SubscriptionEntryVO toSubscriptionEntry(SubscriptionData
subscription) {
+ return SubscriptionEntryVO.builder()
+ .topic(subscription.getTopic())
+ .expression(subscription.getSubString())
+ .type(subscription.getExpressionType())
+ .filterMode(subscription.getExpressionType())
+ .consistency(subscription.getConsistency() == null ? null :
String.valueOf(subscription.getConsistency()))
+ .build();
+ }
+
+ private static void validateConsumerGroup(ConsumerGroupVO group) {
+ if (group == null || !StringUtils.hasText(group.getName())) {
+ throw new BusinessException(400, "Consumer group name is
required");
+ }
+ }
+
+ private static void requireGroupName(String groupName) {
+ if (!StringUtils.hasText(groupName)) {
+ throw new BusinessException(400, "Consumer group name is
required");
+ }
+ }
+
+ private static int retryMaxTimes(ConsumerGroupVO group) {
+ return group.getRetryMaxTimes() > 0 ? group.getRetryMaxTimes() :
DEFAULT_MAX_RETRY_TIMES;
+ }
+
+ private static boolean isOrderly(ConsumerGroupVO group) {
+ String deliveryOrderType = group.getDeliveryOrderType();
+ return deliveryOrderType != null
+ && (deliveryOrderType.toUpperCase(Locale.ROOT).contains("FIFO")
+ ||
deliveryOrderType.toUpperCase(Locale.ROOT).contains("ORDER"));
+ }
+
+ private static int toInt(Long value) {
+ return value == null ? 0 : Math.toIntExact(value);
+ }
+
+ private static ConsumeType toConsumeType(Boolean consumeMessageOrderly) {
+ // Tencent consumer groups use clustering consumption; orderly only
affects delivery order.
+ return ConsumeType.CLUSTERING;
+ }
+
+ private static ConsumeType toConsumeType(String raw) {
+ if (!StringUtils.hasText(raw)) {
+ return null;
+ }
+ if (raw.toUpperCase(Locale.ROOT).contains("BROADCAST")) {
+ return ConsumeType.BROADCASTING;
+ }
+ if (raw.toUpperCase(Locale.ROOT).contains("CLUSTER")) {
+ return ConsumeType.CLUSTERING;
+ }
+ return null;
+ }
+
private static TopicType toTopicType(String raw) {
if (!StringUtils.hasText(raw)) {
return null;
@@ -374,6 +918,13 @@ public class TencentInstanceProvider implements
InstanceProvider {
return contains(topic.getName(), needle) ||
contains(topic.getRemark(), needle);
}
+ private static boolean matchesSearch(String search, String value) {
+ if (!StringUtils.hasText(search)) {
+ return true;
+ }
+ return contains(value, search.trim().toLowerCase(Locale.ROOT));
+ }
+
private static boolean contains(String value, String needle) {
return value != null &&
value.toLowerCase(Locale.ROOT).contains(needle);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
index 36a87ee5..e1207158 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageControllerTest.java
@@ -74,12 +74,13 @@ class MessageControllerTest {
@Test
void messageTraceShouldPassInstanceId() throws Exception {
TraceRecordVO trace =
TraceRecordVO.builder().nodes(List.of()).consumerStatus(List.of()).build();
- when(messageService.getMessageTrace("instance-a",
"msg-001")).thenReturn(trace);
+ when(messageService.getMessageTrace("instance-a", "msg-001",
"orders")).thenReturn(trace);
- mockMvc.perform(get("/api/messages/msg-001/trace").param("instanceId",
"instance-a"))
+ mockMvc.perform(get("/api/messages/msg-001/trace").param("instanceId",
"instance-a")
+ .param("topic", "orders"))
.andExpect(status().isOk())
.andExpect(jsonPath("$.code").value(200));
- verify(messageService).getMessageTrace("instance-a", "msg-001");
+ verify(messageService).getMessageTrace("instance-a", "msg-001",
"orders");
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
index be2fa14f..d919a820 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageProviderStubTest.java
@@ -37,7 +37,7 @@ class MessageProviderStubTest {
@Test
void getMessageTraceShouldFailExplicitlyWhenRealProviderIsMissing() {
- assertThatThrownBy(() -> provider.getMessageTrace("instance-a",
"msg-001"))
+ assertThatThrownBy(() -> provider.getMessageTrace("instance-a",
"msg-001", "orders"))
.isInstanceOf(BusinessException.class)
.hasMessage("Message query provider is not configured")
.extracting("code")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
index 8d96cc8b..754f2647 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
@@ -52,7 +52,7 @@ class MessageServiceTest {
InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
MessageService service = new MessageService(provider, registry);
- assertThatThrownBy(() -> service.getMessageTrace("instance-a", " "))
+ assertThatThrownBy(() -> service.getMessageTrace("instance-a", " ",
null))
.isInstanceOf(BusinessException.class)
.hasMessage("msgId is required");
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
index 0d8ef275..17c4f72f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunInstanceProviderTest.java
@@ -345,7 +345,7 @@ class AliyunInstanceProviderTest {
.build();
when(asyncClient.getTrace(any())).thenReturn(CompletableFuture.completedFuture(response));
- TraceRecordVO trace = provider.getMessageTrace(STUDIO_INSTANCE_ID,
"msg-1");
+ TraceRecordVO trace = provider.getMessageTrace(STUDIO_INSTANCE_ID,
"msg-1", "orders");
assertThat(trace.getNodes()).hasSize(3);
TraceNodeVO producer = trace.getNodes().get(0);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index 159572c1..1fa3d303 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -225,7 +225,7 @@ class RocketMQMessageProviderTest {
when(adminExt.queryMessage(anyString(), anyString(), anyInt(),
anyLong(), anyLong()))
.thenReturn(queryResult);
- TraceRecordVO record = provider.getMessageTrace("instance-a",
"msg-123");
+ TraceRecordVO record = provider.getMessageTrace("instance-a",
"msg-123", "orders");
assertThat(record.getNodes()).hasSize(2);
verify(queryHistoryService).recordTraceQuery(eq("instance-a"),
eq("msg-123"), eq(null), eq(2), eq(1));
@@ -257,7 +257,7 @@ class RocketMQMessageProviderTest {
when(adminExt.queryMessage(anyString(), anyString(), anyInt(),
anyLong(), anyLong()))
.thenReturn(queryResult);
- TraceRecordVO record = provider.getMessageTrace("instance-a",
"msg-tx");
+ TraceRecordVO record = provider.getMessageTrace("instance-a",
"msg-tx", "orders");
assertThat(record.getNodes()).hasSize(1);
TraceNodeVO transaction = record.getNodes().get(0);
@@ -274,7 +274,7 @@ class RocketMQMessageProviderTest {
when(adminExt.queryMessage(anyString(), anyString(), anyInt(),
anyLong(), anyLong()))
.thenThrow(new IllegalStateException("broker unavailable"));
- assertThatThrownBy(() -> provider.getMessageTrace("instance-a",
"msg-123"))
+ assertThatThrownBy(() -> provider.getMessageTrace("instance-a",
"msg-123", "orders"))
.isInstanceOf(BusinessException.class)
.hasMessage("Failed to query message trace: broker
unavailable")
.satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(502));
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 d1870695..e4c204c6 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
@@ -32,7 +32,6 @@ import
org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.persistence.entity.RmqGroup;
import org.apache.rocketmq.studio.persistence.mapper.RmqGroupMapper;
import org.apache.rocketmq.studio.persistence.mapper.RmqTopicMapper;
-import org.apache.rocketmq.remoting.protocol.body.GroupList;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
@@ -40,7 +39,6 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.util.HashSet;
import java.util.List;
-import java.util.HashSet;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
index de79ba53..8dd8fb3b 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProviderTest.java
@@ -16,19 +16,42 @@
*/
package org.apache.rocketmq.studio.provider.tencent;
+import com.tencentcloudapi.trocket.v20230308.models.ConsumeGroupItem;
+import com.tencentcloudapi.trocket.v20230308.models.CreateConsumerGroupRequest;
import com.tencentcloudapi.trocket.v20230308.models.CreateTopicRequest;
+import com.tencentcloudapi.trocket.v20230308.models.DeleteConsumerGroupRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeConsumerGroupListResponse;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeConsumerGroupResponse;
+import com.tencentcloudapi.trocket.v20230308.models.DescribeMessageListRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeMessageListResponse;
+import com.tencentcloudapi.trocket.v20230308.models.DescribeMessageRequest;
+import com.tencentcloudapi.trocket.v20230308.models.DescribeMessageResponse;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeMessageTraceRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeMessageTraceResponse;
+import com.tencentcloudapi.trocket.v20230308.models.MessageItem;
+import com.tencentcloudapi.trocket.v20230308.models.MessageTraceItem;
+import
com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListByGroupResponse;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicListResponse;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicRequest;
import com.tencentcloudapi.trocket.v20230308.models.DescribeTopicResponse;
import com.tencentcloudapi.trocket.v20230308.models.ModifyTopicRequest;
+import
com.tencentcloudapi.trocket.v20230308.models.ResetConsumerGroupOffsetRequest;
import com.tencentcloudapi.trocket.v20230308.models.SubscriptionData;
import com.tencentcloudapi.trocket.v20230308.models.TopicItem;
import com.tencentcloudapi.trocket.v20230308.TrocketClient;
+import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
+import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
import org.apache.rocketmq.studio.common.domain.enums.TopicType;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.instance.InstanceRepository;
import org.apache.rocketmq.studio.instance.InstanceVO;
+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.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
+import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
import org.apache.rocketmq.studio.instance.topic.TopicVO;
import org.junit.jupiter.api.BeforeEach;
@@ -270,6 +293,103 @@ class TencentInstanceProviderTest {
.containsOnly(100L);
}
+ @Test
+ void listConsumerGroupsShouldMapAndFilterTest() throws Exception {
+ ConsumeGroupItem one = new ConsumeGroupItem();
+ one.setConsumerGroup("GID_test");
+ one.setMaxRetryTimes(10L);
+ one.setConsumeMessageOrderly(false);
+ ConsumeGroupItem two = new ConsumeGroupItem();
+ two.setConsumerGroup("GID_orders");
+ two.setMaxRetryTimes(16L);
+ DescribeConsumerGroupListResponse response = new
DescribeConsumerGroupListResponse();
+ response.setData(new ConsumeGroupItem[]{one, two});
+ when(client.DescribeConsumerGroupList(any())).thenReturn(response);
+ DescribeConsumerGroupResponse detail = new
DescribeConsumerGroupResponse();
+ detail.setCreatedTime(1600000000000L);
+ detail.setConsumeModel("CLUSTERING");
+ when(client.DescribeConsumerGroup(any())).thenReturn(detail);
+
+ List<ConsumerGroupVO> groups =
provider.listConsumerGroups(STUDIO_INSTANCE_ID, "orders");
+
+ assertThat(groups).hasSize(1);
+ assertThat(groups.get(0).getName()).isEqualTo("GID_orders");
+
assertThat(groups.get(0).getInstanceId()).isEqualTo(STUDIO_INSTANCE_ID);
+ assertThat(groups.get(0).getRetryMaxTimes()).isEqualTo(16);
+ assertThat(groups.get(0).getCreatedAt()).isNotNull();
+
assertThat(groups.get(0).getConsumeType()).isEqualTo(ConsumeType.CLUSTERING);
+ assertThat(groups.get(0).getInstances()).isNotNull().isEmpty();
+ }
+
+ @Test
+ void createConsumerGroupShouldCallTencentOpenApiTest() throws Exception {
+ when(client.CreateConsumerGroup(any())).thenReturn(null);
+ ConsumerGroupVO group = new ConsumerGroupVO();
+ group.setName("GID_new");
+ group.setRetryMaxTimes(20);
+
+ ConsumerGroupVO created =
provider.createConsumerGroup(STUDIO_INSTANCE_ID, group);
+
+ ArgumentCaptor<CreateConsumerGroupRequest> captor =
ArgumentCaptor.forClass(CreateConsumerGroupRequest.class);
+ verify(client).CreateConsumerGroup(captor.capture());
+
assertThat(captor.getValue().getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(captor.getValue().getConsumerGroup()).isEqualTo("GID_new");
+ assertThat(captor.getValue().getMaxRetryTimes()).isEqualTo(20L);
+ assertThat(created.getInstanceId()).isEqualTo(STUDIO_INSTANCE_ID);
+ assertThat(created.getRetryMaxTimes()).isEqualTo(20);
+ }
+
+ @Test
+ void deleteConsumerGroupShouldCallTencentOpenApiTest() throws Exception {
+ when(client.DeleteConsumerGroup(any())).thenReturn(null);
+
+ provider.deleteConsumerGroup(STUDIO_INSTANCE_ID, "GID_test");
+
+ ArgumentCaptor<DeleteConsumerGroupRequest> captor =
ArgumentCaptor.forClass(DeleteConsumerGroupRequest.class);
+ verify(client).DeleteConsumerGroup(captor.capture());
+
assertThat(captor.getValue().getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(captor.getValue().getConsumerGroup()).isEqualTo("GID_test");
+ }
+
+ @Test
+ void getGroupProgressAndSubscriptionsShouldMapSubscriptionDataTest()
throws Exception {
+ SubscriptionData subscription = new SubscriptionData();
+ subscription.setTopic("orders");
+ subscription.setSubString("*");
+ subscription.setExpressionType("TAG");
+ subscription.setConsumerLag(42L);
+ subscription.setConsistency(0L);
+ DescribeTopicListByGroupResponse response = new
DescribeTopicListByGroupResponse();
+ response.setData(new SubscriptionData[]{subscription});
+ when(client.DescribeTopicListByGroup(any())).thenReturn(response);
+
+ List<QueueProgressVO> progress =
provider.getGroupProgress(STUDIO_INSTANCE_ID, "GID_test");
+ assertThat(progress).hasSize(1);
+ assertThat(progress.get(0).getBroker()).isEqualTo("topic:orders");
+ assertThat(progress.get(0).getDiffTotal()).isEqualTo(42L);
+
+ List<SubscriptionEntryVO> subscriptions =
provider.getGroupSubscriptions(STUDIO_INSTANCE_ID, "GID_test");
+ assertThat(subscriptions).hasSize(1);
+ assertThat(subscriptions.get(0).getTopic()).isEqualTo("orders");
+ assertThat(subscriptions.get(0).getExpression()).isEqualTo("*");
+ assertThat(subscriptions.get(0).getType()).isEqualTo("TAG");
+ }
+
+ @Test
+ void resetOffsetShouldCallTencentOpenApiTest() throws Exception {
+ when(client.ResetConsumerGroupOffset(any())).thenReturn(null);
+
+ provider.resetOffset(STUDIO_INSTANCE_ID, "GID_test", 1600000000000L,
"orders");
+
+ ArgumentCaptor<ResetConsumerGroupOffsetRequest> captor =
+ ArgumentCaptor.forClass(ResetConsumerGroupOffsetRequest.class);
+ verify(client).ResetConsumerGroupOffset(captor.capture());
+
assertThat(captor.getValue().getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(captor.getValue().getConsumerGroup()).isEqualTo("GID_test");
+ assertThat(captor.getValue().getTopic()).isEqualTo("orders");
+
assertThat(captor.getValue().getResetTimestamp()).isEqualTo(1600000000000L);
+ }
+
private static TopicItem topicItem(String name, String type, long
queueNum) {
TopicItem item = new TopicItem();
item.setTopic(name);
@@ -285,4 +405,145 @@ class TencentInstanceProviderTest {
subscription.setMessageModel("CLUSTERING");
return subscription;
}
+
+ @Test
+ void queryMessagesByMsgIdShouldReturnDetailWithBodyTest() throws Exception
{
+ DescribeMessageResponse detail = new DescribeMessageResponse();
+ detail.setMessageId("MSG-1");
+ detail.setShowTopicName("orders");
+ detail.setBody("hello body");
+ detail.setProducerAddr("1.2.3.4:5000");
+ detail.setProduceTime("2024-09-12 14:06:55,591");
+
detail.setProperties("{\"UNIQ_KEY\":\"MSG-1\",\"TAGS\":\"tagA\",\"KEYS\":\"keyA\",\"__CLIENT_HOST\":\"1.2.3.4\"}");
+ when(client.DescribeMessage(any())).thenReturn(detail);
+
+ List<MessageRecordVO> messages =
+ provider.queryMessages(STUDIO_INSTANCE_ID, "orders", "MSG-1",
null, null, null, null);
+
+ assertThat(messages).hasSize(1);
+ MessageRecordVO record = messages.get(0);
+ assertThat(record.getMsgId()).isEqualTo("MSG-1");
+ assertThat(record.getTopic()).isEqualTo("orders");
+ assertThat(record.getBody()).isEqualTo("hello body");
+ assertThat(record.getBornHost()).isEqualTo("1.2.3.4:5000");
+ assertThat(record.getTag()).isEqualTo("tagA");
+ assertThat(record.getKey()).isEqualTo("keyA");
+ assertThat(record.getProperties()).containsEntry("UNIQ_KEY", "MSG-1");
+ ArgumentCaptor<DescribeMessageRequest> captor =
ArgumentCaptor.forClass(DescribeMessageRequest.class);
+ verify(client).DescribeMessage(captor.capture());
+
assertThat(captor.getValue().getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(captor.getValue().getTopic()).isEqualTo("orders");
+ assertThat(captor.getValue().getMsgId()).isEqualTo("MSG-1");
+ }
+
+ @Test
+ void queryMessagesByTopicShouldUseMessageListTest() throws Exception {
+ MessageItem one = new MessageItem();
+ one.setMsgId("MSG-A");
+ one.setTags("tagA");
+ one.setKeys("keyA");
+ one.setProducerAddr("1.2.3.4:5000");
+ one.setProduceTime("2024-09-12 14:06:55,591");
+ DescribeMessageListResponse response = new
DescribeMessageListResponse();
+ response.setData(new MessageItem[]{one});
+ response.setTotalCount(1L);
+ when(client.DescribeMessageList(any())).thenReturn(response);
+
+ List<MessageRecordVO> messages =
provider.queryMessages(STUDIO_INSTANCE_ID, "orders", null,
+ "tagA", "keyA", 1600000000000L, 1600001000000L);
+
+ assertThat(messages).hasSize(1);
+ MessageRecordVO record = messages.get(0);
+ assertThat(record.getMsgId()).isEqualTo("MSG-A");
+ assertThat(record.getTag()).isEqualTo("tagA");
+ assertThat(record.getKey()).isEqualTo("keyA");
+ assertThat(record.getBornHost()).isEqualTo("1.2.3.4:5000");
+ assertThat(record.getStoreTime()).isGreaterThan(0L);
+ ArgumentCaptor<DescribeMessageListRequest> captor =
ArgumentCaptor.forClass(DescribeMessageListRequest.class);
+ verify(client).DescribeMessageList(captor.capture());
+
assertThat(captor.getValue().getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(captor.getValue().getTopic()).isEqualTo("orders");
+ assertThat(captor.getValue().getMsgKey()).isEqualTo("keyA");
+ assertThat(captor.getValue().getTag()).isEqualTo("tagA");
+ assertThat(captor.getValue().getTaskRequestId()).isNotBlank();
+ assertThat(captor.getValue().getOffset()).isEqualTo(0L);
+ }
+
+ @Test
+ void queryMessagesByTopicShouldPageWithSameTaskRequestIdTest() throws
Exception {
+ // Page 1 returns a full page (MESSAGE_LIMIT), forcing a second page;
page 2 is a short
+ // page, which terminates the loop.
+ MessageItem[] page1Items = new
MessageItem[TencentInstanceProvider.MESSAGE_LIMIT];
+ for (int i = 0; i < page1Items.length; i++) {
+ MessageItem item = new MessageItem();
+ item.setMsgId("MSG-" + (i + 1));
+ item.setProduceTime("2024-09-12 14:06:55,591");
+ page1Items[i] = item;
+ }
+ MessageItem last = new MessageItem();
+ last.setMsgId("MSG-LAST");
+ last.setProduceTime("2024-09-12 14:06:56,591");
+ DescribeMessageListResponse page1 = new DescribeMessageListResponse();
+ page1.setData(page1Items);
+ DescribeMessageListResponse page2 = new DescribeMessageListResponse();
+ page2.setData(new MessageItem[]{last});
+ when(client.DescribeMessageList(any()))
+ .thenReturn(page1)
+ .thenReturn(page2);
+
+ List<MessageRecordVO> messages =
provider.queryMessages(STUDIO_INSTANCE_ID, "orders", null,
+ null, null, 1600000000000L, 1600001000000L);
+
+ assertThat(messages).hasSize(TencentInstanceProvider.MESSAGE_LIMIT +
1);
+ assertThat(messages.get(0).getMsgId()).isEqualTo("MSG-1");
+ assertThat(messages.get(TencentInstanceProvider.MESSAGE_LIMIT -
1).getMsgId())
+ .isEqualTo("MSG-" + TencentInstanceProvider.MESSAGE_LIMIT);
+
assertThat(messages.get(TencentInstanceProvider.MESSAGE_LIMIT).getMsgId()).isEqualTo("MSG-LAST");
+ ArgumentCaptor<DescribeMessageListRequest> captor =
ArgumentCaptor.forClass(DescribeMessageListRequest.class);
+ verify(client,
org.mockito.Mockito.times(2)).DescribeMessageList(captor.capture());
+ java.util.List<DescribeMessageListRequest> requests =
captor.getAllValues();
+ assertThat(requests).hasSize(2);
+ assertThat(requests.get(0).getOffset()).isEqualTo(0L);
+ assertThat(requests.get(1).getOffset()).isEqualTo((long)
TencentInstanceProvider.MESSAGE_LIMIT);
+ // A single logical query reuses the same TaskRequestId across pages;
only offset advances.
+ assertThat(requests.get(1).getTaskRequestId())
+ .isEqualTo(requests.get(0).getTaskRequestId())
+ .isNotBlank();
+ }
+
+ @Test
+ void getMessageTraceShouldMapStagesTest() throws Exception {
+ MessageTraceItem produce = new MessageTraceItem();
+ produce.setStage("produce");
+
produce.setData("{\"MsgId\":\"MSG-1\",\"Status\":0,\"ProduceTime\":\"2024-09-12
14:06:55,591\","
+ + "\"ProducerAddr\":\"1.2.3.4:5000\",\"Duration\":2}");
+ MessageTraceItem consume = new MessageTraceItem();
+ consume.setStage("consume");
+
consume.setData("{\"TotalCount\":1,\"RocketMqConsumeLogs\":[{\"MsgId\":\"MSG-1\",\"Status\":2,"
+ + "\"PushTime\":\"2024-09-12
14:06:55,600\",\"ConsumerGroup\":\"GID_test\",\"RetryTimes\":1}]}");
+ DescribeMessageTraceResponse response = new
DescribeMessageTraceResponse();
+ response.setData(new MessageTraceItem[]{produce, consume});
+ when(client.DescribeMessageTrace(any())).thenReturn(response);
+
+ TraceRecordVO trace = provider.getMessageTrace(STUDIO_INSTANCE_ID,
"MSG-1", "orders");
+
+ assertThat(trace.getNodes()).hasSize(2);
+ TraceNodeVO produceNode = trace.getNodes().get(0);
+ assertThat(produceNode.getTitle()).isEqualTo("produce");
+ assertThat(produceNode.getStatus()).isEqualTo("finish");
+ assertThat(produceNode.getCostTime()).isEqualTo(2L);
+ assertThat(produceNode.getTimestamp()).isGreaterThan(0L);
+ TraceNodeVO consumeNode = trace.getNodes().get(1);
+ assertThat(consumeNode.getTitle()).isEqualTo("consume");
+ assertThat(consumeNode.getStatus()).isEqualTo("finish");
+ assertThat(trace.getConsumerStatus()).hasSize(1);
+
assertThat(trace.getConsumerStatus().get(0).getGroup()).isEqualTo("GID_test");
+
assertThat(trace.getConsumerStatus().get(0).getDeliveryStatus()).isEqualTo(DeliveryStatus.success);
+
assertThat(trace.getConsumerStatus().get(0).getRetryCount()).isEqualTo(1);
+ ArgumentCaptor<DescribeMessageTraceRequest> captor =
ArgumentCaptor.forClass(DescribeMessageTraceRequest.class);
+ verify(client).DescribeMessageTrace(captor.capture());
+
assertThat(captor.getValue().getInstanceId()).isEqualTo(CLOUD_INSTANCE_ID);
+ assertThat(captor.getValue().getTopic()).isEqualTo("orders");
+ assertThat(captor.getValue().getMsgId()).isEqualTo("MSG-1");
+ }
}
diff --git a/web/src/api/message.ts b/web/src/api/message.ts
index f70d15d2..83f535a7 100644
--- a/web/src/api/message.ts
+++ b/web/src/api/message.ts
@@ -80,10 +80,10 @@ export async function queryMessages(params: MessageQuery) {
return sortMessagesByStoreTimeDesc(res.data.data);
}
-export async function getMessageTrace(msgId: string, instanceId?: string) {
+export async function getMessageTrace(msgId: string, instanceId?: string,
topic?: string) {
const res = await client.get<{ data: TraceRecord }>(
`/messages/${encodeURIComponent(msgId)}/trace`,
- { params: { instanceId } },
+ { params: { instanceId, topic } },
);
return res.data.data;
}
diff --git a/web/src/pages/instance/message.tsx
b/web/src/pages/instance/message.tsx
index 5de86438..cc2eaf5e 100644
--- a/web/src/pages/instance/message.tsx
+++ b/web/src/pages/instance/message.tsx
@@ -458,7 +458,7 @@ const MessagePageContent = ({
setTraceLoading(true);
setTraceError(null);
try {
- const result = await getMessageTrace(record.msgId, selectedInstanceId);
+ const result = await getMessageTrace(record.msgId, selectedInstanceId,
record.topic);
if (traceGenerationRef.current !== requestGeneration) return;
setTraceData(result);
setTraceError(null);
diff --git a/web/src/services/messageService.ts
b/web/src/services/messageService.ts
index 9e08e612..81a0de42 100644
--- a/web/src/services/messageService.ts
+++ b/web/src/services/messageService.ts
@@ -51,12 +51,13 @@ export async function queryMessages(params: MessageQuery):
Promise<MessageRecord
export async function getMessageTrace(
msgId: string,
instanceId?: string,
+ topic?: string,
): Promise<TraceRecord | null> {
if (isMockMode()) {
const trace = mockMessageTraces[msgId] as unknown as TraceRecord |
undefined;
return trace ? cloneTrace(trace) : null;
}
- return messageApi.getMessageTrace(msgId, instanceId);
+ return messageApi.getMessageTrace(msgId, instanceId, topic);
}
export async function listDLQGroups(instanceId: string): Promise<DLQGroup[]> {