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