RockteMQ-AI commented on code in PR #1641:
URL:
https://github.com/apache/rocketmq-dashboard/pull/1641#discussion_r3755771277
##########
server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageController.java:
##########
@@ -46,7 +46,8 @@ public Result<List<MessageRecordVO>> queryMessages(
}
@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,
Review Comment:
The `topic` parameter is `required = false`, which is correct for backward
compatibility with Aliyun/Apache providers. However, the Tencent provider will
throw if `topic` is null (`requireTopic`). This is fine since the frontend
always passes it, but consider documenting this contract in the Javadoc or
adding a more descriptive error message when topic is missing for Tencent
instances.
##########
server/src/main/java/org/apache/rocketmq/studio/provider/tencent/TencentInstanceProvider.java:
##########
@@ -229,43 +274,439 @@ public List<TopicConsumerVO> getTopicConsumers(String
instanceId, String topicNa
@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);
Review Comment:
Minor readability suggestion: the break condition `returned == 0 || returned
< MESSAGE_LIMIT || total > 0 && result.size() >= total` relies on `&&` having
higher precedence than `||`. Consider adding parentheses around `total > 0 &&
result.size() >= total` to make the intent explicit:
```java
if (returned == 0 || returned < MESSAGE_LIMIT || (total > 0 && result.size()
>= total)) {
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]