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 3c2639df fix: bound message query inputs and payloads (#1026)
3c2639df is described below

commit 3c2639df4e4e6487989b0ea3378dce569659be54
Author: aias00 <[email protected]>
AuthorDate: Wed Aug 5 02:46:38 2026 -0700

    fix: bound message query inputs and payloads (#1026)
---
 .../studio/instance/message/MessageRecordVO.java   |  3 +
 .../studio/instance/message/MessageService.java    | 18 ++++++
 .../studio/rocketmq/RocketMQMessageProvider.java   | 70 +++++++++++++++++++++-
 .../instance/message/MessageServiceTest.java       | 46 ++++++++++++++
 .../rocketmq/RocketMQMessageProviderTest.java      | 40 ++++++++++++-
 5 files changed, 173 insertions(+), 4 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageRecordVO.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageRecordVO.java
index fc999fbd..6908f916 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageRecordVO.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/message/MessageRecordVO.java
@@ -33,9 +33,12 @@ public class MessageRecordVO {
     private String tag;
     private String key;
     private String body;
+    private String bodyEncoding;
+    private boolean bodyTruncated;
     private long storeTime;
     private String bornHost;
     private String storeHost;
     private Map<String, String> properties;
+    private boolean propertiesTruncated;
     private int size;
 }
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 a2bdafd2..6e7b7030 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
@@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.instance.message;
 
 import lombok.RequiredArgsConstructor;
 import lombok.extern.slf4j.Slf4j;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.springframework.stereotype.Service;
 
 import java.util.List;
@@ -27,10 +28,13 @@ import java.util.List;
 @Slf4j
 public class MessageService {
 
+    private static final long MAX_TOPIC_QUERY_WINDOW_MILLIS = 7L * 24 * 60 * 
60 * 1000;
+
     private final MessageProvider messageProvider;
 
     public List<MessageRecordVO> queryMessages(
             String topic, String msgId, String tag, String key, Long 
startTime, Long endTime) {
+        validateTopicQueryWindow(topic, msgId, key, startTime, endTime);
         log.info("Querying messages: topic={}, msgId={}, tag={}, key={}", 
topic, msgId, tag, key);
         return messageProvider.queryMessages(topic, msgId, tag, key, 
startTime, endTime);
     }
@@ -39,4 +43,18 @@ public class MessageService {
         log.info("Getting message trace: msgId={}", msgId);
         return messageProvider.getMessageTrace(msgId);
     }
+
+    private void validateTopicQueryWindow(String topic, String msgId, String 
key, Long startTime, Long endTime) {
+        if (msgId != null && !msgId.isBlank() || key != null && !key.isBlank() 
|| topic == null || topic.isBlank()) {
+            return;
+        }
+        long end = endTime == null ? System.currentTimeMillis() : endTime;
+        long start = startTime == null ? end - 60 * 60 * 1000L : startTime;
+        if (start > end) {
+            throw new BusinessException(400, "startTime must not be after 
endTime");
+        }
+        if (end - start > MAX_TOPIC_QUERY_WINDOW_MILLIS) {
+            throw new BusinessException(400, "topic query time range must not 
exceed 7 days");
+        }
+    }
 }
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
index 1657d4e6..c8626813 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
@@ -43,10 +43,17 @@ import org.springframework.util.StringUtils;
 
 import java.net.InetSocketAddress;
 import java.net.SocketAddress;
+import java.nio.ByteBuffer;
 import java.nio.charset.StandardCharsets;
+import java.nio.charset.CharacterCodingException;
+import java.nio.charset.CodingErrorAction;
 import java.util.ArrayList;
 import java.util.Collections;
+import java.util.Comparator;
+import java.util.Base64;
+import java.util.LinkedHashMap;
 import java.util.List;
+import java.util.Map;
 import java.util.Set;
 
 /**
@@ -66,6 +73,10 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     private static final int TRACE_QUERY_MAX = 64;
     private static final int DEFAULT_TOPIC_LIMIT = 200;
     private static final int TOPIC_QUERY_HARD_CAP = 2000;
+    private static final int MAX_BODY_DISPLAY_BYTES = 64 * 1024;
+    private static final int MAX_BINARY_BODY_DISPLAY_BYTES = 48 * 1024;
+    private static final int MAX_PROPERTIES = 64;
+    private static final int MAX_PROPERTY_VALUE_CHARS = 1024;
     private static final long ONE_HOUR_MILLIS = 3600_000L;
     private static final long ONE_DAY_MILLIS = 24 * ONE_HOUR_MILLIS;
 
@@ -385,22 +396,75 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         }
     }
 
-    private MessageRecordVO toRecordVO(MessageExt messageExt) {
+    MessageRecordVO toRecordVO(MessageExt messageExt) {
         byte[] body = messageExt.getBody();
+        DisplayBody displayBody = displayBody(body);
+        Map<String, String> properties = messageExt.getProperties();
+        Map<String, String> displayProperties = limitProperties(properties);
         return MessageRecordVO.builder()
                 .msgId(messageExt.getMsgId())
                 .topic(messageExt.getTopic())
                 .tag(messageExt.getTags())
                 .key(messageExt.getKeys())
-                .body(body == null ? null : new String(body, 
StandardCharsets.UTF_8))
+                .body(displayBody.value())
+                .bodyEncoding(displayBody.encoding())
+                .bodyTruncated(displayBody.truncated())
                 .storeTime(messageExt.getStoreTimestamp())
                 .bornHost(String.valueOf(messageExt.getBornHost()))
                 .storeHost(String.valueOf(messageExt.getStoreHost()))
-                .properties(messageExt.getProperties())
+                .properties(displayProperties)
+                .propertiesTruncated(properties != null && 
(displayProperties.size() < properties.size()
+                        || hasOversizedProperty(properties)))
                 .size(messageExt.getStoreSize())
                 .build();
     }
 
+    private DisplayBody displayBody(byte[] body) {
+        if (body == null) {
+            return new DisplayBody(null, null, false);
+        }
+        int textLength = Math.min(body.length, MAX_BODY_DISPLAY_BYTES);
+        try {
+            String value = StandardCharsets.UTF_8.newDecoder()
+                    .onMalformedInput(CodingErrorAction.REPORT)
+                    .onUnmappableCharacter(CodingErrorAction.REPORT)
+                    .decode(ByteBuffer.wrap(body, 0, textLength))
+                    .toString();
+            return new DisplayBody(value, "UTF-8", body.length > textLength);
+        } catch (CharacterCodingException ignored) {
+            int binaryLength = Math.min(body.length, 
MAX_BINARY_BODY_DISPLAY_BYTES);
+            return new DisplayBody(Base64.getEncoder().encodeToString(
+                    java.util.Arrays.copyOf(body, binaryLength)), "BASE64", 
body.length > binaryLength);
+        }
+    }
+
+    private Map<String, String> limitProperties(Map<String, String> 
properties) {
+        if (properties == null || properties.isEmpty()) {
+            return Collections.emptyMap();
+        }
+        Map<String, String> limited = new LinkedHashMap<>();
+        properties.entrySet().stream()
+                
.sorted(Map.Entry.comparingByKey(Comparator.nullsLast(String::compareTo)))
+                .limit(MAX_PROPERTIES)
+                .forEach(entry -> limited.put(entry.getKey(), 
abbreviate(entry.getValue(), MAX_PROPERTY_VALUE_CHARS)));
+        return limited;
+    }
+
+    private boolean hasOversizedProperty(Map<String, String> properties) {
+        return properties != null && properties.values().stream()
+                .anyMatch(value -> value != null && value.length() > 
MAX_PROPERTY_VALUE_CHARS);
+    }
+
+    private String abbreviate(String value, int maxLength) {
+        if (value == null || value.length() <= maxLength) {
+            return value;
+        }
+        return value.substring(0, maxLength) + "...";
+    }
+
+    private record DisplayBody(String value, String encoding, boolean 
truncated) {
+    }
+
     private boolean matchesTag(MessageExt messageExt, String tag) {
         if (!StringUtils.hasText(tag) || "*".equals(tag)) {
             return true;
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
new file mode 100644
index 00000000..d4045388
--- /dev/null
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/message/MessageServiceTest.java
@@ -0,0 +1,46 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ */
+package org.apache.rocketmq.studio.instance.message;
+
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+import org.junit.jupiter.api.Test;
+
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verifyNoInteractions;
+
+class MessageServiceTest {
+
+    @Test
+    void rejectsReversedTopicQueryWindowBeforeCallingProvider() {
+        MessageProvider provider = mock(MessageProvider.class);
+        MessageService service = new MessageService(provider);
+
+        assertThatThrownBy(() -> service.queryMessages("TopicA", null, null, 
null, 200L, 100L))
+                .isInstanceOf(BusinessException.class)
+                .hasMessage("startTime must not be after endTime");
+
+        verifyNoInteractions(provider);
+    }
+
+    @Test
+    void rejectsTopicQueryWindowLongerThanSevenDaysBeforeCallingProvider() {
+        MessageProvider provider = mock(MessageProvider.class);
+        MessageService service = new MessageService(provider);
+
+        assertThatThrownBy(() -> service.queryMessages("TopicA", null, null, 
null, 0L,
+                8L * 24 * 60 * 60 * 1000))
+                .isInstanceOf(BusinessException.class)
+                .hasMessage("topic query time range must not exceed 7 days");
+
+        verifyNoInteractions(provider);
+    }
+}
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
index c94a8505..40520901 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProviderTest.java
@@ -17,6 +17,7 @@
 package org.apache.rocketmq.studio.rocketmq;
 
 import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
+import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageQueue;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
 import org.apache.rocketmq.studio.queryhistory.QueryHistoryService;
@@ -29,6 +30,7 @@ import org.mockito.MockedConstruction;
 import org.mockito.junit.jupiter.MockitoExtension;
 import org.springframework.beans.factory.ObjectProvider;
 
+import java.nio.charset.StandardCharsets;
 import java.util.List;
 
 import static org.assertj.core.api.Assertions.assertThat;
@@ -37,6 +39,7 @@ import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.ArgumentMatchers.anyLong;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.lenient;
 import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
@@ -58,7 +61,7 @@ class RocketMQMessageProviderTest {
 
     @BeforeEach
     void setUp() {
-        when(adminExtProvider.getIfAvailable()).thenReturn(adminExt);
+        lenient().when(adminExtProvider.getIfAvailable()).thenReturn(adminExt);
         provider = new RocketMQMessageProvider(adminExtProvider, 
queryHistoryService, new RocketMQProperties());
     }
 
@@ -82,4 +85,39 @@ class RocketMQMessageProviderTest {
         }
         verify(queryHistoryService).recordMessageQuery("TOPIC", "TopicA", 
null, null, null, 100L, 200L, 0);
     }
+
+    @Test
+    void toRecordVOBoundsMessageBodyAndProperties() {
+        MessageExt message = new MessageExt();
+        message.setMsgId("msg-1");
+        message.setTopic("TopicA");
+        message.setBody("x".repeat(70 * 
1024).getBytes(StandardCharsets.UTF_8));
+        message.putUserProperty("large", "v".repeat(2 * 1024));
+        for (int index = 0; index < 70; index++) {
+            message.putUserProperty("property-" + index, "value");
+        }
+
+        MessageRecordVO record = provider.toRecordVO(message);
+
+        assertThat(record.isBodyTruncated()).isTrue();
+        assertThat(record.getBody()).hasSize(64 * 1024);
+        assertThat(record.getBodyEncoding()).isEqualTo("UTF-8");
+        assertThat(record.isPropertiesTruncated()).isTrue();
+        assertThat(record.getProperties()).hasSize(64);
+        assertThat(record.getProperties().get("large")).endsWith("...");
+    }
+
+    @Test
+    void toRecordVOBase64EncodesBinaryPayloads() {
+        MessageExt message = new MessageExt();
+        message.setMsgId("msg-binary");
+        message.setTopic("TopicA");
+        message.setBody(new byte[] {(byte) 0xC3, (byte) 0x28});
+
+        MessageRecordVO record = provider.toRecordVO(message);
+
+        assertThat(record.getBodyEncoding()).isEqualTo("BASE64");
+        assertThat(record.getBody()).isEqualTo("wyg=");
+        assertThat(record.isBodyTruncated()).isFalse();
+    }
 }

Reply via email to