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 c7329ae1 [ISSUE #1522] Validate message size before producer startup 
(#1523)
c7329ae1 is described below

commit c7329ae1b5bf3a729dc6eff7fad9c79a244d8a21
Author: youngkermit8-coder <[email protected]>
AuthorDate: Tue Aug 11 20:21:29 2026 +0800

    [ISSUE #1522] Validate message size before producer startup (#1523)
    
    Signed-off-by: youngkermit8-coder <[email protected]>
---
 .../provider/apache/RocketMQAdminClientImpl.java   | 22 +++++-----
 .../apache/RocketMQAdminClientImplTest.java        | 50 ++++++++++++++++++++++
 2 files changed, 62 insertions(+), 10 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
index ccb3f18f..d22accd8 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
@@ -319,6 +319,18 @@ public class RocketMQAdminClientImpl implements 
AdminClient {
 
     @Override
     public SendMessageVO sendMessage(SendMessageDTO request) {
+        String topic = request.getTopic();
+        String tag = request.getTag() != null ? request.getTag() : "";
+        String key = request.getKey() != null ? request.getKey() : "";
+        String body = request.getBody() != null ? request.getBody() : "";
+        byte[] bodyBytes = body.getBytes(StandardCharsets.UTF_8);
+        if (bodyBytes.length > MAX_MESSAGE_SIZE) {
+            String message = "Message body size " + bodyBytes.length
+                    + " exceeds the maximum of " + MAX_MESSAGE_SIZE + " bytes";
+            recordAudit("SEND_MESSAGE", topic, message, "FAILED");
+            throw new BusinessException(400, message);
+        }
+
         String namesrvAddr = namesrvAddr(request.getInstanceId());
 
         DefaultMQProducer producer = new 
DefaultMQProducer(nextMessageSenderGroup());
@@ -328,16 +340,6 @@ public class RocketMQAdminClientImpl implements 
AdminClient {
         try {
             producer.start();
 
-            String topic = request.getTopic();
-            String tag = request.getTag() != null ? request.getTag() : "";
-            String key = request.getKey() != null ? request.getKey() : "";
-            String body = request.getBody() != null ? request.getBody() : "";
-            byte[] bodyBytes = body.getBytes(StandardCharsets.UTF_8);
-            if (bodyBytes.length > MAX_MESSAGE_SIZE) {
-                throw new BusinessException(400, "Message body size " + 
bodyBytes.length
-                        + " exceeds the maximum of " + MAX_MESSAGE_SIZE + " 
bytes");
-            }
-
             Message msg = new Message(topic, tag, key, bodyBytes);
 
             // Add custom properties
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
index f80e57a6..b935b742 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
@@ -64,6 +64,7 @@ import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verifyNoInteractions;
 import static org.mockito.Mockito.when;
 
 @ExtendWith(MockitoExtension.class)
@@ -376,6 +377,55 @@ class RocketMQAdminClientImplTest {
         }
     }
 
+    @Test
+    void 
sendMessageShouldRejectOversizedUtf8BodyBeforeResolvingEndpointOrCreatingProducer()
 {
+        SendMessageDTO request = new SendMessageDTO();
+        request.setInstanceId("instance-a");
+        request.setTopic("TopicA");
+        request.setBody("\u754c".repeat((4 * 1024 * 1024 / 3) + 1));
+
+        try (MockedConstruction<DefaultMQProducer> mockedProducers =
+                     mockConstruction(DefaultMQProducer.class)) {
+            assertThatThrownBy(() -> adminClient.sendMessage(request))
+                    .isInstanceOf(BusinessException.class)
+                    .hasMessageContaining("exceeds the maximum")
+                    .satisfies(exception -> assertThat(((BusinessException) 
exception).getCode()).isEqualTo(400));
+
+            assertThat(mockedProducers.constructed()).isEmpty();
+        }
+        verifyNoInteractions(runtimeAdminClientResolver);
+        verify(properties, never()).getNamesrvAddr();
+        verify(auditService).record("SEND_MESSAGE", "TopicA",
+                "Message body size 4194306 exceeds the maximum of 4194304 
bytes", "FAILED");
+    }
+
+    @Test
+    void sendMessageShouldAllowBodyAtMaximumSize() throws Exception {
+        when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
+        try (MockedConstruction<DefaultMQProducer> mockedProducers =
+                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
+                         doNothing().when(producer).start();
+                         SendResult sendResult = new SendResult();
+                         sendResult.setSendStatus(SendStatus.SEND_OK);
+                         sendResult.setMsgId("msg-1");
+                         sendResult.setOffsetMsgId("offset-1");
+                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
+                         doNothing().when(producer).shutdown();
+                     })) {
+            SendMessageDTO request = new SendMessageDTO();
+            request.setTopic("TopicA");
+            request.setBody("x".repeat(4 * 1024 * 1024));
+
+            SendMessageVO result = adminClient.sendMessage(request);
+
+            assertThat(result.getMsgId()).isEqualTo("msg-1");
+            DefaultMQProducer producer = 
mockedProducers.constructed().getFirst();
+            ArgumentCaptor<Message> messageCaptor = 
ArgumentCaptor.forClass(Message.class);
+            verify(producer).send(messageCaptor.capture());
+            assertThat(messageCaptor.getValue().getBody()).hasSize(4 * 1024 * 
1024);
+        }
+    }
+
     @Test
     void sendMessageUsesSelectedInstanceEndpoint() throws Exception {
         
when(runtimeAdminClientResolver.resolveEndpoint("instance-a")).thenReturn("10.0.0.2:9876");

Reply via email to