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();

Reply via email to