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 99692b307 fix(message): report truncation for Apache key queries 
(#4185)
99692b307 is described below

commit 99692b30732d33f9949a931d1d98a5e3585adc64
Author: Zhao Jianing <[email protected]>
AuthorDate: Tue Sep 15 19:40:53 2026 +0800

    fix(message): report truncation for Apache key queries (#4185)
    
    RocketMQMessageProvider never overrides queryMessagesDetailed, so every
    Apache message query is wrapped as complete(...) by the interface
    default: a key query capped at the broker-side budget (KEY_QUERY_MAX =
    64) returns exactly 64 rows with resultMayBeTruncated=false and no UI
    warning, while Tencent (#3048) and Aliyun (#4164) report the same
    condition. Apache is the default deployment, so the widest audience is
    left with silently incomplete key query results.
    
    MQAdminImpl fans the key query out to every route broker with a
    per-broker budget of KEY_QUERY_MAX and merges the responses without a
    client-side cap, so an exact verdict is not derivable from the merged
    list; report mayBeTruncated when the merged count reaches the budget
    (checked before the client-side tag filter, which can otherwise hide the
    capped broker results). ApacheInstanceProvider now delegates
    queryMessagesDetailed so registered Apache instances reach the signal
    instead of the complete(...) default.
    
    Signed-off-by: zjncs <[email protected]>
---
 .../provider/apache/ApacheInstanceProvider.java    |  7 +++
 .../provider/apache/RocketMQMessageProvider.java   | 35 ++++++++----
 .../apache/ApacheInstanceProviderTest.java         | 15 ++++++
 .../apache/RocketMQMessageProviderTest.java        | 63 ++++++++++++++++++++++
 4 files changed, 109 insertions(+), 11 deletions(-)

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 29847bc01..213cb54f9 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
@@ -24,6 +24,7 @@ import 
org.apache.rocketmq.studio.instance.group.QueueProgressVO;
 import org.apache.rocketmq.studio.instance.group.ResetConsumerOffsetPreviewVO;
 import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
 import org.apache.rocketmq.studio.instance.message.MessageProvider;
+import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
 import org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
 import 
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageResultVO;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
@@ -171,6 +172,12 @@ public class ApacheInstanceProvider implements 
InstanceProvider {
         return messageProvider.queryMessages(instanceId, topic, msgId, tag, 
key, startTime, endTime);
     }
 
+    @Override
+    public MessageQueryResult queryMessagesDetailed(String instanceId, String 
topic, String msgId,
+                                                    String tag, String key, 
Long startTime, Long endTime) {
+        return messageProvider.queryMessagesDetailed(instanceId, topic, msgId, 
tag, key, startTime, endTime);
+    }
+
     @Override
     public TraceRecordVO getMessageTrace(String instanceId, String msgId, 
String topic) {
         return messageProvider.getMessageTrace(instanceId, msgId, topic);
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 640f500bd..b22ebee38 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
@@ -36,6 +36,7 @@ import org.apache.rocketmq.studio.common.util.MqResponseCodes;
 import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
 import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
 import org.apache.rocketmq.studio.instance.message.MessageProvider;
+import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
 import org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
 import 
org.apache.rocketmq.studio.instance.message.DirectConsumeMessageResultVO;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
@@ -103,17 +104,23 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     @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) {
         return runtimeAdminClientResolver.execute(instanceId,
                 adminExt -> queryMessages(instanceId, (DefaultMQAdminExt) 
adminExt, topic, msgId, tag, key,
                         startTime, endTime));
     }
 
-    private List<MessageRecordVO> queryMessages(String instanceId, 
DefaultMQAdminExt adminExt,
-                                                 String topic, String msgId, 
String tag, String key,
-                                                 Long startTime, Long endTime) 
{
+    private MessageQueryResult queryMessages(String instanceId, 
DefaultMQAdminExt adminExt,
+                                             String topic, String msgId, 
String tag, String key,
+                                             Long startTime, Long endTime) {
 
         if (StringUtils.hasText(msgId)) {
-            return queryByMsgId(adminExt, topic, msgId);
+            return MessageQueryResult.complete(queryByMsgId(adminExt, topic, 
msgId));
         }
 
         long end = endTime != null ? endTime : System.currentTimeMillis();
@@ -128,11 +135,12 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             if (begin >= 0 && end >= 0 && end - begin > 
MAX_TOPIC_QUERY_WINDOW_MILLIS) {
                 throw new BusinessException(400, "Topic message query time 
range must not exceed 7 days");
             }
-            return queryByTopic(instanceId, topic, tag, begin, end, 
DEFAULT_TOPIC_LIMIT);
+            return MessageQueryResult.complete(
+                    queryByTopic(instanceId, topic, tag, begin, end, 
DEFAULT_TOPIC_LIMIT));
         }
 
         log.warn("queryMessages requires at least one of msgId/topic, 
returning empty list");
-        return Collections.emptyList();
+        return MessageQueryResult.complete(Collections.emptyList());
     }
 
     private List<MessageRecordVO> queryByMsgId(DefaultMQAdminExt adminExt, 
String topic, String msgId) {
@@ -176,27 +184,32 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         }
     }
 
-    private List<MessageRecordVO> queryByKey(DefaultMQAdminExt adminExt, 
String topic, String key,
-                                             String tag, long begin, long end) 
{
+    private MessageQueryResult queryByKey(DefaultMQAdminExt adminExt, String 
topic, String key,
+                                          String tag, long begin, long end) {
         try {
             QueryResult queryResult = adminExt.queryMessage(topic, key, 
KEY_QUERY_MAX, begin, end);
             if (queryResult == null || queryResult.getMessageList() == null) {
-                return Collections.emptyList();
+                return MessageQueryResult.complete(Collections.emptyList());
             }
+            // MQAdminImpl fans the key query out to every route broker with a 
per-broker
+            // budget of KEY_QUERY_MAX and merges the responses without a 
client-side cap,
+            // so a merged count at or above the budget means some broker may 
have stopped
+            // at its cap. Surface that instead of silently dropping further 
matches.
+            boolean mayBeTruncated = queryResult.getMessageList().size() >= 
KEY_QUERY_MAX;
             List<MessageRecordVO> result = new ArrayList<>();
             for (MessageExt messageExt : queryResult.getMessageList()) {
                 if (matchesTag(messageExt, tag)) {
                     result.add(toRecordVO(messageExt));
                 }
             }
-            return result;
+            return mayBeTruncated ? MessageQueryResult.truncated(result) : 
MessageQueryResult.complete(result);
         } catch (Exception e) {
             if (MqResponseCodes.hasResponseCode(e, ResponseCode.NO_MESSAGE)) {
                 // MQAdminImpl.queryMessage throws 
MQClientException(NO_MESSAGE) instead of
                 // returning an empty QueryResult when the key matches 
nothing: the query
                 // completed, so the correct response is an empty list, not a 
gateway error.
                 log.info("queryMessage(topic={}, key={}) matched nothing", 
topic, key);
-                return Collections.emptyList();
+                return MessageQueryResult.complete(Collections.emptyList());
             }
             log.warn("queryMessage(topic={}, key={}) failed: {}", topic, key, 
e.getMessage());
             throw new BusinessException(502, "Failed to query messages by key: 
" + e.getMessage());
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
index e1fc6994c..b562ec3f7 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ApacheInstanceProviderTest.java
@@ -23,6 +23,7 @@ import org.apache.rocketmq.studio.instance.InstanceVO;
 import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
 import org.apache.rocketmq.studio.instance.topic.TopicVO;
 import org.apache.rocketmq.studio.instance.message.MessageProvider;
+import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
 import org.apache.rocketmq.studio.provider.InstanceCapability;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
@@ -89,6 +90,20 @@ class ApacheInstanceProviderTest {
         assertThat(provider.vendor()).isEqualTo(InstanceVendor.APACHE);
     }
 
+    @Test
+    void queryMessagesDetailedShouldDelegateToMessageProviderTest() {
+        MessageQueryResult result = 
MessageQueryResult.truncated(java.util.List.of());
+        when(messageProvider.queryMessagesDetailed(
+                "inst-1", "TopicA", null, null, "order-1", 100L, 200L))
+                .thenReturn(result);
+
+        assertThat(provider.queryMessagesDetailed(
+                "inst-1", "TopicA", null, null, "order-1", 100L, 200L))
+                .isSameAs(result);
+        verify(messageProvider).queryMessagesDetailed(
+                "inst-1", "TopicA", null, null, "order-1", 100L, 200L);
+    }
+
     @Test
     void countTopicsShouldDelegateToRepositoryTest() {
         InstanceVO instance = InstanceVO.builder().name("inst-1").build();
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 c794ed919..da307f9b1 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
@@ -42,6 +42,7 @@ import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
+import org.apache.rocketmq.studio.instance.message.MessageQueryResult;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
 import org.apache.rocketmq.studio.instance.message.DirectConsumeMessageDTO;
 import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
@@ -303,6 +304,68 @@ class RocketMQMessageProviderTest {
         return mqAdmin;
     }
 
+    @Test
+    void queryByKeyReportsTruncationWhenTheBrokerBudgetIsReachedTest() throws 
Exception {
+        when(adminExt.queryMessage("TopicA", "order-1", 64, 100L, 200L))
+                .thenReturn(new QueryResult(0L, keyQueryMatches(64)));
+
+        MessageQueryResult result = provider.queryMessagesDetailed(
+                "instance-a", "TopicA", null, null, "order-1", 100L, 200L);
+
+        assertThat(result.mayBeTruncated()).isTrue();
+        assertThat(result.messages()).hasSize(64);
+    }
+
+    @Test
+    void queryByKeyStaysCompleteBelowTheBrokerBudgetTest() throws Exception {
+        when(adminExt.queryMessage("TopicA", "order-1", 64, 100L, 200L))
+                .thenReturn(new QueryResult(0L, keyQueryMatches(63)));
+
+        MessageQueryResult result = provider.queryMessagesDetailed(
+                "instance-a", "TopicA", null, null, "order-1", 100L, 200L);
+
+        assertThat(result.mayBeTruncated()).isFalse();
+        assertThat(result.messages()).hasSize(63);
+    }
+
+    @Test
+    void queryByKeyKeepsTheTruncationSignalWhenTagFilteringDropsEveryRowTest() 
throws Exception {
+        when(adminExt.queryMessage("TopicA", "order-1", 64, 100L, 200L))
+                .thenReturn(new QueryResult(0L, keyQueryMatches(64)));
+
+        MessageQueryResult result = provider.queryMessagesDetailed(
+                "instance-a", "TopicA", null, "TagA", "order-1", 100L, 200L);
+
+        assertThat(result.mayBeTruncated()).isTrue();
+        assertThat(result.messages()).isEmpty();
+    }
+
+    @Test
+    void queryByKeyReportsCompleteWhenTheClientReportsNoMessageTest() throws 
Exception {
+        // The NO_MESSAGE degradation must travel through the Detailed path 
too, otherwise a
+        // key that matches nothing would surface as possibly truncated 
instead of complete.
+        when(adminExt.queryMessage("TopicA", "order-1", 64, 100L, 200L))
+                .thenThrow(new MQClientException(ResponseCode.NO_MESSAGE,
+                        "query message by key finished, but no message."));
+
+        MessageQueryResult result = provider.queryMessagesDetailed(
+                "instance-a", "TopicA", null, null, "order-1", 100L, 200L);
+
+        assertThat(result.mayBeTruncated()).isFalse();
+        assertThat(result.messages()).isEmpty();
+    }
+
+    private static List<MessageExt> keyQueryMatches(int count) {
+        return IntStream.range(0, count)
+                .mapToObj(index -> {
+                    MessageExt message = new MessageExt();
+                    message.setTopic("TopicA");
+                    message.setMsgId("msg-" + index);
+                    return message;
+                })
+                .toList();
+    }
+
     @Test
     void queryByMsgIdUsesDecodedPhysicalOffsetForFallback() throws Exception {
         String msgId = "AC1E0A6400002A9F0000000001A3F2B1";

Reply via email to