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 c16dbcaa fix: harden infrastructure operation boundaries (#945)
c16dbcaa is described below
commit c16dbcaaf29ea7c66307cbe6d23bcb016d0df7c1
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 4 03:06:19 2026 -0700
fix: harden infrastructure operation boundaries (#945)
* [ISSUE #848] Validate null instance requests
* [ISSUE #850] Validate null nameserver requests
* [ISSUE #852] Validate null K8s certificate requests
* [ISSUE #860] Validate null cluster config requests
* [ISSUE #784] Return unavailable for unimplemented cluster operations
---
.../studio/cluster/broker/ClusterController.java | 10 ++-
.../studio/cluster/broker/ClusterService.java | 40 ++++++------
.../studio/cluster/k8s/K8sCertController.java | 19 ++++--
.../studio/cluster/k8s/K8sCertService.java | 10 +++
.../cluster/nameserver/NameServerController.java | 25 ++++++--
.../studio/instance/InstanceController.java | 13 +++-
.../rocketmq/studio/instance/InstanceService.java | 8 +++
.../cluster/broker/ClusterControllerTest.java | 12 ++++
.../studio/cluster/broker/ClusterServiceTest.java | 75 ++++++++++++++++++----
.../studio/cluster/k8s/K8sCertControllerTest.java | 21 ++++++
.../studio/cluster/k8s/K8sCertServiceTest.java | 23 +++++++
.../nameserver/NameServerControllerTest.java | 22 +++++++
.../studio/instance/InstanceControllerTest.java | 24 +++++++
.../studio/instance/InstanceServiceTest.java | 21 ++++++
14 files changed, 281 insertions(+), 42 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
index 3ba3161e..36bd36d1 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterController.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.cluster.broker;
import org.apache.rocketmq.studio.cluster.config.UpdateConfigDTO;
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;
@@ -55,7 +56,8 @@ public class ClusterController {
}
@PostMapping("/config/update")
- public Result<ClusterVO> updateClusterConfig(@Valid @RequestBody
UpdateConfigDTO command) {
+ public Result<ClusterVO> updateClusterConfig(@Valid @RequestBody(required
= false) UpdateConfigDTO command) {
+ requireUpdateConfigCommand(command);
return Result.ok(clusterService.updateClusterConfig(command));
}
@@ -68,4 +70,10 @@ public class ClusterController {
"message", "Broker restart initiated for " + name
));
}
+
+ private void requireUpdateConfigCommand(UpdateConfigDTO command) {
+ if (command == null) {
+ throw new BusinessException(400, "Cluster config update request is
required");
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
index a2f0eec3..d50c8499 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterService.java
@@ -26,7 +26,6 @@ import
org.apache.rocketmq.studio.cluster.nameserver.UpdateNameServerDTO;
import org.apache.rocketmq.studio.cluster.nameserver.UpgradeNameServerDTO;
import org.apache.rocketmq.studio.cluster.proxy.RestartProxyDTO;
-import org.apache.rocketmq.studio.common.domain.enums.ClusterStatus;
import org.apache.rocketmq.studio.common.domain.enums.FlushDiskType;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.apache.rocketmq.studio.rocketmq.RocketMQBrokerConfigService;
@@ -202,56 +201,52 @@ public class ClusterService {
.noneMatch(broker -> brokerName.equals(broker.getName()))) {
throw new BusinessException(404, "Broker not found: " +
brokerName);
}
- log.info("Broker restart initiated for: {} in cluster: {}",
brokerName, clusterId);
- return true;
+ throw unsupportedOperation("Broker restart");
}
public NameServerVO createNameServer(CreateNameServerDTO command) {
+ requireNameServerCommand(command);
log.info("Creating NameServer for cluster: {}",
command.getClusterId());
clusterRepository.findById(command.getClusterId())
.orElseThrow(() -> new BusinessException(404, "Cluster not
found: " + command.getClusterId()));
- NameServerVO ns = NameServerVO.builder()
- .addr(command.getAddr())
- .status(ClusterStatus.healthy)
- .build();
- log.info("NameServer created: {}", command.getAddr());
- return ns;
+ throw unsupportedOperation("NameServer create");
}
public void updateNameServer(UpdateNameServerDTO command) {
+ requireNameServerCommand(command);
log.info("Updating NameServer: {} in cluster: {}", command.getAddr(),
command.getClusterId());
ClusterVO cluster = clusterRepository.findById(command.getClusterId())
.orElseThrow(() -> new BusinessException(404, "Cluster not
found: " + command.getClusterId()));
requireNameServer(cluster, command.getAddr());
- log.info("NameServer updated: {}", command.getAddr());
+ throw unsupportedOperation("NameServer update");
}
public boolean restartNameServer(RestartNameServerDTO command) {
+ requireNameServerCommand(command);
log.info("Restarting NameServer: {} in cluster: {}",
command.getAddr(), command.getClusterId());
ClusterVO cluster = clusterRepository.findById(command.getClusterId())
.orElseThrow(() -> new BusinessException(404, "Cluster not
found: " + command.getClusterId()));
requireNameServer(cluster, command.getAddr());
- log.info("NameServer restart initiated: {}", command.getAddr());
- return true;
+ throw unsupportedOperation("NameServer restart");
}
public boolean upgradeNameServer(UpgradeNameServerDTO command) {
+ requireNameServerCommand(command);
log.info("Upgrading NameServer: {} to version: {} in cluster: {}",
command.getAddr(), command.getTargetVersion(),
command.getClusterId());
ClusterVO cluster = clusterRepository.findById(command.getClusterId())
.orElseThrow(() -> new BusinessException(404, "Cluster not
found: " + command.getClusterId()));
requireNameServer(cluster, command.getAddr());
- log.info("NameServer upgrade initiated: {}", command.getAddr());
- return true;
+ throw unsupportedOperation("NameServer upgrade");
}
public boolean deleteNameServer(DeleteNameServerDTO command) {
+ requireNameServerCommand(command);
log.info("Deleting NameServer: {} from cluster: {}",
command.getAddr(), command.getClusterId());
ClusterVO cluster = clusterRepository.findById(command.getClusterId())
.orElseThrow(() -> new BusinessException(404, "Cluster not
found: " + command.getClusterId()));
requireNameServer(cluster, command.getAddr());
- log.info("NameServer deleted: {}", command.getAddr());
- return true;
+ throw unsupportedOperation("NameServer delete");
}
public boolean restartProxy(RestartProxyDTO command) {
@@ -259,8 +254,13 @@ public class ClusterService {
ClusterVO cluster = clusterRepository.findById(command.getClusterId())
.orElseThrow(() -> new BusinessException(404, "Cluster not
found: " + command.getClusterId()));
requireProxy(cluster, command.getAddr());
- log.info("Proxy restart initiated: {}", command.getAddr());
- return true;
+ throw unsupportedOperation("Proxy restart");
+ }
+
+ private void requireNameServerCommand(Object command) {
+ if (command == null) {
+ throw new BusinessException(400, "NameServer request is required");
+ }
}
private void requireNameServer(ClusterVO cluster, String addr) {
@@ -276,4 +276,8 @@ public class ClusterService {
throw new BusinessException(404, "Proxy not found: " + addr);
}
}
+
+ private BusinessException unsupportedOperation(String operation) {
+ return new BusinessException(501, operation + " is not implemented by
the current cluster provider");
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertController.java
index 10ce3c58..2446e6c9 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertController.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.cluster.k8s;
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;
@@ -40,23 +41,33 @@ public class K8sCertController {
}
@PostMapping("/create")
- public Result<K8sCertVO> createCert(@Valid @RequestBody CreateCertDTO
command) {
+ public Result<K8sCertVO> createCert(@Valid @RequestBody(required = false)
CreateCertDTO command) {
+ requireCommand(command);
return Result.ok(k8sCertService.createCert(command));
}
@PostMapping("/update")
- public Result<K8sCertVO> updateCert(@Valid @RequestBody UpdateCertDTO
command) {
+ public Result<K8sCertVO> updateCert(@Valid @RequestBody(required = false)
UpdateCertDTO command) {
+ requireCommand(command);
return Result.ok(k8sCertService.updateCert(command));
}
@PostMapping("/renew")
- public Result<K8sCertVO> renewCert(@Valid @RequestBody RenewCertDTO
command) {
+ public Result<K8sCertVO> renewCert(@Valid @RequestBody(required = false)
RenewCertDTO command) {
+ requireCommand(command);
return Result.ok(k8sCertService.renewCert(command));
}
@PostMapping("/delete")
- public Result<Void> deleteCert(@Valid @RequestBody DeleteCertDTO command) {
+ public Result<Void> deleteCert(@Valid @RequestBody(required = false)
DeleteCertDTO command) {
+ requireCommand(command);
k8sCertService.deleteCert(command);
return Result.ok();
}
+
+ private void requireCommand(Object command) {
+ if (command == null) {
+ throw new BusinessException(400, "K8s certificate request is
required");
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java
index 88206e44..6ca6debf 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertService.java
@@ -59,6 +59,7 @@ public class K8sCertService {
}
public K8sCertVO createCert(CreateCertDTO command) {
+ requireCommand(command);
log.info("Creating K8s certificate: {}", command.getName());
LocalDateTime now = LocalDateTime.now(clock);
@@ -86,6 +87,7 @@ public class K8sCertService {
}
public K8sCertVO updateCert(UpdateCertDTO command) {
+ requireCommand(command);
log.info("Updating K8s certificate: {}", command.getId());
K8sCertVO existing = k8sCertRepository.findById(command.getId())
.orElseThrow(() -> new BusinessException(404, "Certificate not
found: " + command.getId()));
@@ -118,6 +120,7 @@ public class K8sCertService {
}
public K8sCertVO renewCert(RenewCertDTO command) {
+ requireCommand(command);
log.info("Renewing K8s certificate: {}", command.getId());
K8sCertVO existing = k8sCertRepository.findById(command.getId())
.orElseThrow(() -> new BusinessException(404, "Certificate not
found: " + command.getId()));
@@ -138,6 +141,7 @@ public class K8sCertService {
}
public void deleteCert(DeleteCertDTO command) {
+ requireCommand(command);
log.info("Deleting K8s certificate: {}", command.getId());
k8sCertRepository.findById(command.getId())
.orElseThrow(() -> new BusinessException(404, "Certificate not
found: " + command.getId()));
@@ -145,6 +149,12 @@ public class K8sCertService {
log.info("K8s certificate deleted: {}", command.getId());
}
+ private void requireCommand(Object command) {
+ if (command == null) {
+ throw new BusinessException(400, "K8s certificate request is
required");
+ }
+ }
+
private K8sCertVO refreshExpirationState(K8sCertVO cert, LocalDateTime
now) {
K8sCertVO refreshed = copyOf(cert);
LocalDateTime notAfter = refreshed.getNotAfter();
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerController.java
index 0e1b8cc2..2b6ee2d7 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerController.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.cluster.nameserver;
import org.apache.rocketmq.studio.cluster.broker.ClusterService;
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.PostMapping;
@@ -36,31 +37,45 @@ public class NameServerController {
private final ClusterService clusterService;
@PostMapping("/create")
- public Result<NameServerVO> createNameServer(@Valid @RequestBody
CreateNameServerDTO command) {
+ public Result<NameServerVO> createNameServer(@Valid @RequestBody(required
= false) CreateNameServerDTO command) {
+ requireCommand(command);
return Result.ok(clusterService.createNameServer(command));
}
@PostMapping("/update")
- public Result<Void> updateNameServer(@Valid @RequestBody
UpdateNameServerDTO command) {
+ public Result<Void> updateNameServer(@Valid @RequestBody(required = false)
UpdateNameServerDTO command) {
+ requireCommand(command);
clusterService.updateNameServer(command);
return Result.ok();
}
@PostMapping("/restart")
- public Result<Map<String, Boolean>> restartNameServer(@Valid @RequestBody
RestartNameServerDTO command) {
+ public Result<Map<String, Boolean>> restartNameServer(
+ @Valid @RequestBody(required = false) RestartNameServerDTO
command) {
+ requireCommand(command);
boolean success = clusterService.restartNameServer(command);
return Result.ok(Map.of("success", success));
}
@PostMapping("/upgrade")
- public Result<Map<String, Boolean>> upgradeNameServer(@Valid @RequestBody
UpgradeNameServerDTO command) {
+ public Result<Map<String, Boolean>> upgradeNameServer(
+ @Valid @RequestBody(required = false) UpgradeNameServerDTO
command) {
+ requireCommand(command);
boolean success = clusterService.upgradeNameServer(command);
return Result.ok(Map.of("success", success));
}
@PostMapping("/delete")
- public Result<Map<String, Boolean>> deleteNameServer(@Valid @RequestBody
DeleteNameServerDTO command) {
+ public Result<Map<String, Boolean>> deleteNameServer(
+ @Valid @RequestBody(required = false) DeleteNameServerDTO command)
{
+ requireCommand(command);
boolean success = clusterService.deleteNameServer(command);
return Result.ok(Map.of("success", success));
}
+
+ private void requireCommand(Object command) {
+ if (command == null) {
+ throw new BusinessException(400, "NameServer request is required");
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
index e5db352a..ff9d0e83 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceController.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.instance;
import org.apache.rocketmq.studio.common.domain.Result;
import org.apache.rocketmq.studio.common.domain.enums.InstanceType;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.GetMapping;
@@ -45,12 +46,14 @@ public class InstanceController {
}
@PostMapping("/create")
- public Result<InstanceVO> createInstance(@RequestBody InstanceVO instance)
{
+ public Result<InstanceVO> createInstance(@RequestBody(required = false)
InstanceVO instance) {
+ requireInstance(instance);
return Result.ok(instanceService.createInstance(instance));
}
@PostMapping("/update")
- public Result<InstanceVO> updateInstance(@RequestBody InstanceVO instance)
{
+ public Result<InstanceVO> updateInstance(@RequestBody(required = false)
InstanceVO instance) {
+ requireInstance(instance);
return Result.ok(instanceService.updateInstance(instance));
}
@@ -59,4 +62,10 @@ public class InstanceController {
instanceService.deleteInstance(request.getId());
return Result.ok();
}
+
+ private void requireInstance(InstanceVO instance) {
+ if (instance == null) {
+ throw new BusinessException(400, "Instance request is required");
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
index 0e17a59c..b95b48f2 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/InstanceService.java
@@ -49,6 +49,7 @@ public class InstanceService {
}
public InstanceVO createInstance(InstanceVO instance) {
+ requireInstance(instance);
log.info("Creating instance: {}", instance.getName());
if (instance.getName() == null || instance.getName().isBlank()) {
@@ -65,6 +66,7 @@ public class InstanceService {
}
public InstanceVO updateInstance(InstanceVO instance) {
+ requireInstance(instance);
log.info("Updating instance: {}", instance.getId());
if (instance.getId() == null || instance.getId().isBlank()) {
@@ -111,6 +113,12 @@ public class InstanceService {
instanceRepository.deleteById(id);
}
+ private void requireInstance(InstanceVO instance) {
+ if (instance == null) {
+ throw new BusinessException(400, "Instance request is required");
+ }
+ }
+
private InstanceVO copyOf(InstanceVO instance) {
InstanceVO copy = InstanceVO.builder()
.name(instance.getName())
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
index cb33bbcd..04c438df 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterControllerTest.java
@@ -187,6 +187,18 @@ class ClusterControllerTest {
.andExpect(jsonPath("$.data.config.readQueueNums").value(16));
}
+ @Test
+ void updateConfigShouldRejectNullRequestBody() throws Exception {
+ mockMvc.perform(post("/api/clusters/config/update")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("null"))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("Cluster config update
request is required"));
+
+ verifyNoInteractions(clusterService);
+ }
+
@Test
void updateConfigShouldRejectMissingId() throws Exception {
UpdateConfigDTO command = UpdateConfigDTO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
index 5cb0b741..a8bd097f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterServiceTest.java
@@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.cluster.broker;
import org.apache.rocketmq.studio.cluster.config.ClusterConfigVO;
import org.apache.rocketmq.studio.cluster.config.UpdateConfigDTO;
+import org.apache.rocketmq.studio.cluster.nameserver.CreateNameServerDTO;
import org.apache.rocketmq.studio.cluster.nameserver.DeleteNameServerDTO;
import org.apache.rocketmq.studio.cluster.nameserver.NameServerVO;
import org.apache.rocketmq.studio.cluster.nameserver.RestartNameServerDTO;
@@ -41,15 +42,16 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
+import org.assertj.core.api.ThrowableAssert;
import static org.assertj.core.api.Assertions.assertThat;
-import static org.assertj.core.api.Assertions.assertThatCode;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@@ -320,12 +322,11 @@ class ClusterServiceTest {
}
@Test
- void restartBrokerShouldReturnTrue() {
+ void restartBrokerShouldThrowUnsupportedWhenBrokerExists() {
when(clusterRepository.findById("cluster-1")).thenReturn(Optional.of(sampleCluster));
- boolean result = clusterService.restartBroker("cluster-1", "broker-0");
-
- assertThat(result).isTrue();
+ assertUnsupportedOperation(() ->
clusterService.restartBroker("cluster-1", "broker-0"),
+ "Broker restart is not implemented");
}
@Test
@@ -348,7 +349,19 @@ class ClusterServiceTest {
}
@Test
- void nameServerOperationsShouldAcceptExistingNameServer() {
+ void createNameServerShouldThrowUnsupportedWhenClusterExists() {
+
when(clusterRepository.findById("cluster-1")).thenReturn(Optional.of(sampleCluster));
+ CreateNameServerDTO command = CreateNameServerDTO.builder()
+ .clusterId("cluster-1")
+ .addr("10.0.0.21:9876")
+ .build();
+
+ assertUnsupportedOperation(() ->
clusterService.createNameServer(command),
+ "NameServer create is not implemented");
+ }
+
+ @Test
+ void nameServerOperationsShouldThrowUnsupportedWhenNameServerExists() {
when(clusterRepository.findById("cluster-1")).thenReturn(Optional.of(sampleCluster));
UpdateNameServerDTO update = UpdateNameServerDTO.builder()
.clusterId("cluster-1")
@@ -368,21 +381,52 @@ class ClusterServiceTest {
.addr("10.0.0.20:9876")
.build();
- assertThatCode(() ->
clusterService.updateNameServer(update)).doesNotThrowAnyException();
- assertThat(clusterService.restartNameServer(restart)).isTrue();
- assertThat(clusterService.upgradeNameServer(upgrade)).isTrue();
- assertThat(clusterService.deleteNameServer(delete)).isTrue();
+ assertUnsupportedOperation(() ->
clusterService.updateNameServer(update),
+ "NameServer update is not implemented");
+ assertUnsupportedOperation(() ->
clusterService.restartNameServer(restart),
+ "NameServer restart is not implemented");
+ assertUnsupportedOperation(() ->
clusterService.upgradeNameServer(upgrade),
+ "NameServer upgrade is not implemented");
+ assertUnsupportedOperation(() ->
clusterService.deleteNameServer(delete),
+ "NameServer delete is not implemented");
}
@Test
- void restartProxyShouldReturnTrueWhenProxyExists() {
+ void nameServerOperationsShouldRejectNullCommand() {
+ assertThatThrownBy(() ->
clusterService.createNameServer((CreateNameServerDTO) null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("NameServer request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> clusterService.updateNameServer(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("NameServer request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> clusterService.restartNameServer(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("NameServer request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> clusterService.upgradeNameServer(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("NameServer request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> clusterService.deleteNameServer(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("NameServer request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verifyNoInteractions(clusterRepository);
+ }
+
+ @Test
+ void restartProxyShouldThrowUnsupportedWhenProxyExists() {
when(clusterRepository.findById("cluster-1")).thenReturn(Optional.of(sampleCluster));
RestartProxyDTO command = RestartProxyDTO.builder()
.clusterId("cluster-1")
.addr("10.0.0.10:8081")
.build();
- assertThat(clusterService.restartProxy(command)).isTrue();
+ assertUnsupportedOperation(() -> clusterService.restartProxy(command),
+ "Proxy restart is not implemented");
}
@Test
@@ -455,4 +499,11 @@ class ClusterServiceTest {
.hasMessage("Proxy not found: missing:8081")
.satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(404));
}
+
+ private void assertUnsupportedOperation(ThrowableAssert.ThrowingCallable
callable, String message) {
+ assertThatThrownBy(callable)
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining(message)
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(501));
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertControllerTest.java
index 10abd3ae..25da3f4f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertControllerTest.java
@@ -139,6 +139,27 @@ class K8sCertControllerTest {
.andExpect(jsonPath("$.data.id").value("cert-min"));
}
+ @Test
+ void certWriteEndpointsShouldRejectNullRequestBody() throws Exception {
+ String[] paths = {
+ "/api/k8s-certs/create",
+ "/api/k8s-certs/update",
+ "/api/k8s-certs/renew",
+ "/api/k8s-certs/delete"
+ };
+
+ for (String path : paths) {
+ mockMvc.perform(post(path)
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("null"))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("K8s certificate
request is required"));
+ }
+
+ verifyNoInteractions(k8sCertService);
+ }
+
@Test
void createCertShouldRejectMissingName() throws Exception {
mockMvc.perform(post("/api/k8s-certs/create")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java
index 3bd4d515..d93bc32e 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/K8sCertServiceTest.java
@@ -39,6 +39,7 @@ 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)
@@ -170,6 +171,28 @@ class K8sCertServiceTest {
verify(k8sCertRepository).save(any(K8sCertVO.class));
}
+ @Test
+ void certWriteOperationsShouldRejectNullCommand() {
+ assertThatThrownBy(() -> k8sCertService.createCert(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("K8s certificate request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> k8sCertService.updateCert(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("K8s certificate request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> k8sCertService.renewCert(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("K8s certificate request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+ assertThatThrownBy(() -> k8sCertService.deleteCert(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("K8s certificate request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verifyNoInteractions(k8sCertRepository);
+ }
+
@Test
void createCertShouldSetCorrectValidityPeriod() {
CreateCertDTO command = CreateCertDTO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerControllerTest.java
index eb4de252..6f969c93 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NameServerControllerTest.java
@@ -71,6 +71,28 @@ class NameServerControllerTest {
verify(clusterService).createNameServer(any(CreateNameServerDTO.class));
}
+ @Test
+ void nameServerWriteEndpointsShouldRejectNullRequestBody() throws
Exception {
+ String[] paths = {
+ "/api/nameservers/create",
+ "/api/nameservers/update",
+ "/api/nameservers/restart",
+ "/api/nameservers/upgrade",
+ "/api/nameservers/delete"
+ };
+
+ for (String path : paths) {
+ mockMvc.perform(post(path)
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("null"))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("NameServer request
is required"));
+ }
+
+ verifyNoInteractions(clusterService);
+ }
+
@Test
void updateNameServerShouldRejectBlankAddr() throws Exception {
UpdateNameServerDTO request = UpdateNameServerDTO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
index 37f14fdf..c10b9fae 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceControllerTest.java
@@ -133,6 +133,18 @@ class InstanceControllerTest {
.andExpect(jsonPath("$.data.type").value("DIRECT"));
}
+ @Test
+ void createInstanceShouldRejectNullRequestBody() throws Exception {
+ mockMvc.perform(post("/api/instances/create")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("null"))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("Instance request is
required"));
+
+ verifyNoInteractions(instanceService);
+ }
+
@Test
void updateInstanceShouldReturnUpdatedInstance() throws Exception {
InstanceVO update = InstanceVO.builder()
@@ -160,6 +172,18 @@ class InstanceControllerTest {
.andExpect(jsonPath("$.data.name").value("updated-name"));
}
+ @Test
+ void updateInstanceShouldRejectNullRequestBody() throws Exception {
+ mockMvc.perform(post("/api/instances/update")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("null"))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("Instance request is
required"));
+
+ verifyNoInteractions(instanceService);
+ }
+
@Test
void deleteInstanceShouldReturnSuccess() throws Exception {
doNothing().when(instanceService).deleteInstance("inst-1");
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
index bae0806e..64cbc05d 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/InstanceServiceTest.java
@@ -34,6 +34,7 @@ import static
org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@@ -135,6 +136,16 @@ class InstanceServiceTest {
verify(instanceRepository).findAll();
}
+ @Test
+ void createInstanceShouldThrowWhenRequestIsNull() {
+ assertThatThrownBy(() -> instanceService.createInstance(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Instance request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verifyNoInteractions(instanceRepository);
+ }
+
@Test
void createInstanceShouldSetIdAndTimestamps() {
InstanceVO input = InstanceVO.builder()
@@ -287,6 +298,16 @@ class InstanceServiceTest {
assertThat(stored.getUpdatedAt()).isEqualTo(originalUpdatedAt);
}
+ @Test
+ void updateInstanceShouldThrowWhenRequestIsNull() {
+ assertThatThrownBy(() -> instanceService.updateInstance(null))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Instance request is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verifyNoInteractions(instanceRepository);
+ }
+
@Test
void updateInstanceShouldThrowWhenIdIsNull() {
InstanceVO input = InstanceVO.builder().name("test").build();