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 25132a0d8 fix(aliyun): report truncated message queries (#4164)
25132a0d8 is described below
commit 25132a0d87af3c3ab6e8a1727573bd321954e9a3
Author: aias00 <[email protected]>
AuthorDate: Tue Sep 15 19:16:03 2026 +0800
fix(aliyun): report truncated message queries (#4164)
Signed-off-by: liuhy <[email protected]>
---
.../provider/alibaba/AliyunInstanceProvider.java | 20 ++++-
.../alibaba/AliyunInstanceProviderTest.java | 90 ++++++++++++++++++++++
2 files changed, 108 insertions(+), 2 deletions(-)
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 08d7e8070..876e42261 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
@@ -54,6 +54,7 @@ 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.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
import org.apache.rocketmq.studio.instance.topic.TopicConsumerVO;
@@ -455,8 +456,16 @@ public class AliyunInstanceProvider implements
InstanceProvider {
@Override
public List<MessageRecordVO> queryMessages(String instanceId, String
topic, String msgId,
String tag, String key, Long
startTime, Long endTime) {
+ return queryMessagesDetailed(instanceId, topic, msgId, tag, key,
startTime, endTime).messages();
+ }
+
+ @Override
+ public MessageQueryResult queryMessagesDetailed(String instanceId, String
topic, String msgId,
+ String tag, String key,
Long startTime, Long endTime) {
Context ctx = resolve(instanceId);
List<MessageRecordVO> records = new ArrayList<>();
+ int fetched = 0;
+ boolean mayBeTruncated = false;
for (int page = 1; page <= AliyunConverters.MESSAGE_MAX_PAGES; page++)
{
ListMessagesRequest.Builder builder = ListMessagesRequest.builder()
.instanceId(ctx.cloudInstanceId())
@@ -486,6 +495,7 @@ public class AliyunInstanceProvider implements
InstanceProvider {
if (list == null || list.isEmpty()) {
break;
}
+ fetched += list.size();
for (ListMessagesResponseBody.List item : list) {
if (item == null) {
continue;
@@ -495,11 +505,17 @@ public class AliyunInstanceProvider implements
InstanceProvider {
records.add(vo);
}
}
- if (list.size() < AliyunConverters.MESSAGE_PAGE_SIZE) {
+ Long totalCount = data.getTotalCount();
+ boolean shortPage = list.size() <
AliyunConverters.MESSAGE_PAGE_SIZE;
+ boolean allFetched = totalCount != null && totalCount > 0L &&
fetched >= totalCount;
+ if (shortPage || allFetched) {
break;
}
+ if (page == AliyunConverters.MESSAGE_MAX_PAGES) {
+ mayBeTruncated = true;
+ }
}
- return records;
+ return mayBeTruncated ? MessageQueryResult.truncated(records) :
MessageQueryResult.complete(records);
}
@Override
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 0ae1136cb..59203a9df 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
@@ -28,6 +28,7 @@ import
com.aliyun.sdk.service.rocketmq20220801.models.GetTraceResponseBody;
import
com.aliyun.sdk.service.rocketmq20220801.models.ListConsumerGroupsRequest;
import
com.aliyun.sdk.service.rocketmq20220801.models.ListConsumerGroupsResponse;
import
com.aliyun.sdk.service.rocketmq20220801.models.ListConsumerGroupsResponseBody;
+import com.aliyun.sdk.service.rocketmq20220801.models.ListMessagesRequest;
import com.aliyun.sdk.service.rocketmq20220801.models.ListMessagesResponse;
import com.aliyun.sdk.service.rocketmq20220801.models.ListMessagesResponseBody;
import com.aliyun.sdk.service.rocketmq20220801.models.ListTopicsRequest;
@@ -47,6 +48,7 @@ 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.ResetConsumerOffsetPreviewVO;
+import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
@@ -497,6 +499,73 @@ class AliyunInstanceProviderTest {
assertThat(filtered.get(0).getMsgId()).isEqualTo("msg-2");
}
+ @Test
+ void queryMessagesDetailedShouldReportTruncationAtPageBudgetTest() {
+ stubInstance();
+ stubCallThrough();
+ when(asyncClient.listMessages(any())).thenAnswer(invocation -> {
+ ListMessagesRequest request = invocation.getArgument(0);
+ return CompletableFuture.completedFuture(messagesResponse(
+ 101L, request.getPageNumber(),
AliyunConverters.MESSAGE_PAGE_SIZE, "tagA"));
+ });
+
+ MessageQueryResult result = provider.queryMessagesDetailed(
+ STUDIO_INSTANCE_ID, "topic-a", null, null, null, null, null);
+
+ assertThat(result.messages()).hasSize(100);
+ assertThat(result.mayBeTruncated()).isTrue();
+ verify(asyncClient,
times(AliyunConverters.MESSAGE_MAX_PAGES)).listMessages(any());
+ }
+
+ @Test
+ void
queryMessagesDetailedShouldPreserveTruncationAfterLocalTagFilterTest() {
+ stubInstance();
+ stubCallThrough();
+ when(asyncClient.listMessages(any())).thenAnswer(invocation -> {
+ ListMessagesRequest request = invocation.getArgument(0);
+ return CompletableFuture.completedFuture(messagesResponse(
+ null, request.getPageNumber(),
AliyunConverters.MESSAGE_PAGE_SIZE, "other-tag"));
+ });
+
+ MessageQueryResult result = provider.queryMessagesDetailed(
+ STUDIO_INSTANCE_ID, "topic-a", null, "wanted-tag", null, null,
null);
+
+ assertThat(result.messages()).isEmpty();
+ assertThat(result.mayBeTruncated()).isTrue();
+ }
+
+ @Test
+ void
queryMessagesDetailedShouldRemainCompleteWhenTotalCountEndsAtBudgetTest() {
+ stubInstance();
+ stubCallThrough();
+ when(asyncClient.listMessages(any())).thenAnswer(invocation -> {
+ ListMessagesRequest request = invocation.getArgument(0);
+ return CompletableFuture.completedFuture(messagesResponse(
+ 100L, request.getPageNumber(),
AliyunConverters.MESSAGE_PAGE_SIZE, "tagA"));
+ });
+
+ MessageQueryResult result = provider.queryMessagesDetailed(
+ STUDIO_INSTANCE_ID, "topic-a", null, null, null, null, null);
+
+ assertThat(result.messages()).hasSize(100);
+ assertThat(result.mayBeTruncated()).isFalse();
+ }
+
+ @Test
+ void queryMessagesDetailedShouldRemainCompleteOnShortPageTest() {
+ stubInstance();
+ stubCallThrough();
+
when(asyncClient.listMessages(any())).thenReturn(CompletableFuture.completedFuture(
+ messagesResponse(null, 1, 3, "tagA")));
+
+ MessageQueryResult result = provider.queryMessagesDetailed(
+ STUDIO_INSTANCE_ID, "topic-a", null, null, null, null, null);
+
+ assertThat(result.messages()).hasSize(3);
+ assertThat(result.mayBeTruncated()).isFalse();
+ verify(asyncClient).listMessages(any());
+ }
+
@Test
void createConsumerGroupShouldApplyDefaultsTest() {
stubInstance();
@@ -763,6 +832,27 @@ class AliyunInstanceProviderTest {
.build();
}
+ private static ListMessagesResponse messagesResponse(Long totalCount, int
pageNumber, int count, String tag) {
+ List<ListMessagesResponseBody.List> rows = IntStream.range(0, count)
+ .mapToObj(index -> ListMessagesResponseBody.List.builder()
+ .messageId("msg-" + pageNumber + "-" + index)
+ .topicName("topic-a")
+ .messageTag(tag)
+ .build())
+ .toList();
+ return ListMessagesResponse.create().toBuilder()
+ .statusCode(200)
+ .body(ListMessagesResponseBody.builder()
+ .data(ListMessagesResponseBody.Data.builder()
+ .list(rows)
+ .pageNumber((long) pageNumber)
+ .pageSize((long)
AliyunConverters.MESSAGE_PAGE_SIZE)
+ .totalCount(totalCount)
+ .build())
+ .build())
+ .build();
+ }
+
@Test
void countTopicsShouldUseTotalCountWithoutFetchingEveryTopicTest() {
stubInstance();