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 8e2c1c1a fix: validate metadata management requests (#930)
8e2c1c1a is described below

commit 8e2c1c1a4f337bbac0f8e248fa07abb25a88e3b7
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 4 03:04:43 2026 -0700

    fix: validate metadata management requests (#930)
    
    * [ISSUE #842] Validate consumer group delete requests
    
    * [ISSUE #858] Validate null topic requests
    
    * fix(server): validate consumer group create requests
    
    * fix(server): validate topic create requests
---
 .../instance/group/ConsumerGroupController.java    |   9 +-
 .../instance/group/CreateConsumerGroupDTO.java     |  57 ++++++++++++
 .../instance/group/DeleteConsumerGroupDTO.java     |  32 +++++++
 .../studio/instance/topic/CreateTopicDTO.java      |  55 +++++++++++
 .../studio/instance/topic/MetadataService.java     |  16 ++++
 .../studio/instance/topic/TopicController.java     |  39 +++++++-
 .../group/ConsumerGroupControllerTest.java         | 103 +++++++++++++++++++++
 .../studio/instance/topic/MetadataServiceTest.java |  21 +++++
 .../studio/instance/topic/TopicControllerTest.java |  76 +++++++++++++++
 9 files changed, 398 insertions(+), 10 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
index 10644eb3..a9a2715b 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupController.java
@@ -30,7 +30,6 @@ import org.springframework.web.bind.annotation.RequestParam;
 import org.springframework.web.bind.annotation.RestController;
 
 import java.util.List;
-import java.util.Map;
 
 @RestController
 @RequestMapping("/api/groups")
@@ -70,13 +69,13 @@ public class ConsumerGroupController {
     }
 
     @PostMapping("/create")
-    public Result<ConsumerGroupVO> createConsumerGroup(@RequestBody 
ConsumerGroupVO group) {
-        return Result.ok(metadataService.createConsumerGroup(group));
+    public Result<ConsumerGroupVO> createConsumerGroup(@Valid @RequestBody 
CreateConsumerGroupDTO group) {
+        return 
Result.ok(metadataService.createConsumerGroup(group.toConsumerGroupVO()));
     }
 
     @PostMapping("/delete")
-    public Result<Void> deleteConsumerGroup(@RequestBody Map<String, String> 
request) {
-        metadataService.deleteConsumerGroup(request.get("name"));
+    public Result<Void> deleteConsumerGroup(@Valid @RequestBody 
DeleteConsumerGroupDTO request) {
+        metadataService.deleteConsumerGroup(request.getName());
         return Result.ok();
     }
 
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/CreateConsumerGroupDTO.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/CreateConsumerGroupDTO.java
new file mode 100644
index 00000000..8175a889
--- /dev/null
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/CreateConsumerGroupDTO.java
@@ -0,0 +1,57 @@
+/*
+ * 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
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.group;
+
+import jakarta.validation.constraints.NotBlank;
+import jakarta.validation.constraints.PositiveOrZero;
+import lombok.Data;
+import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
+import org.apache.rocketmq.studio.common.domain.enums.SubscriptionMode;
+
+@Data
+public class CreateConsumerGroupDTO {
+    @NotBlank(message = "name is required")
+    private String name;
+    private String namespace;
+    private String clusterId;
+    private SubscriptionMode subscriptionMode;
+    private ConsumeType consumeType;
+    private String subscriptionDataType;
+    private String deliveryOrderType;
+    @PositiveOrZero(message = "retryMaxTimes must be zero or positive")
+    private Integer retryMaxTimes;
+    @PositiveOrZero(message = "delaySeconds must be zero or positive")
+    private Integer delaySeconds;
+
+    public ConsumerGroupVO toConsumerGroupVO() {
+        ConsumerGroupVO group = new ConsumerGroupVO();
+        group.setName(name);
+        group.setNamespace(namespace);
+        group.setClusterId(clusterId);
+        group.setSubscriptionMode(subscriptionMode);
+        group.setConsumeType(consumeType);
+        group.setSubscriptionDataType(subscriptionDataType);
+        group.setDeliveryOrderType(deliveryOrderType);
+        if (retryMaxTimes != null) {
+            group.setRetryMaxTimes(retryMaxTimes);
+        }
+        if (delaySeconds != null) {
+            group.setDelaySeconds(delaySeconds);
+        }
+        return group;
+    }
+}
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/group/DeleteConsumerGroupDTO.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/DeleteConsumerGroupDTO.java
new file mode 100644
index 00000000..69cc2ec2
--- /dev/null
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/group/DeleteConsumerGroupDTO.java
@@ -0,0 +1,32 @@
+/*
+ * 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
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.group;
+
+import jakarta.validation.constraints.NotBlank;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class DeleteConsumerGroupDTO {
+    @NotBlank(message = "name is required")
+    private String name;
+}
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/CreateTopicDTO.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/CreateTopicDTO.java
new file mode 100644
index 00000000..2f0146d5
--- /dev/null
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/CreateTopicDTO.java
@@ -0,0 +1,55 @@
+/*
+ * 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
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.instance.topic;
+
+import jakarta.validation.constraints.NotBlank;
+import jakarta.validation.constraints.PositiveOrZero;
+import lombok.Data;
+import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
+import org.apache.rocketmq.studio.common.domain.enums.TopicType;
+
+@Data
+public class CreateTopicDTO {
+    @NotBlank(message = "name is required")
+    private String name;
+    private String namespace;
+    private String clusterId;
+    private TopicType type;
+    @PositiveOrZero(message = "writeQueues must be zero or positive")
+    private Integer writeQueues;
+    @PositiveOrZero(message = "readQueues must be zero or positive")
+    private Integer readQueues;
+    private TopicPerm perm;
+    private String remark;
+
+    public TopicVO toTopicVO() {
+        TopicVO topic = new TopicVO();
+        topic.setName(name);
+        topic.setNamespace(namespace);
+        topic.setClusterId(clusterId);
+        topic.setType(type);
+        if (writeQueues != null) {
+            topic.setWriteQueues(writeQueues);
+        }
+        if (readQueues != null) {
+            topic.setReadQueues(readQueues);
+        }
+        topic.setPerm(perm);
+        topic.setRemark(remark);
+        return topic;
+    }
+}
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
index f2c4e541..4470d7e6 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
@@ -16,6 +16,7 @@
  */
 package org.apache.rocketmq.studio.instance.topic;
 
+import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
 import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
 import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
@@ -46,11 +47,13 @@ public class MetadataService {
 
 
     public TopicVO createTopic(TopicVO topic) {
+        requireTopic(topic);
         return adminClient.createTopic(topic);
     }
 
 
     public TopicVO updateTopic(TopicVO topic) {
+        requireTopic(topic);
         return adminClient.updateTopic(topic);
     }
 
@@ -71,6 +74,7 @@ public class MetadataService {
 
 
     public SendMessageVO sendMessage(SendMessageDTO request) {
+        requireSendMessageRequest(request);
         return adminClient.sendMessage(request);
     }
 
@@ -121,4 +125,16 @@ public class MetadataService {
     private String normalizeFilter(String value) {
         return value == null || value.isBlank() ? null : value.trim();
     }
+
+    private void requireTopic(TopicVO topic) {
+        if (topic == null) {
+            throw new BusinessException(400, "Topic request is required");
+        }
+    }
+
+    private void requireSendMessageRequest(SendMessageDTO request) {
+        if (request == null) {
+            throw new BusinessException(400, "Topic send message request is 
required");
+        }
+    }
 }
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
index 07c059d2..706e095b 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/TopicController.java
@@ -17,6 +17,7 @@
 package org.apache.rocketmq.studio.instance.topic;
 
 import org.apache.rocketmq.studio.common.domain.Result;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
 import jakarta.validation.Valid;
 import lombok.RequiredArgsConstructor;
 import org.springframework.web.bind.annotation.GetMapping;
@@ -45,17 +46,20 @@ public class TopicController {
     }
 
     @PostMapping("/create")
-    public Result<TopicVO> createTopic(@RequestBody TopicVO topic) {
-        return Result.ok(metadataService.createTopic(topic));
+    public Result<TopicVO> createTopic(@Valid @RequestBody(required = false) 
CreateTopicDTO topic) {
+        requireCreateTopicRequest(topic);
+        return Result.ok(metadataService.createTopic(topic.toTopicVO()));
     }
 
     @PostMapping("/update")
-    public Result<TopicVO> updateTopic(@RequestBody TopicVO topic) {
+    public Result<TopicVO> updateTopic(@RequestBody(required = false) TopicVO 
topic) {
+        requireTopicRequest(topic);
         return Result.ok(metadataService.updateTopic(topic));
     }
 
     @PostMapping("/delete")
-    public Result<Void> deleteTopic(@Valid @RequestBody DeleteTopicDTO 
request) {
+    public Result<Void> deleteTopic(@Valid @RequestBody(required = false) 
DeleteTopicDTO request) {
+        requireDeleteTopicRequest(request);
         metadataService.deleteTopic(request.getName());
         return Result.ok();
     }
@@ -71,7 +75,32 @@ public class TopicController {
     }
 
     @PostMapping("/send")
-    public Result<SendMessageVO> sendMessage(@RequestBody SendMessageDTO 
request) {
+    public Result<SendMessageVO> sendMessage(@RequestBody(required = false) 
SendMessageDTO request) {
+        requireSendMessageRequest(request);
         return Result.ok(metadataService.sendMessage(request));
     }
+
+    private void requireTopicRequest(TopicVO topic) {
+        if (topic == null) {
+            throw new BusinessException(400, "Topic request is required");
+        }
+    }
+
+    private void requireCreateTopicRequest(CreateTopicDTO topic) {
+        if (topic == null) {
+            throw new BusinessException(400, "Topic request is required");
+        }
+    }
+
+    private void requireDeleteTopicRequest(DeleteTopicDTO request) {
+        if (request == null) {
+            throw new BusinessException(400, "Topic delete request is 
required");
+        }
+    }
+
+    private void requireSendMessageRequest(SendMessageDTO request) {
+        if (request == null) {
+            throw new BusinessException(400, "Topic send message request is 
required");
+        }
+    }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
index 4dbd2d60..82ad24da 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/group/ConsumerGroupControllerTest.java
@@ -20,6 +20,7 @@ package org.apache.rocketmq.studio.instance.group;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import org.apache.rocketmq.studio.instance.topic.MetadataService;
 import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
 import org.springframework.beans.factory.annotation.Autowired;
 import 
org.springframework.boot.test.autoconfigure.web.servlet.AutoConfigureMockMvc;
 import org.springframework.boot.test.autoconfigure.web.servlet.WebMvcTest;
@@ -31,6 +32,8 @@ import java.time.LocalDateTime;
 import java.util.List;
 import java.util.Map;
 
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.when;
 import static org.mockito.Mockito.verify;
@@ -56,6 +59,70 @@ class ConsumerGroupControllerTest {
     @MockBean
     private ConsumerDiagnosticsService consumerDiagnosticsService;
 
+    @Test
+    void createConsumerGroupShouldPassValidatedRequest() throws Exception {
+        Map<String, Object> body = Map.of(
+                "name", "cg-orders",
+                "clusterId", "cluster-a",
+                "retryMaxTimes", 8,
+                "delaySeconds", 0
+        );
+        ConsumerGroupVO created = new ConsumerGroupVO();
+        created.setName("cg-orders");
+        created.setClusterId("cluster-a");
+        created.setRetryMaxTimes(8);
+
+        
when(metadataService.createConsumerGroup(any(ConsumerGroupVO.class))).thenReturn(created);
+
+        mockMvc.perform(post("/api/groups/create")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content(objectMapper.writeValueAsString(body)))
+                .andExpect(status().isOk())
+                .andExpect(jsonPath("$.code").value(200))
+                .andExpect(jsonPath("$.data.name").value("cg-orders"))
+                .andExpect(jsonPath("$.data.retryMaxTimes").value(8));
+
+        ArgumentCaptor<ConsumerGroupVO> captor = 
ArgumentCaptor.forClass(ConsumerGroupVO.class);
+        verify(metadataService).createConsumerGroup(captor.capture());
+        assertThat(captor.getValue().getName()).isEqualTo("cg-orders");
+        assertThat(captor.getValue().getClusterId()).isEqualTo("cluster-a");
+        assertThat(captor.getValue().getRetryMaxTimes()).isEqualTo(8);
+    }
+
+    @Test
+    void createConsumerGroupShouldRejectMissingName() throws Exception {
+        Map<String, Object> body = Map.of(
+                "clusterId", "cluster-a",
+                "retryMaxTimes", 8
+        );
+
+        mockMvc.perform(post("/api/groups/create")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content(objectMapper.writeValueAsString(body)))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("name is required"));
+
+        verifyNoInteractions(metadataService);
+    }
+
+    @Test
+    void createConsumerGroupShouldRejectNegativeRetryMaxTimes() throws 
Exception {
+        Map<String, Object> body = Map.of(
+                "name", "cg-orders",
+                "retryMaxTimes", -1
+        );
+
+        mockMvc.perform(post("/api/groups/create")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content(objectMapper.writeValueAsString(body)))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("retryMaxTimes must be 
zero or positive"));
+
+        verifyNoInteractions(metadataService);
+    }
+
     @Test
     void getConsumerStackShouldReturnStackTrace() throws Exception {
         ConsumerThreadStackVO thread = ConsumerThreadStackVO.builder()
@@ -105,6 +172,42 @@ class ConsumerGroupControllerTest {
         verify(metadataService).resetOffset(eq("cg-orders"), 
eq(1784246400000L), eq("orders"));
     }
 
+    @Test
+    void deleteConsumerGroupShouldReturnSuccess() throws Exception {
+        mockMvc.perform(post("/api/groups/delete")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        
.content(objectMapper.writeValueAsString(Map.of("name", "cg-orders"))))
+                .andExpect(status().isOk())
+                .andExpect(jsonPath("$.code").value(200))
+                .andExpect(jsonPath("$.message").value("success"));
+
+        verify(metadataService).deleteConsumerGroup("cg-orders");
+    }
+
+    @Test
+    void deleteConsumerGroupShouldRejectMissingName() throws Exception {
+        mockMvc.perform(post("/api/groups/delete")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content("{}"))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("name is required"));
+
+        verifyNoInteractions(metadataService);
+    }
+
+    @Test
+    void deleteConsumerGroupShouldRejectBlankName() throws Exception {
+        mockMvc.perform(post("/api/groups/delete")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        
.content(objectMapper.writeValueAsString(Map.of("name", " "))))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("name is required"));
+
+        verifyNoInteractions(metadataService);
+    }
+
     @Test
     void resetOffsetShouldRejectMissingName() throws Exception {
         Map<String, Object> body = Map.of(
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
index 554335b7..1a059b7c 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/MetadataServiceTest.java
@@ -17,6 +17,7 @@
 
 package org.apache.rocketmq.studio.instance.topic;
 
+import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
@@ -27,8 +28,10 @@ import org.mockito.junit.jupiter.MockitoExtension;
 import java.util.List;
 
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
 import static org.mockito.Mockito.when;
 
 @ExtendWith(MockitoExtension.class)
@@ -90,6 +93,24 @@ class MetadataServiceTest {
         verify(metadataProvider).listTopics(null, null, null);
     }
 
+    @Test
+    void topicWriteOperationsShouldRejectNullRequest() {
+        assertThatThrownBy(() -> metadataService.createTopic(null))
+                .isInstanceOf(BusinessException.class)
+                .hasMessage("Topic request is required")
+                .satisfies(ex -> assertThat(((BusinessException) 
ex).getCode()).isEqualTo(400));
+        assertThatThrownBy(() -> metadataService.updateTopic(null))
+                .isInstanceOf(BusinessException.class)
+                .hasMessage("Topic request is required")
+                .satisfies(ex -> assertThat(((BusinessException) 
ex).getCode()).isEqualTo(400));
+        assertThatThrownBy(() -> metadataService.sendMessage(null))
+                .isInstanceOf(BusinessException.class)
+                .hasMessage("Topic send message request is required")
+                .satisfies(ex -> assertThat(((BusinessException) 
ex).getCode()).isEqualTo(400));
+
+        verifyNoInteractions(adminClient);
+    }
+
     @Test
     void createTopicShouldDelegateToAdminClient() {
         TopicVO input = new TopicVO();
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
index 27c7bc16..d39af4cc 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/instance/topic/TopicControllerTest.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.instance.topic;
 
 import com.fasterxml.jackson.databind.ObjectMapper;
 import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
 import org.springframework.beans.factory.annotation.Autowired;
 import 
org.springframework.boot.test.autoconfigure.web.servlet.AutoConfigureMockMvc;
 import org.springframework.boot.test.autoconfigure.web.servlet.WebMvcTest;
@@ -29,6 +30,7 @@ import org.springframework.test.web.servlet.MockMvc;
 import java.util.List;
 import java.util.Map;
 
+import static org.assertj.core.api.Assertions.assertThat;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.ArgumentMatchers.isNull;
@@ -106,6 +108,80 @@ class TopicControllerTest {
                 .andExpect(jsonPath("$.code").value(200))
                 .andExpect(jsonPath("$.data.name").value("new-topic"))
                 .andExpect(jsonPath("$.data.writeQueues").value(16));
+
+        ArgumentCaptor<TopicVO> captor = 
ArgumentCaptor.forClass(TopicVO.class);
+        verify(metadataService).createTopic(captor.capture());
+        assertThat(captor.getValue().getName()).isEqualTo("new-topic");
+        assertThat(captor.getValue().getWriteQueues()).isEqualTo(16);
+        assertThat(captor.getValue().getReadQueues()).isEqualTo(16);
+    }
+
+    @Test
+    void createTopicShouldRejectMissingName() throws Exception {
+        mockMvc.perform(post("/api/topics/create")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content("""
+                                {
+                                  "writeQueues": 8,
+                                  "readQueues": 8
+                                }
+                                """))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("name is required"));
+
+        verifyNoInteractions(metadataService);
+    }
+
+    @Test
+    void createTopicShouldRejectNegativeQueueCount() throws Exception {
+        mockMvc.perform(post("/api/topics/create")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content("""
+                                {
+                                  "name": "new-topic",
+                                  "writeQueues": -1,
+                                  "readQueues": 8
+                                }
+                                """))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("writeQueues must be 
zero or positive"));
+
+        verifyNoInteractions(metadataService);
+    }
+
+    @Test
+    void topicWriteEndpointsShouldRejectNullRequestBody() throws Exception {
+        mockMvc.perform(post("/api/topics/create")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content("null"))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("Topic request is 
required"));
+
+        mockMvc.perform(post("/api/topics/update")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content("null"))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("Topic request is 
required"));
+
+        mockMvc.perform(post("/api/topics/delete")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content("null"))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("Topic delete request 
is required"));
+
+        mockMvc.perform(post("/api/topics/send")
+                        .contentType(MediaType.APPLICATION_JSON)
+                        .content("null"))
+                .andExpect(status().isBadRequest())
+                .andExpect(jsonPath("$.code").value(400))
+                .andExpect(jsonPath("$.message").value("Topic send message 
request is required"));
+
+        verifyNoInteractions(metadataService);
     }
 
     @Test

Reply via email to