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 571da93a fix(server): harden validation, normalization and null-safety
(#1114)
571da93a is described below
commit 571da93a1c2228a11bbea3e4ca80db01e92eb0b4
Author: 0 <[email protected]>
AuthorDate: Thu Aug 13 16:11:50 2026 +0800
fix(server): harden validation, normalization and null-safety (#1114)
* fix(web): migrate deprecated Ant Design props across Studio pages
Consolidates 17 per-page migration PRs by 123123213weqw: message (#1835),
certificate (#1836), client (#1837), cluster (#1838), ACL (#1839),
consumer (#1840), DLQ (#1841), instance (#1842), topic (#1843),
alert rule (#1844), settings modal (#1845), alert asset modal (#1846),
group management (#1847), LiteTopic (#1848), producer (#1849),
proxy (#1850), SSL settings (#1851). Behavior unchanged.
* fix(server): harden validation, normalization and null-safety
Consolidates 40 backend robustness PRs by 123123213weqw (#1114,
#1733-#1735, #1737-#1739, #1742, #1744-#1745, #1854-#1882, #1884):
block new admin connections during shutdown, Locale.ROOT normalization
for LLM engines and agent providers, reject null or malformed persisted
JSON, clamp cloud catalog counts and retry values, skip null catalog,
trace and connection entries, report lost concurrent updates for alert
rules, ACL rules, credentials, instances and certificates, normalize
stored enum values, validate alert identifiers and stored auth modes.
---
.../studio/cluster/broker/MqAdminExtFactory.java | 9 ++-
.../client/ProducerConnectionSummaryVO.java | 4 +-
.../cluster/k8s/MybatisPlusK8sCertRepository.java | 24 ++++--
.../AbstractPrometheusCompatibleMetricsSource.java | 20 ++++-
.../studio/cluster/metrics/MetricsService.java | 3 +-
.../metrics/grafana/GrafanaDashboardService.java | 3 +
.../rocketmq/studio/common/domain/PageResult.java | 2 +-
.../studio/common/util/CredentialUtils.java | 15 +++-
.../rocketmq/studio/common/util/UrlHostGuard.java | 25 +++++--
.../rocketmq/studio/instance/InstanceService.java | 5 ++
.../instance/MybatisPlusInstanceRepository.java | 5 +-
.../instance/acl/MybatisPlusAclRepository.java | 4 +-
.../studio/ops/ai/AgentProviderRegistry.java | 3 +-
.../apache/rocketmq/studio/ops/ai/LlmConfigVO.java | 4 +-
.../studio/ops/alert/AlertRuleAssetService.java | 14 +++-
.../rocketmq/studio/ops/alert/AlertService.java | 26 +++++--
.../ops/alert/MybatisPlusAlertRepository.java | 5 +-
.../rocketmq/studio/ops/audit/AuditCleanupDTO.java | 2 +
.../rocketmq/studio/ops/audit/AuditService.java | 3 +
.../persistence/MybatisPlusSettingsRepository.java | 14 +++-
.../provider/alibaba/AliyunCatalogService.java | 5 +-
.../provider/apache/RocketMQMessageProvider.java | 1 +
.../credential/CloudCredentialService.java | 4 +-
.../MybatisPlusCloudCredentialRepository.java | 5 +-
.../cluster/broker/MqAdminExtFactoryTest.java | 13 ++++
.../client/ProducerConnectionSummaryVOTest.java | 10 +++
.../k8s/MybatisPlusK8sCertRepositoryTest.java | 46 ++++++++++++
.../studio/cluster/metrics/MetricsServiceTest.java | 22 ++++++
.../metrics/PrometheusMetricsSourceTest.java | 26 +++++++
.../grafana/GrafanaDashboardServiceTest.java | 11 +++
.../studio/common/domain/PageResultTest.java} | 25 +++----
.../studio/common/util/CredentialUtilsTest.java} | 27 +++----
.../common/util/UrlHostGuardMulticastTest.java} | 24 +++---
.../studio/common/util/UrlHostGuardTest.java | 43 +++++++++++
.../studio/instance/InstanceServiceTest.java | 47 ++++++++++++
.../MybatisPlusInstanceRepositoryTest.java | 14 ++++
.../instance/acl/MybatisPlusAclRepositoryTest.java | 17 +++++
.../studio/ops/ai/AgentProviderRegistryTest.java} | 40 +++++-----
.../rocketmq/studio/ops/ai/LlmConfigVOTest.java} | 33 +++++----
.../ops/alert/AlertRuleAssetServiceTest.java | 17 +++++
.../studio/ops/alert/AlertServiceTest.java | 85 ++++++++++++++++++++++
.../ops/alert/MybatisPlusAlertRepositoryTest.java | 22 ++++++
.../studio/ops/audit/AuditControllerTest.java | 12 +++
.../studio/ops/audit/AuditServiceTest.java | 7 ++
.../MybatisPlusSettingsRepositoryTest.java | 44 +++++++++++
.../provider/alibaba/AliyunCatalogServiceTest.java | 30 ++++++++
.../alibaba/AliyunConvertersMessageKeysTest.java} | 30 ++++----
.../alibaba/AliyunConvertersNullEndpointTest.java | 46 ++++++++++++
.../provider/alibaba/AliyunConvertersTest.java} | 36 ++++-----
.../alibaba/AliyunConvertersTraceElementsTest.java | 47 ++++++++++++
.../apache/RocketMQMessageProviderTest.java | 29 ++++++++
.../credential/CloudCredentialServiceTest.java | 17 +++++
.../MybatisPlusCloudCredentialRepositoryTest.java | 15 ++++
53 files changed, 885 insertions(+), 155 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
index 2bffa05e..0377f867 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactory.java
@@ -91,7 +91,14 @@ public class MqAdminExtFactory {
AdminClientCacheKey cacheKey = new
AdminClientCacheKey(normalizedNamesrvAddr,
normalizeAuthenticationIdentity(authenticationIdentity));
DefaultMQAdminExt admin = cache.computeIfAbsent(cacheKey,
- key -> createAndStart(key.namesrvAddr(), rpcHook));
+ key -> {
+ // Re-check under the cache lock so a request that passed
the initial closed check
+ // cannot create a fresh connection while the factory is
shutting down.
+ if (closed) {
+ throw new BusinessException(503, "Admin factory is
shutting down");
+ }
+ return createAndStart(key.namesrvAddr(), rpcHook);
+ });
try {
return action.apply(admin);
} catch (BusinessException ex) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
index 7294a2e1..75362d5e 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVO.java
@@ -49,7 +49,9 @@ public class ProducerConnectionSummaryVO {
private String readiness = READY;
public static ProducerConnectionSummaryVO from(List<ProducerConnectionVO>
connections) {
- List<ProducerConnectionVO> safeConnections = connections == null ?
List.of() : connections;
+ List<ProducerConnectionVO> safeConnections = connections == null
+ ? List.of()
+ : connections.stream().filter(Objects::nonNull).toList();
ProducerConnectionSummaryVO summary = new
ProducerConnectionSummaryVO();
summary.totalConnections = safeConnections.size();
summary.uniqueClientCount = countDistinct(safeConnections,
ProducerConnectionVO::getClientId);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepository.java
index 155b576f..11bb8b33 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepository.java
@@ -32,6 +32,7 @@ import lombok.RequiredArgsConstructor;
import java.time.LocalDateTime;
import java.util.List;
+import java.util.Locale;
import java.util.Optional;
import java.util.stream.Collectors;
@@ -63,7 +64,10 @@ public class MybatisPlusK8sCertRepository implements
K8sCertRepository {
public K8sCertVO save(K8sCertVO cert) {
RmqK8sCertificate entity = toEntity(cert);
if (entity.getId() != null && certMapper.selectById(entity.getId()) !=
null) {
- certMapper.updateById(entity);
+ if (certMapper.updateById(entity) == 0) {
+ throw new BusinessException(409,
+ "Certificate update was not applied: " +
entity.getId());
+ }
} else {
certMapper.insert(entity);
cert.setId(entity.getId());
@@ -116,11 +120,13 @@ public class MybatisPlusK8sCertRepository implements
K8sCertRepository {
if (!StringUtils.hasText(value)) {
return null;
}
- try {
- return CertType.valueOf(value);
- } catch (IllegalArgumentException exception) {
- throw invalidPersistedValue(certificateId, "type", value);
+ String normalized = value.trim();
+ for (CertType type : CertType.values()) {
+ if (type.name().equalsIgnoreCase(normalized)) {
+ return type;
+ }
}
+ throw invalidPersistedValue(certificateId, "type", value);
}
private CertStatus parseCertStatus(String certificateId, String value) {
@@ -128,7 +134,7 @@ public class MybatisPlusK8sCertRepository implements
K8sCertRepository {
return null;
}
try {
- return CertStatus.valueOf(value);
+ return CertStatus.valueOf(value.trim().toLowerCase(Locale.ROOT));
} catch (IllegalArgumentException exception) {
throw invalidPersistedValue(certificateId, "status", value);
}
@@ -139,8 +145,12 @@ public class MybatisPlusK8sCertRepository implements
K8sCertRepository {
return List.of();
}
try {
- return objectMapper.readValue(json, new
TypeReference<List<String>>() {
+ List<String> san = objectMapper.readValue(json, new
TypeReference<List<String>>() {
});
+ if (san == null) {
+ throw invalidPersistedValue(certificateId, "SAN JSON", json);
+ }
+ return san;
} catch (JsonProcessingException exception) {
throw invalidPersistedValue(certificateId, "SAN JSON", json);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/AbstractPrometheusCompatibleMetricsSource.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/AbstractPrometheusCompatibleMetricsSource.java
index 7d13ed9e..4bf11e7d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/AbstractPrometheusCompatibleMetricsSource.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/AbstractPrometheusCompatibleMetricsSource.java
@@ -268,6 +268,10 @@ public abstract class
AbstractPrometheusCompatibleMetricsSource implements Metri
Iterator<Map.Entry<String, JsonNode>> fields = metric.fields();
fields.forEachRemaining(entry -> {
String key = entry.getKey();
+ if (!entry.getValue().isTextual()) {
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ backendLabel() + " returned a malformed time-series
label");
+ }
String value = entry.getValue().asText();
if (key.length() > MAX_LABEL_KEY_LENGTH || value.length() >
MAX_LABEL_VALUE_LENGTH) {
throw new
PrometheusException(HttpStatus.PAYLOAD_TOO_LARGE.value(),
@@ -291,12 +295,17 @@ public abstract class
AbstractPrometheusCompatibleMetricsSource implements Metri
private MetricDataVO.MetricSampleVO parseSample(JsonNode sampleNode) {
if (sampleNode == null || !sampleNode.isArray() || sampleNode.size()
!= 2
- || !sampleNode.get(0).isNumber()) {
+ || !sampleNode.get(0).isNumber() ||
!sampleNode.get(1).isTextual()) {
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ backendLabel() + " returned a malformed sample");
+ }
+ double timestamp = sampleNode.get(0).asDouble();
+ if (!Double.isFinite(timestamp)) {
throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
backendLabel() + " returned a malformed sample");
}
return MetricDataVO.MetricSampleVO.builder()
- .timestamp(sampleNode.get(0).asDouble())
+ .timestamp(timestamp)
.value(sampleNode.get(1).asText())
.build();
}
@@ -307,8 +316,13 @@ public abstract class
AbstractPrometheusCompatibleMetricsSource implements Metri
throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
backendLabel() + " returned a malformed histogram sample");
}
+ double timestamp = sampleNode.get(0).asDouble();
+ if (!Double.isFinite(timestamp)) {
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ backendLabel() + " returned a malformed histogram sample");
+ }
return MetricDataVO.MetricHistogramSampleVO.builder()
- .timestamp(sampleNode.get(0).asDouble())
+ .timestamp(timestamp)
.histogram(sampleNode.get(1))
.build();
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
index 04ceec22..2ceb35f0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
@@ -121,9 +121,10 @@ public class MetricsService {
return "none";
}
return switch (auth.trim().toLowerCase(Locale.ROOT)) {
+ case "none" -> "none";
case "basic auth", "basic" -> "basic";
case "bearer token", "bearer" -> "bearer";
- default -> "none";
+ default -> throw badRequest("Unsupported data source
authentication mode: " + auth.trim());
};
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardService.java
index e0c5dfdd..c28bb588 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardService.java
@@ -100,6 +100,9 @@ public class GrafanaDashboardService {
try (InputStream in = resource.getInputStream()) {
@SuppressWarnings("unchecked")
Map<String, Object> model = objectMapper.readValue(in, Map.class);
+ if (model == null) {
+ throw new BusinessException(500, "Failed to read Grafana
dashboard: " + uid);
+ }
return model;
} catch (IOException e) {
throw new BusinessException(500, "Failed to read Grafana
dashboard: " + uid);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
b/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
index ae7507cd..34fa1f57 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
@@ -31,7 +31,7 @@ public class PageResult<T> {
public static <T> PageResult<T> of(List<T> items, long total, int page,
int size) {
PageResult<T> result = new PageResult<>();
- result.items = items;
+ result.items = items == null ? Collections.emptyList() : items;
result.total = total;
result.page = page;
result.size = size;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/util/CredentialUtils.java
b/server/src/main/java/org/apache/rocketmq/studio/common/util/CredentialUtils.java
index f665f800..718a5757 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/util/CredentialUtils.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/common/util/CredentialUtils.java
@@ -16,6 +16,9 @@
*/
package org.apache.rocketmq.studio.common.util;
+import java.nio.ByteBuffer;
+import java.nio.charset.CharacterCodingException;
+import java.nio.charset.CodingErrorAction;
import java.nio.charset.StandardCharsets;
import java.util.Base64;
@@ -44,9 +47,15 @@ public final class CredentialUtils {
return null;
}
try {
- return new String(Base64.getDecoder().decode(stored),
StandardCharsets.UTF_8);
- } catch (IllegalArgumentException ex) {
- // tolerate legacy values that were stored without encoding
+ byte[] decoded = Base64.getDecoder().decode(stored);
+ return StandardCharsets.UTF_8.newDecoder()
+ .onMalformedInput(CodingErrorAction.REPORT)
+ .onUnmappableCharacter(CodingErrorAction.REPORT)
+ .decode(ByteBuffer.wrap(decoded))
+ .toString();
+ } catch (IllegalArgumentException | CharacterCodingException ex) {
+ // Tolerate legacy plaintext, including text that is syntactically
Base64 but does
+ // not decode to valid UTF-8. Returning replacement characters
would corrupt it.
return stored;
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/util/UrlHostGuard.java
b/server/src/main/java/org/apache/rocketmq/studio/common/util/UrlHostGuard.java
index 76d68c48..43c9ef3d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/util/UrlHostGuard.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/common/util/UrlHostGuard.java
@@ -99,17 +99,26 @@ public final class UrlHostGuard {
return allowLoopback;
}
try {
- InetAddress address = InetAddress.getByName(normalized);
- if (address.isAnyLocalAddress() || address.isLinkLocalAddress()) {
- return false;
- }
- if (address.isLoopbackAddress()) {
- return allowLoopback;
- }
- return true;
+ return areAllowed(InetAddress.getAllByName(normalized),
allowLoopback);
} catch (UnknownHostException exception) {
// Fail closed: an unresolvable host must not be handed to the
connection layer.
return false;
}
}
+
+ static boolean areAllowed(InetAddress[] addresses, boolean allowLoopback) {
+ if (addresses == null || addresses.length == 0) {
+ return false;
+ }
+ for (InetAddress address : addresses) {
+ if (address == null || address.isAnyLocalAddress() ||
address.isLinkLocalAddress()
+ || address.isMulticastAddress()) {
+ return false;
+ }
+ if (address.isLoopbackAddress() && !allowLoopback) {
+ return false;
+ }
+ }
+ return true;
+ }
}
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 163f443b..ffdff159 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
@@ -148,6 +148,10 @@ public class InstanceService {
}
CloudInstanceDetailVO detail = providerRegistry.catalogFor(vendor)
.getCloudInstance(instance.getCredentialId(),
instance.getRegionId(), instance.getCloudInstanceId());
+ if (detail == null) {
+ throw new BusinessException(502,
+ "Cloud instance details unavailable: " +
instance.getCloudInstanceId());
+ }
if (!StringUtils.hasText(instance.getName())) {
instance.setName(detail.getInstanceName() != null &&
!detail.getInstanceName().isBlank()
? detail.getInstanceName() : detail.getInstanceId());
@@ -162,6 +166,7 @@ public class InstanceService {
throw new BusinessException(502, "Cloud instance has no endpoint:
" + detail.getInstanceId());
}
return detail.getEndpoints().stream()
+ .filter(Objects::nonNull)
.filter(endpoint -> endpoint.getEndpointUrl() != null &&
!endpoint.getEndpointUrl().isBlank())
.sorted((a, b) ->
Integer.compare(endpointPriority(a.getEndpointType()),
endpointPriority(b.getEndpointType())))
.map(CloudInstanceDetailVO.CloudEndpoint::getEndpointUrl)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
index 961f89fa..fed034c9 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepository.java
@@ -112,7 +112,10 @@ public class MybatisPlusInstanceRepository implements
InstanceRepository {
public InstanceVO save(InstanceVO instance) {
RmqInstance entity = toEntity(instance);
if (instanceMapper.selectById(entity.getId()) != null) {
- instanceMapper.updateById(entity);
+ if (instanceMapper.updateById(entity) == 0) {
+ throw new BusinessException(409,
+ "Instance update was not applied: " + entity.getId());
+ }
} else {
instanceMapper.insert(entity);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
index f8243b72..3101ed8c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepository.java
@@ -77,7 +77,9 @@ public class MybatisPlusAclRepository implements
AclRepository {
}
RmqAclRule entity = toRuleEntity(rule);
entity.setCreatedAt(existing.getCreatedAt());
- ruleMapper.updateById(entity);
+ if (ruleMapper.updateById(entity) == 0) {
+ return Optional.empty();
+ }
rule.setCreatedAt(existing.getCreatedAt());
return Optional.of(rule);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/AgentProviderRegistry.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/AgentProviderRegistry.java
index 091f6f24..7b86cf9f 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/AgentProviderRegistry.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/AgentProviderRegistry.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.ops.ai;
import org.springframework.stereotype.Component;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.function.Function;
import java.util.stream.Collectors;
@@ -37,7 +38,7 @@ public class AgentProviderRegistry {
}
public AgentProvider forEngine(String engine) {
- AgentProvider provider = providers.get(engine == null ? "" :
engine.trim().toLowerCase());
+ AgentProvider provider = providers.get(engine == null ? "" :
engine.trim().toLowerCase(Locale.ROOT));
if (provider == null) {
throw new LlmGatewayException(400, "llm.config.unsupported_engine",
"Agent engine is not supported: " + engine,
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigVO.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigVO.java
index 23c222c8..f5f1587b 100644
--- a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigVO.java
+++ b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/LlmConfigVO.java
@@ -25,6 +25,8 @@ import lombok.NoArgsConstructor;
import lombok.ToString;
import org.springframework.util.StringUtils;
+import java.util.Locale;
+
@Data
@Builder
@NoArgsConstructor
@@ -67,6 +69,6 @@ public class LlmConfigVO {
}
public String normalizeEngine() {
- return StringUtils.hasText(engine) ? engine.trim().toLowerCase() :
ENGINE_HTTP;
+ return StringUtils.hasText(engine) ?
engine.trim().toLowerCase(Locale.ROOT) : ENGINE_HTTP;
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetService.java
index d56eb3b5..85e0a529 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetService.java
@@ -46,7 +46,15 @@ public class AlertRuleAssetService {
private static final String LOCATION_PATTERN = "classpath*:alerts/*.yaml";
private final ObjectMapper yamlMapper = new ObjectMapper(new
YAMLFactory());
- private final ResourcePatternResolver resourceResolver = new
PathMatchingResourcePatternResolver();
+ private final ResourcePatternResolver resourceResolver;
+
+ public AlertRuleAssetService() {
+ this(new PathMatchingResourcePatternResolver());
+ }
+
+ AlertRuleAssetService(ResourcePatternResolver resourceResolver) {
+ this.resourceResolver = resourceResolver;
+ }
/**
* Lists metadata for every bundled alert rule asset.
@@ -166,8 +174,8 @@ public class AlertRuleAssetService {
try {
return resourceResolver.getResources(LOCATION_PATTERN);
} catch (IOException e) {
- log.warn("Unable to resolve alert rule assets: {}",
e.getMessage());
- return new Resource[0];
+ log.error("Unable to resolve alert rule assets", e);
+ throw new BusinessException(500, "Failed to resolve bundled alert
rule assets");
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
index 250b73f0..175dd72c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertService.java
@@ -28,6 +28,7 @@ import java.util.List;
import java.util.HashSet;
import java.util.Locale;
import java.util.Map;
+import java.util.Objects;
import java.util.Set;
import java.util.UUID;
import java.util.regex.Pattern;
@@ -120,9 +121,10 @@ public class AlertService {
public AlertRuleVO toggleRule(String id, boolean enabled) {
log.info("Toggling alert rule id={}, enabled={}", id, enabled);
+ validateRuleId(id);
List<AlertRuleVO> rules = alertRepository.findAllRules();
AlertRuleVO rule = rules.stream()
- .filter(r -> r.getId().equals(id))
+ .filter(r -> Objects.equals(r.getId(), id))
.findFirst()
.orElseThrow(() -> new
org.apache.rocketmq.studio.common.exception.BusinessException(404, "Alert rule
not found: " + id));
rule.setEnabled(enabled);
@@ -134,6 +136,7 @@ public class AlertService {
public void deleteRule(String id) {
log.info("Deleting alert rule id={}", id);
+ validateRuleId(id);
if (!alertRepository.deleteRule(id)) {
throw ruleNotFound(id);
}
@@ -149,9 +152,12 @@ public class AlertService {
public SystemAlertVO acknowledgeAlert(String id) {
log.info("Acknowledging system alert id={}", id);
+ if (id == null || id.isBlank()) {
+ throw new BusinessException(400, "System alert ID is required");
+ }
List<SystemAlertVO> alerts = alertRepository.findAlerts(null);
SystemAlertVO alert = alerts.stream()
- .filter(a -> a.getId().equals(id))
+ .filter(a -> Objects.equals(a.getId(), id))
.findFirst()
.orElseThrow(() -> new
org.apache.rocketmq.studio.common.exception.BusinessException(404, "System
alert not found: " + id));
alert.setAcknowledged(true);
@@ -290,16 +296,19 @@ public class AlertService {
if (!hasText(metric)) {
return "broker";
}
- if (metric.contains("replication") || metric.contains("fall_behind")
|| metric.contains("slave")) {
+ String normalizedMetric = metric.toLowerCase(Locale.ROOT);
+ if (normalizedMetric.contains("replication") ||
normalizedMetric.contains("fall_behind")
+ || normalizedMetric.contains("slave")) {
return "broker";
}
- if (metric.contains("consumer") || metric.contains("lag")) {
+ if (normalizedMetric.contains("consumer") ||
normalizedMetric.contains("lag")) {
return "consumer";
}
- if (metric.contains("producer") || metric.contains("client")) {
+ if (normalizedMetric.contains("producer") ||
normalizedMetric.contains("client")) {
return "client";
}
- if (metric.contains("topic") || metric.contains("messages_in") ||
metric.contains("messages_out")) {
+ if (normalizedMetric.contains("topic") ||
normalizedMetric.contains("messages_in")
+ || normalizedMetric.contains("messages_out")) {
return "topic";
}
return "broker";
@@ -308,7 +317,10 @@ public class AlertService {
private String summary(AlertRuleVO rule) {
String description = rule.getDescription();
if (hasText(description) && description.contains(" - ")) {
- return description.substring(0, description.indexOf(" - "));
+ String candidate = description.substring(0, description.indexOf("
- "));
+ if (hasText(candidate)) {
+ return candidate;
+ }
}
return hasText(rule.getName()) ? rule.getName() : "RocketMQ alert";
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
index b67c370e..2c04157c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepository.java
@@ -67,8 +67,7 @@ public class MybatisPlusAlertRepository implements
AlertRepository {
if (ruleMapper.selectById(rule.getId()) == null) {
return false;
}
- ruleMapper.updateById(toRuleEntity(rule));
- return true;
+ return ruleMapper.updateById(toRuleEntity(rule)) > 0;
}
@Override
@@ -172,7 +171,7 @@ public class MybatisPlusAlertRepository implements
AlertRepository {
return AlertLevel.info;
}
try {
- return AlertLevel.valueOf(level);
+ return AlertLevel.valueOf(level.trim().toLowerCase(Locale.ROOT));
} catch (IllegalArgumentException exception) {
return AlertLevel.info;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
index 7c3fd411..1f6f4434 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.studio.ops.audit;
+import jakarta.validation.constraints.Max;
import jakarta.validation.constraints.Positive;
import lombok.AllArgsConstructor;
import lombok.Builder;
@@ -28,5 +29,6 @@ import lombok.NoArgsConstructor;
@AllArgsConstructor
public class AuditCleanupDTO {
@Positive(message = "beforeDays must be greater than 0")
+ @Max(value = 365, message = "beforeDays must not exceed 365")
private Integer beforeDays;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
index 4179ffeb..956206cf 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditService.java
@@ -106,6 +106,9 @@ public class AuditService {
if (beforeDays <= 0) {
throw new BusinessException(400, "beforeDays must be greater than
0");
}
+ if (beforeDays > 365) {
+ throw new BusinessException(400, "beforeDays must not exceed 365");
+ }
log.info("Cleaning up audit logs older than {} days", beforeDays);
LocalDateTime cutoff = LocalDateTime.now().minusDays(beforeDays);
return auditRepository.deleteBefore(cutoff);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepository.java
index 45e63a45..68843f37 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepository.java
@@ -71,7 +71,11 @@ public class MybatisPlusSettingsRepository implements
SettingsRepository {
.build();
}
try {
- return objectMapper.readValue(entity.getJson(),
GeneralSettingsVO.class);
+ GeneralSettingsVO settings =
objectMapper.readValue(entity.getJson(), GeneralSettingsVO.class);
+ if (settings == null) {
+ throw new BusinessException(500, "Persisted general settings
are invalid");
+ }
+ return settings;
} catch (JsonProcessingException e) {
log.error("Failed to deserialize general settings", e);
throw new BusinessException(500, "Persisted general settings are
invalid");
@@ -127,8 +131,7 @@ public class MybatisPlusSettingsRepository implements
SettingsRepository {
}
existing.setJson(toJson(dataSource));
existing.setUpdatedAt(LocalDateTime.now());
- dataSourceMapper.updateById(existing);
- return true;
+ return dataSourceMapper.updateById(existing) > 0;
}
@Override
@@ -148,9 +151,12 @@ public class MybatisPlusSettingsRepository implements
SettingsRepository {
private DataSourceVO toDataSourceVO(RmqDataSource entity) {
try {
DataSourceVO vo = objectMapper.readValue(entity.getJson(),
DataSourceVO.class);
+ if (vo == null) {
+ throw new BusinessException(500, "Persisted data source is
invalid: " + entity.getDsKey());
+ }
vo.setKey(entity.getDsKey());
return vo;
- } catch (JsonProcessingException e) {
+ } catch (JsonProcessingException | IllegalArgumentException e) {
log.error("Failed to deserialize data source: {}",
entity.getDsKey(), e);
throw new BusinessException(500, "Persisted data source is
invalid: " + entity.getDsKey());
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
index 1e6ce14b..b8477ee7 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogService.java
@@ -70,7 +70,7 @@ public class AliyunCatalogService implements
CloudCatalogProvider {
return regions;
}
for (ListRegionsResponseBody.Data item : data) {
- if (Boolean.TRUE.equals(item.getSupportRocketmqV5())) {
+ if (item != null &&
Boolean.TRUE.equals(item.getSupportRocketmqV5())) {
regions.add(AliyunConverters.toRegionVO(item));
}
}
@@ -86,6 +86,9 @@ public class AliyunCatalogService implements
CloudCatalogProvider {
List<ListInstancesResponseBody.List> all =
fetchAllInstances(credentialId, regionId);
List<CloudInstanceOptionVO> options = new ArrayList<>();
for (ListInstancesResponseBody.List item : all) {
+ if (item == null) {
+ continue;
+ }
CloudInstanceOptionVO vo =
AliyunConverters.toInstanceOptionVO(item);
if (matchesSearch(search, vo)) {
options.add(vo);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index 511d66cd..86702559 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -260,6 +260,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
} finally {
consumer.shutdown();
}
+
result.sort(Comparator.comparingLong(MessageRecordVO::getStoreTime).reversed());
return result;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java
index 29b93db6..e8ea14a3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialService.java
@@ -116,7 +116,9 @@ public class CloudCredentialService {
if (instanceRepository.existsByCredentialId(id)) {
throw new BusinessException(400, "Cloud credential is referenced
by existing instances");
}
- credentialRepository.deleteById(id);
+ if (!credentialRepository.deleteById(id)) {
+ throw new BusinessException(404, "Cloud credential not found: " +
id);
+ }
invalidateCloudClients(existing);
recordAudit("DELETE_CLOUD_CREDENTIAL", "CLOUD_CREDENTIAL", id, null,
credentialAuditDetail(existing));
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepository.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepository.java
index fbd806c5..aef890fd 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepository.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepository.java
@@ -68,7 +68,10 @@ public class MybatisPlusCloudCredentialRepository implements
CloudCredentialRepo
public CloudCredentialVO save(CloudCredentialVO credential) {
RmqCloudCredential entity = toEntity(credential);
if (credentialMapper.selectById(entity.getId()) != null) {
- credentialMapper.updateById(entity);
+ if (credentialMapper.updateById(entity) == 0) {
+ throw new BusinessException(409,
+ "Cloud credential update was not applied: " +
entity.getId());
+ }
} else {
credentialMapper.insert(entity);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
index 31965e12..d04e0f7d 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/MqAdminExtFactoryTest.java
@@ -177,4 +177,17 @@ class MqAdminExtFactoryTest {
.isInstanceOf(BusinessException.class)
.hasMessageContaining("shutting down");
}
+
+ @Test
+ void executeShouldNotCreateConnectionAfterShutdown() {
+ DefaultMQAdminExt admin = mock(DefaultMQAdminExt.class);
+ RecordingFactory factory = new RecordingFactory(admin);
+ factory.shutdown();
+
+ assertThatThrownBy(() -> factory.execute("10.0.0.1:9876", null, a ->
"unused"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("shutting down");
+ // No fresh admin connection may be established while the factory is
shut down.
+ assertThat(factory.created.get()).isZero();
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
index 24c5eb68..14065d5c 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/client/ProducerConnectionSummaryVOTest.java
@@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.cluster.client;
import org.junit.jupiter.api.Test;
+import java.util.Arrays;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
@@ -33,6 +34,15 @@ class ProducerConnectionSummaryVOTest {
assertThat(summary.getWarnings()).containsExactly(ProducerConnectionSummaryVO.NO_CONNECTIONS);
}
+ @Test
+ void fromShouldSkipNullConnectionEntries() {
+ ProducerConnectionSummaryVO summary =
ProducerConnectionSummaryVO.from(Arrays.asList(
+ null, connection("producer-a", "10.0.0.1:38888", "Java",
"5.1.0")));
+
+ assertThat(summary.getTotalConnections()).isEqualTo(1);
+
assertThat(summary.getReadiness()).isEqualTo(ProducerConnectionSummaryVO.READY);
+ }
+
@Test
void fromShouldReportMixedVersionsAndDuplicateClients() {
ProducerConnectionSummaryVO summary =
ProducerConnectionSummaryVO.from(List.of(
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepositoryTest.java
index 4f3ce18f..1ea2aa0c 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/k8s/MybatisPlusK8sCertRepositoryTest.java
@@ -22,12 +22,46 @@ import
org.apache.rocketmq.studio.persistence.entity.RmqK8sCertificate;
import org.apache.rocketmq.studio.persistence.mapper.RmqK8sCertificateMapper;
import org.junit.jupiter.api.Test;
+import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.Mockito.mock;
+import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.when;
class MybatisPlusK8sCertRepositoryTest {
+ @Test
+ void findByIdNormalizesPersistedCertificateEnums() {
+ RmqK8sCertificateMapper mapper = mock(RmqK8sCertificateMapper.class);
+ RmqK8sCertificate entity = certificate();
+ entity.setCertType(" mtls ");
+ entity.setStatus(" EXPIRING ");
+ when(mapper.selectById("cert-1")).thenReturn(entity);
+
+ assertThat(repository(mapper).findById("cert-1")).get()
+ .satisfies(cert -> {
+ assertThat(cert.getType()).isEqualTo(
+
org.apache.rocketmq.studio.common.domain.enums.CertType.mTLS);
+ assertThat(cert.getStatus()).isEqualTo(
+
org.apache.rocketmq.studio.common.domain.enums.CertStatus.expiring);
+ });
+ }
+
+ @Test
+ void saveShouldReportALostConcurrentUpdate() {
+ RmqK8sCertificateMapper mapper = mock(RmqK8sCertificateMapper.class);
+ when(mapper.selectById("cert-1")).thenReturn(certificate());
+ when(mapper.updateById(any(RmqK8sCertificate.class))).thenReturn(0);
+ K8sCertVO cert = K8sCertVO.builder().name("broker").build();
+ cert.setId("cert-1");
+
+ assertThatThrownBy(() -> repository(mapper).save(cert))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Certificate update was not applied: cert-1")
+ .satisfies(error -> org.assertj.core.api.Assertions.assertThat(
+ ((BusinessException) error).getCode()).isEqualTo(409));
+ }
+
@Test
void findByIdSurfacesInvalidPersistedCertificateType() {
RmqK8sCertificateMapper mapper = mock(RmqK8sCertificateMapper.class);
@@ -52,6 +86,18 @@ class MybatisPlusK8sCertRepositoryTest {
.hasMessageContaining("SAN JSON");
}
+ @Test
+ void findByIdSurfacesNullPersistedSanJson() {
+ RmqK8sCertificateMapper mapper = mock(RmqK8sCertificateMapper.class);
+ RmqK8sCertificate entity = certificate();
+ entity.setSan("null");
+ when(mapper.selectById("cert-1")).thenReturn(entity);
+
+ assertThatThrownBy(() -> repository(mapper).findById("cert-1"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessageContaining("SAN JSON");
+ }
+
@Test
void findByIdSurfacesInvalidPersistedCertificateStatus() {
RmqK8sCertificateMapper mapper = mock(RmqK8sCertificateMapper.class);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
index 353d0ab1..1c9be5b2 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
@@ -422,6 +422,28 @@ class MetricsServiceTest {
verifyNoInteractions(metricsSourceFactory, metricsSource);
}
+ @Test
+ void queryByDataSourceShouldRejectUnsupportedStoredAuthMode() {
+ MetricsDataSourceQueryRequest request = dataSourceRequest(null);
+ DataSourceVO dataSource = DataSourceVO.builder()
+ .key("ds-1")
+ .name("prometheus-a")
+ .type("prometheus")
+ .url("http://prometheus:9090")
+ .auth("digest")
+ .build();
+ when(settingsService.getDataSource("ds-1")).thenReturn(dataSource);
+
+ assertThatExceptionOfType(PrometheusException.class)
+ .isThrownBy(() -> metricsService.queryByDataSource("ds-1",
request))
+ .satisfies(exception -> {
+ assertThat(exception.getStatusCode()).isEqualTo(400);
+ assertThat(exception.getMessage())
+ .isEqualTo("Unsupported data source authentication
mode: digest");
+ });
+ verifyNoInteractions(metricsSourceFactory, metricsSource);
+ }
+
@Test
void queryByDataSourceShouldNormalizeAuthIndependentlyOfDefaultLocale() {
MetricQueryDTO query = MetricQueryDTO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSourceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSourceTest.java
index bb31ed8b..b228393c 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSourceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSourceTest.java
@@ -229,6 +229,32 @@ class PrometheusMetricsSourceTest {
.hasMessage("Prometheus returned a malformed response");
}
+ @Test
+ void queryShouldRejectNonTextSampleValues() {
+ server.createContext("/api/v1/query_range", exchange ->
respond(exchange, 200, """
+ {"status":"success","data":{"resultType":"matrix","result":[
+ {"metric":{"topic":"orders"},"values":[[1784107658,1.25]]}
+ ]}}
+ """));
+
+ assertThatThrownBy(() -> source(Duration.ofSeconds(2)).query(query()))
+ .isInstanceOf(PrometheusException.class)
+ .hasMessage("Prometheus returned a malformed sample");
+ }
+
+ @Test
+ void queryShouldRejectNonTextSeriesLabels() {
+ server.createContext("/api/v1/query_range", exchange ->
respond(exchange, 200, """
+ {"status":"success","data":{"resultType":"matrix","result":[
+
{"metric":{"topic":{"name":"orders"}},"values":[[1784107658,"1.25"]]}
+ ]}}
+ """));
+
+ assertThatThrownBy(() -> source(Duration.ofSeconds(2)).query(query()))
+ .isInstanceOf(PrometheusException.class)
+ .hasMessage("Prometheus returned a malformed time-series
label");
+ }
+
@Test
void queryShouldRejectResponseWithTooManySeries() {
server.createContext("/api/v1/query_range", exchange ->
respond(exchange, 200, responseWithSeries(1_001)));
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardServiceTest.java
index 7e7fc028..b28644d5 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/grafana/GrafanaDashboardServiceTest.java
@@ -65,6 +65,17 @@ class GrafanaDashboardServiceTest {
assertTrue(model.containsKey("panels"), "dashboard should contain
panels");
}
+ @Test
+ void getDashboardShouldRejectNullJsonAsset() {
+ GrafanaDashboardService service =
serviceWithResources(resource("null.json", "null"));
+
+ BusinessException exception = assertThrows(BusinessException.class,
+ () -> service.getDashboard("null"));
+
+ assertEquals(500, exception.getCode());
+ assertEquals("Failed to read Grafana dashboard: null",
exception.getMessage());
+ }
+
@Test
void getDashboardShouldThrowWhenUidUnknown() {
BusinessException exception = assertThrows(BusinessException.class,
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
b/server/src/test/java/org/apache/rocketmq/studio/common/domain/PageResultTest.java
similarity index 67%
copy from
server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
copy to
server/src/test/java/org/apache/rocketmq/studio/common/domain/PageResultTest.java
index 7c3fd411..7bcb3dad 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/domain/PageResultTest.java
@@ -14,19 +14,18 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.ops.audit;
+package org.apache.rocketmq.studio.common.domain;
-import jakarta.validation.constraints.Positive;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
+import org.junit.jupiter.api.Test;
-@Data
-@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class AuditCleanupDTO {
- @Positive(message = "beforeDays must be greater than 0")
- private Integer beforeDays;
+import static org.assertj.core.api.Assertions.assertThat;
+
+class PageResultTest {
+
+ @Test
+ void ofShouldNormalizeNullItemsToAnEmptyList() {
+ PageResult<String> result = PageResult.of(null, 0, 1, 20);
+
+ assertThat(result.getItems()).isNotNull().isEmpty();
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
b/server/src/test/java/org/apache/rocketmq/studio/common/util/CredentialUtilsTest.java
similarity index 63%
copy from
server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
copy to
server/src/test/java/org/apache/rocketmq/studio/common/util/CredentialUtilsTest.java
index 7c3fd411..698b6c56 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/util/CredentialUtilsTest.java
@@ -14,19 +14,20 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.ops.audit;
+package org.apache.rocketmq.studio.common.util;
-import jakarta.validation.constraints.Positive;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
+import org.junit.jupiter.api.Test;
-@Data
-@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class AuditCleanupDTO {
- @Positive(message = "beforeDays must be greater than 0")
- private Integer beforeDays;
+import java.util.Base64;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class CredentialUtilsTest {
+
+ @Test
+ void decodeBase64ShouldPreserveLegacyTextWhenDecodedBytesAreNotUtf8() {
+ String legacyValue = Base64.getEncoder().encodeToString(new
byte[]{(byte) 0xC3, 0x28});
+
+
assertThat(CredentialUtils.decodeBase64(legacyValue)).isEqualTo(legacyValue);
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
b/server/src/test/java/org/apache/rocketmq/studio/common/util/UrlHostGuardMulticastTest.java
similarity index 66%
copy from
server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
copy to
server/src/test/java/org/apache/rocketmq/studio/common/util/UrlHostGuardMulticastTest.java
index 7c3fd411..df0c21d3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/util/UrlHostGuardMulticastTest.java
@@ -14,19 +14,17 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.ops.audit;
+package org.apache.rocketmq.studio.common.util;
-import jakarta.validation.constraints.Positive;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
+import org.junit.jupiter.api.Test;
-@Data
-@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class AuditCleanupDTO {
- @Positive(message = "beforeDays must be greater than 0")
- private Integer beforeDays;
+import static org.assertj.core.api.Assertions.assertThat;
+
+class UrlHostGuardMulticastTest {
+
+ @Test
+ void isAllowedHostShouldRejectIpv4AndIpv6MulticastAddresses() {
+ assertThat(UrlHostGuard.isAllowedHost("224.0.0.1", false)).isFalse();
+ assertThat(UrlHostGuard.isAllowedHost("ff02::1", false)).isFalse();
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/common/util/UrlHostGuardTest.java
b/server/src/test/java/org/apache/rocketmq/studio/common/util/UrlHostGuardTest.java
new file mode 100644
index 00000000..f31ee265
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/util/UrlHostGuardTest.java
@@ -0,0 +1,43 @@
+/*
+ * 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.common.util;
+
+import org.junit.jupiter.api.Test;
+
+import java.net.InetAddress;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class UrlHostGuardTest {
+
+ @Test
+ void areAllowedShouldRejectAHostWhenAnyResolvedAddressIsDisallowed()
throws Exception {
+ InetAddress publicAddress = InetAddress.getByAddress(new byte[]{8, 8,
8, 8});
+ InetAddress loopbackAddress = InetAddress.getByAddress(new byte[]{127,
0, 0, 1});
+
+ assertThat(UrlHostGuard.areAllowed(
+ new InetAddress[]{publicAddress, loopbackAddress},
false)).isFalse();
+ }
+
+ @Test
+ void areAllowedShouldAcceptEveryPublicResolvedAddress() throws Exception {
+ InetAddress first = InetAddress.getByAddress(new byte[]{8, 8, 8, 8});
+ InetAddress second = InetAddress.getByAddress(new byte[]{1, 1, 1, 1});
+
+ assertThat(UrlHostGuard.areAllowed(new InetAddress[]{first, second},
false)).isTrue();
+ }
+}
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 b29a8bd4..b428cc1b 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
@@ -35,6 +35,7 @@ import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import java.time.LocalDateTime;
+import java.util.Arrays;
import java.util.List;
import java.util.Locale;
import java.util.Optional;
@@ -791,6 +792,28 @@ class InstanceServiceTest {
.hasMessageContaining("credentialId");
}
+ @Test
+ void createInstanceShouldRejectMissingCloudCatalogDetails() {
+ InstanceVO instance = InstanceVO.builder()
+ .vendor(InstanceVendor.ALIYUN)
+ .credentialId("cred-1")
+ .cloudInstanceId("rmq-missing")
+ .regionId("cn-hangzhou")
+ .build();
+ CloudCredentialVO credential = new CloudCredentialVO();
+ credential.setVendor(InstanceVendor.ALIYUN);
+
when(cloudCredentialRepository.findById("cred-1")).thenReturn(Optional.of(credential));
+ CloudCatalogProvider catalog =
org.mockito.Mockito.mock(CloudCatalogProvider.class);
+
when(providerRegistry.catalogFor(InstanceVendor.ALIYUN)).thenReturn(catalog);
+ when(catalog.getCloudInstance("cred-1", "cn-hangzhou",
"rmq-missing")).thenReturn(null);
+
+ assertThatThrownBy(() -> instanceService.createInstance(instance))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Cloud instance details unavailable: rmq-missing")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(502));
+ verify(instanceRepository, never()).save(any());
+ }
+
@Test
void createInstanceShouldResolveAliyunEndpointFromCatalogTest() {
InstanceVO instance = InstanceVO.builder()
@@ -821,6 +844,30 @@ class InstanceServiceTest {
assertThat(created.getType()).isEqualTo(InstanceType.PROXY);
}
+ @Test
+ void createInstanceShouldSkipNullCloudEndpointEntries() {
+ InstanceVO instance = InstanceVO.builder()
+ .vendor(InstanceVendor.ALIYUN)
+ .credentialId("cred-1")
+ .cloudInstanceId("rmq-cn-xxx")
+ .regionId("cn-hangzhou")
+ .build();
+ CloudCredentialVO credential = new CloudCredentialVO();
+ credential.setVendor(InstanceVendor.ALIYUN);
+
when(cloudCredentialRepository.findById("cred-1")).thenReturn(Optional.of(credential));
+ CloudCatalogProvider catalog =
org.mockito.Mockito.mock(CloudCatalogProvider.class);
+ CloudInstanceDetailVO detail = new CloudInstanceDetailVO();
+ detail.setInstanceId("rmq-cn-xxx");
+ detail.setInstanceName("prod-mq");
+ detail.setEndpoints(Arrays.asList(null,
+ new CloudInstanceDetailVO.CloudEndpoint("TCP_VPC",
"vpc:8080")));
+
when(providerRegistry.catalogFor(InstanceVendor.ALIYUN)).thenReturn(catalog);
+ when(catalog.getCloudInstance("cred-1", "cn-hangzhou",
"rmq-cn-xxx")).thenReturn(detail);
+ when(instanceRepository.save(any())).thenAnswer(invocation ->
invocation.getArgument(0));
+
+
assertThat(instanceService.createInstance(instance).getEndpoint()).isEqualTo("vpc:8080");
+ }
+
@Test
void
createInstanceShouldPrioritizeEndpointsIndependentlyOfDefaultLocaleTest() {
InstanceVO instance = InstanceVO.builder()
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
index 731b26ba..a280fc68 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/MybatisPlusInstanceRepositoryTest.java
@@ -172,6 +172,7 @@ class MybatisPlusInstanceRepositoryTest {
void saveShouldUpdateWhenInstanceExists() {
InstanceVO vo = vo("instance-proxy-2", InstanceType.PROXY);
when(instanceMapper.selectById("instance-proxy-2")).thenReturn(entity("instance-proxy-2",
InstanceType.PROXY));
+ when(instanceMapper.updateById(any(RmqInstance.class))).thenReturn(1);
repository.save(vo);
@@ -179,6 +180,19 @@ class MybatisPlusInstanceRepositoryTest {
verify(instanceMapper, never()).insert(any(RmqInstance.class));
}
+ @Test
+ void saveShouldReportALostConcurrentUpdate() {
+ InstanceVO vo = vo("instance-proxy-2", InstanceType.PROXY);
+ when(instanceMapper.selectById("instance-proxy-2"))
+ .thenReturn(entity("instance-proxy-2", InstanceType.PROXY));
+ when(instanceMapper.updateById(any(RmqInstance.class))).thenReturn(0);
+
+ assertThatThrownBy(() -> repository.save(vo))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Instance update was not applied:
instance-proxy-2")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(409));
+ }
+
@Test
void deleteByIdShouldDelegateToMapper() {
repository.deleteById("instance-direct-1");
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
index 9365c246..71477e5f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/acl/MybatisPlusAclRepositoryTest.java
@@ -57,6 +57,23 @@ class MybatisPlusAclRepositoryTest {
@InjectMocks
private MybatisPlusAclRepository repository;
+ @Test
+ void replaceRuleShouldReturnEmptyWhenConcurrentDeleteWins() {
+ RmqAclRule existing = new RmqAclRule();
+ existing.setId("rule-a");
+ existing.setCreatedAt(LocalDateTime.of(2026, 1, 1, 0, 0));
+ when(ruleMapper.selectById("rule-a")).thenReturn(existing);
+ when(ruleMapper.updateById(any(RmqAclRule.class))).thenReturn(0);
+
+ AclRuleVO replacement = AclRuleVO.builder()
+ .id("rule-a")
+ .principal("svc-a")
+ .resource("orders")
+ .build();
+
+ assertThat(repository.replaceRule(replacement)).isEmpty();
+ }
+
@Test
void upsertShouldAssignUniqueRuleIdPerPermission() {
when(userMapper.selectOne(any(QueryWrapper.class))).thenReturn(null);
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/AgentProviderRegistryTest.java
similarity index 50%
copy from
server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
copy to
server/src/test/java/org/apache/rocketmq/studio/ops/ai/AgentProviderRegistryTest.java
index ae7507cd..446e1543 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/AgentProviderRegistryTest.java
@@ -14,31 +14,31 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.common.domain;
+package org.apache.rocketmq.studio.ops.ai;
+
+import org.junit.jupiter.api.Test;
-import lombok.Getter;
-import java.util.Collections;
import java.util.List;
+import java.util.Locale;
-@Getter
-public class PageResult<T> {
- private List<T> items;
- private long total;
- private int page;
- private int size;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
- private PageResult() {}
+class AgentProviderRegistryTest {
- public static <T> PageResult<T> of(List<T> items, long total, int page,
int size) {
- PageResult<T> result = new PageResult<>();
- result.items = items;
- result.total = total;
- result.page = page;
- result.size = size;
- return result;
- }
+ @Test
+ void shouldResolveEngineIndependentlyOfDefaultLocale() {
+ AgentProvider provider = mock(AgentProvider.class);
+ when(provider.engine()).thenReturn("cli");
+ AgentProviderRegistry registry = new
AgentProviderRegistry(List.of(provider));
+ Locale original = Locale.getDefault();
+ try {
+ Locale.setDefault(Locale.forLanguageTag("tr-TR"));
- public static <T> PageResult<T> empty(int page, int size) {
- return of(Collections.emptyList(), 0, page, size);
+ assertThat(registry.forEngine(" CLI ")).isSameAs(provider);
+ } finally {
+ Locale.setDefault(original);
+ }
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigVOTest.java
similarity index 56%
copy from
server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
copy to
server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigVOTest.java
index 7c3fd411..4376b384 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/ai/LlmConfigVOTest.java
@@ -14,19 +14,26 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.ops.audit;
+package org.apache.rocketmq.studio.ops.ai;
-import jakarta.validation.constraints.Positive;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
+import org.junit.jupiter.api.Test;
-@Data
-@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class AuditCleanupDTO {
- @Positive(message = "beforeDays must be greater than 0")
- private Integer beforeDays;
+import java.util.Locale;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class LlmConfigVOTest {
+
+ @Test
+ void shouldNormalizeEngineIndependentlyOfDefaultLocale() {
+ Locale original = Locale.getDefault();
+ try {
+ Locale.setDefault(Locale.forLanguageTag("tr-TR"));
+ LlmConfigVO config = LlmConfigVO.builder().engine(" CLI ").build();
+
+ assertThat(config.normalizeEngine()).isEqualTo("cli");
+ } finally {
+ Locale.setDefault(original);
+ }
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetServiceTest.java
index b3eea9fc..f60aae58 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleAssetServiceTest.java
@@ -20,7 +20,9 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import org.junit.jupiter.api.Test;
import org.springframework.core.io.ByteArrayResource;
import org.springframework.core.io.Resource;
+import org.springframework.core.io.support.ResourcePatternResolver;
+import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.List;
@@ -28,6 +30,9 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
class AlertRuleAssetServiceTest {
@@ -76,6 +81,18 @@ class AlertRuleAssetServiceTest {
assertEquals(404, exception.getCode());
}
+ @Test
+ void listAssetsShouldSurfaceResourceDiscoveryFailures() throws IOException
{
+ ResourcePatternResolver resolver = mock(ResourcePatternResolver.class);
+ when(resolver.getResources(anyString())).thenThrow(new
IOException("classpath unavailable"));
+ AlertRuleAssetService failingService = new
AlertRuleAssetService(resolver);
+
+ BusinessException exception = assertThrows(BusinessException.class,
failingService::listAssets);
+
+ assertEquals(500, exception.getCode());
+ assertEquals("Failed to resolve bundled alert rule assets",
exception.getMessage());
+ }
+
@Test
void parseRulesShouldMapSeverityAndTeamLabels() {
List<PrometheusAlertRule> rules = service.loadDefaultRules();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
index ecf8f60a..8c43b4fd 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertServiceTest.java
@@ -178,6 +178,22 @@ class AlertServiceTest {
.contains("- alert: HighLag_2\n");
}
+ @Test
+ void
exportPrometheusRulesYamlShouldInferTeamWithoutMetricCaseSensitivity() {
+ AlertRuleVO rule = AlertRuleVO.builder()
+ .name("Uppercase lag")
+ .metric("ROCKETMQ_CONSUMER_LAG_MESSAGES")
+ .operator(">")
+ .threshold(10)
+ .enabled(true)
+ .build();
+ when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+
+ String result = alertService.exportPrometheusRulesYaml();
+
+ assertThat(result).contains("- name: rocketmq-consumer.rules");
+ }
+
@Test
void
exportPrometheusRulesYamlShouldRenderReplicationLagRuleWithScopeAndSeverity() {
AlertRuleVO rule = AlertRuleVO.builder()
@@ -268,6 +284,27 @@ class AlertServiceTest {
assertThat(result).contains("severity: info");
}
+ @Test
+ void exportPrometheusRulesYamlShouldNotEmitABlankSummary() throws
Exception {
+ AlertRuleVO rule = AlertRuleVO.builder()
+ .name("Lag alert")
+ .metric("rocketmq_consumer_lag_messages")
+ .operator(">")
+ .threshold(1)
+ .description(" - consumer is behind")
+ .enabled(true)
+ .build();
+ when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+
+ JsonNode exportedRule = new ObjectMapper(new YAMLFactory())
+ .readTree(alertService.exportPrometheusRulesYaml())
+ .path("groups").get(0).path("rules").get(0);
+
+
assertThat(exportedRule.path("annotations").path("summary").asText()).isEqualTo("Lag
alert");
+
assertThat(exportedRule.path("annotations").path("description").asText())
+ .isEqualTo("consumer is behind");
+ }
+
@Test
void exportPrometheusRulesYamlShouldEscapeSpecialCharacters() throws
Exception {
AlertRuleVO rule = AlertRuleVO.builder()
@@ -502,6 +539,25 @@ class AlertServiceTest {
assertThat(result.isEnabled()).isFalse();
}
+ @Test
+ void toggleRuleShouldRejectBlankIdBeforeLoadingRules() {
+ assertThatThrownBy(() -> alertService.toggleRule(" ", true))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Alert rule ID is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verify(alertRepository, never()).findAllRules();
+ }
+
+ @Test
+ void toggleRuleShouldIgnorePersistedRulesWithNullIds() {
+
when(alertRepository.findAllRules()).thenReturn(List.of(AlertRuleVO.builder().name("corrupt").build()));
+
+ assertThatThrownBy(() -> alertService.toggleRule("missing", true))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Alert rule not found: missing");
+ }
+
@Test
void toggleRuleShouldThrowWhenRuleNotFound() {
when(alertRepository.findAllRules()).thenReturn(Collections.emptyList());
@@ -511,6 +567,16 @@ class AlertServiceTest {
.hasMessageContaining("Alert rule not found: non-existent");
}
+ @Test
+ void deleteRuleShouldRejectBlankIdBeforeDeleting() {
+ assertThatThrownBy(() -> alertService.deleteRule(" "))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Alert rule ID is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verify(alertRepository, never()).deleteRule(any());
+ }
+
@Test
void deleteRuleShouldCallRepository() {
when(alertRepository.deleteRule("rule-1")).thenReturn(true);
@@ -575,6 +641,25 @@ class AlertServiceTest {
eq(null), eq("acknowledged=true"), eq("SUCCESS"), eq(null));
}
+ @Test
+ void acknowledgeAlertShouldRejectBlankIdBeforeLoadingAlerts() {
+ assertThatThrownBy(() -> alertService.acknowledgeAlert(" "))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("System alert ID is required")
+ .satisfies(ex -> assertThat(((BusinessException)
ex).getCode()).isEqualTo(400));
+
+ verify(alertRepository, never()).findAlerts(any());
+ }
+
+ @Test
+ void acknowledgeAlertShouldIgnorePersistedAlertsWithNullIds() {
+
when(alertRepository.findAlerts(null)).thenReturn(List.of(SystemAlertVO.builder().title("corrupt").build()));
+
+ assertThatThrownBy(() -> alertService.acknowledgeAlert("missing"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("System alert not found: missing");
+ }
+
@Test
void acknowledgeAlertShouldThrowWhenAlertNotFound() {
when(alertRepository.findAlerts(null)).thenReturn(Collections.emptyList());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
index 050b42f5..7b7a0458 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/MybatisPlusAlertRepositoryTest.java
@@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.ops.alert;
import com.baomidou.mybatisplus.core.conditions.Wrapper;
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import org.apache.rocketmq.studio.persistence.entity.RmqAlertRule;
import org.apache.rocketmq.studio.persistence.entity.RmqSystemAlert;
import org.apache.rocketmq.studio.persistence.mapper.RmqAlertRuleMapper;
import org.apache.rocketmq.studio.persistence.mapper.RmqSystemAlertMapper;
@@ -30,6 +31,7 @@ import org.mockito.junit.jupiter.MockitoExtension;
import java.util.List;
import java.util.Locale;
+import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.verify;
@@ -47,6 +49,26 @@ class MybatisPlusAlertRepositoryTest {
@InjectMocks
private MybatisPlusAlertRepository repository;
+ @Test
+ void replaceRuleShouldReportAConcurrentDelete() {
+ AlertRuleVO rule =
AlertRuleVO.builder().id("rule-1").name("Lag").build();
+ when(ruleMapper.selectById("rule-1")).thenReturn(new RmqAlertRule());
+ when(ruleMapper.updateById(any(RmqAlertRule.class))).thenReturn(0);
+
+ assertThat(repository.replaceRule(rule)).isFalse();
+ }
+
+ @Test
+ void findAlertsShouldNormalizeStoredLevelValues() {
+ RmqSystemAlert entity = new RmqSystemAlert();
+ entity.setLevel(" WARNING ");
+ when(alertMapper.selectList(any())).thenReturn(List.of(entity));
+
+ assertThat(repository.findAlerts(null)).singleElement()
+ .satisfies(alert -> assertThat(alert.getLevel()).isEqualTo(
+
org.apache.rocketmq.studio.common.domain.enums.AlertLevel.warning));
+ }
+
@Test
void findAlertsShouldNormalizeLevelIndependentlyOfDefaultLocale() {
when(alertMapper.selectList(any())).thenReturn(List.of());
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditControllerTest.java
index aefc1eba..673d2683 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditControllerTest.java
@@ -134,6 +134,18 @@ class AuditControllerTest {
verifyNoInteractions(auditService);
}
+ @Test
+ void cleanupLogsShouldRejectRetentionBeyondMaximum() throws Exception {
+ mockMvc.perform(post("/api/audit-logs/cleanup")
+ .contentType(MediaType.APPLICATION_JSON)
+
.content(objectMapper.writeValueAsString(Map.of("beforeDays", 366))))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("beforeDays must not
exceed 365"));
+
+ verifyNoInteractions(auditService);
+ }
+
@Test
void cleanupLogsShouldRejectInvalidRetentionType() throws Exception {
mockMvc.perform(post("/api/audit-logs/cleanup")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
index 43593b19..08b8bb51 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/audit/AuditServiceTest.java
@@ -182,4 +182,11 @@ class AuditServiceTest {
.isInstanceOf(BusinessException.class)
.hasMessage("beforeDays must be greater than 0");
}
+
+ @Test
+ void cleanupLogsRejectsRetentionBeyondMaximum() {
+ assertThatThrownBy(() -> auditService.cleanupLogs(366))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("beforeDays must not exceed 365");
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepositoryTest.java
index 84ad68d4..c71a1a77 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/persistence/MybatisPlusSettingsRepositoryTest.java
@@ -12,6 +12,7 @@ import
org.apache.rocketmq.studio.persistence.entity.RmqDataSource;
import org.apache.rocketmq.studio.persistence.entity.RmqSettings;
import org.apache.rocketmq.studio.persistence.mapper.RmqDataSourceMapper;
import org.apache.rocketmq.studio.persistence.mapper.RmqSettingsMapper;
+import org.apache.rocketmq.studio.settings.DataSourceVO;
import org.apache.rocketmq.studio.settings.GeneralSettingsVO;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -58,6 +59,19 @@ class MybatisPlusSettingsRepositoryTest {
.isEqualTo(500);
}
+ @Test
+ void shouldRejectNullPersistedGeneralSettings() {
+ RmqSettings settings = new RmqSettings();
+ settings.setJson("null");
+ when(settingsMapper.selectById("singleton")).thenReturn(settings);
+
+ assertThatThrownBy(repository::loadGeneralSettings)
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Persisted general settings are invalid")
+ .extracting("code")
+ .isEqualTo(500);
+ }
+
@Test
void shouldReadValidPersistedGeneralSettings() {
RmqSettings settings = new RmqSettings();
@@ -70,6 +84,20 @@ class MybatisPlusSettingsRepositoryTest {
assertThat(loaded.isRequireLogin()).isTrue();
}
+ @Test
+ void shouldRejectNullPersistedDataSource() {
+ RmqDataSource dataSource = new RmqDataSource();
+ dataSource.setDsKey("metrics-prod");
+ dataSource.setJson("null");
+
when(dataSourceMapper.selectById("metrics-prod")).thenReturn(dataSource);
+
+ assertThatThrownBy(() ->
repository.findDataSourceByKey("metrics-prod"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Persisted data source is invalid: metrics-prod")
+ .extracting("code")
+ .isEqualTo(500);
+ }
+
@Test
void shouldRejectCorruptPersistedDataSource() {
RmqDataSource dataSource = new RmqDataSource();
@@ -83,4 +111,20 @@ class MybatisPlusSettingsRepositoryTest {
.extracting("code")
.isEqualTo(500);
}
+
+ @Test
+ void shouldReportWhenDataSourceDisappearsDuringReplacement() {
+ RmqDataSource existing = new RmqDataSource();
+ existing.setDsKey("metrics-prod");
+ when(dataSourceMapper.selectById("metrics-prod")).thenReturn(existing);
+ when(dataSourceMapper.updateById(existing)).thenReturn(0);
+ DataSourceVO replacement = DataSourceVO.builder()
+ .key("metrics-prod")
+ .name("Production metrics")
+ .type("prometheus")
+ .url("https://metrics.example.com")
+ .build();
+
+ assertThat(repository.replaceDataSource(replacement)).isFalse();
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
index cacd8388..5be042c6 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunCatalogServiceTest.java
@@ -33,6 +33,7 @@ import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.List;
import static org.assertj.core.api.Assertions.assertThat;
@@ -59,6 +60,25 @@ class AliyunCatalogServiceTest {
service = new AliyunCatalogService(clientFactory);
}
+ @Test
+ void listRegionsShouldSkipNullSdkRecords() {
+ ListRegionsResponse response = ListRegionsResponse.create().toBuilder()
+ .statusCode(200)
+ .body(ListRegionsResponseBody.builder()
+ .data(Arrays.asList(null,
ListRegionsResponseBody.Data.builder()
+ .regionId("cn-hangzhou")
+ .supportRocketmqV5(true)
+ .build()))
+ .build())
+ .build();
+ when(clientFactory.call(eq(CREDENTIAL_ID),
eq(AliyunCatalogService.DEFAULT_REGION), any()))
+ .thenReturn(response);
+
+ assertThat(service.listRegions(CREDENTIAL_ID))
+ .extracting(CloudRegionVO::getRegionId)
+ .containsExactly("cn-hangzhou");
+ }
+
@Test
void listRegionsShouldKeepOnlyRocketmqV5RegionsTest() {
ListRegionsResponse response = ListRegionsResponse.create().toBuilder()
@@ -90,6 +110,16 @@ class AliyunCatalogServiceTest {
assertThat(regions.get(1).getRegionName()).isEqualTo("hangzhou");
}
+ @Test
+ void listCloudInstancesShouldSkipNullSdkRecords() {
+ when(clientFactory.call(eq(CREDENTIAL_ID), eq(REGION), any()))
+ .thenReturn(instancesResponse(Arrays.asList(null,
instanceRow("rmq-a", "A"))));
+
+ assertThat(service.listCloudInstances(CREDENTIAL_ID, REGION, null))
+ .extracting(CloudInstanceOptionVO::getInstanceId)
+ .containsExactly("rmq-a");
+ }
+
@Test
void listCloudInstancesShouldAggregatePagesTest() {
ListInstancesResponse firstPage =
instancesResponse(instanceRows(AliyunConverters.PAGE_SIZE, 0));
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersMessageKeysTest.java
similarity index 55%
copy from
server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
copy to
server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersMessageKeysTest.java
index 7c3fd411..f10d8e14 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/audit/AuditCleanupDTO.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersMessageKeysTest.java
@@ -14,19 +14,23 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.ops.audit;
+package org.apache.rocketmq.studio.provider.alibaba;
-import jakarta.validation.constraints.Positive;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
+import com.aliyun.sdk.service.rocketmq20220801.models.ListMessagesResponseBody;
+import org.junit.jupiter.api.Test;
-@Data
-@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class AuditCleanupDTO {
- @Positive(message = "beforeDays must be greater than 0")
- private Integer beforeDays;
+import java.util.Arrays;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class AliyunConvertersMessageKeysTest {
+
+ @Test
+ void toMessageRecordShouldSkipNullAndBlankMessageKeys() {
+ ListMessagesResponseBody.List data =
ListMessagesResponseBody.List.builder()
+ .messageKeys(Arrays.asList("key-a", null, " ", "key-b"))
+ .build();
+
+
assertThat(AliyunConverters.toMessageRecord(data).getKey()).isEqualTo("key-a
key-b");
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersNullEndpointTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersNullEndpointTest.java
new file mode 100644
index 00000000..20b63651
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersNullEndpointTest.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
+ *
+ * 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.provider.alibaba;
+
+import com.aliyun.sdk.service.rocketmq20220801.models.GetInstanceResponseBody;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class AliyunConvertersNullEndpointTest {
+
+ @Test
+ void toInstanceDetailShouldSkipNullEndpointEntries() {
+ GetInstanceResponseBody.Endpoints endpoint =
GetInstanceResponseBody.Endpoints.builder()
+ .endpointType("TCP_VPC")
+ .endpointUrl("10.0.0.1:8080")
+ .build();
+ GetInstanceResponseBody.NetworkInfo network =
GetInstanceResponseBody.NetworkInfo.builder()
+ .endpoints(Arrays.asList(null, endpoint))
+ .build();
+ GetInstanceResponseBody.Data data =
GetInstanceResponseBody.Data.builder()
+ .instanceId("rmq-a")
+ .networkInfo(network)
+ .build();
+
+ assertThat(AliyunConverters.toInstanceDetailVO(data).getEndpoints())
+ .singleElement()
+ .satisfies(value ->
assertThat(value.getEndpointUrl()).isEqualTo("10.0.0.1:8080"));
+ }
+}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersTest.java
similarity index 52%
copy from
server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
copy to
server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersTest.java
index ae7507cd..a47ec89c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/domain/PageResult.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersTest.java
@@ -14,31 +14,25 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.common.domain;
+package org.apache.rocketmq.studio.provider.alibaba;
-import lombok.Getter;
-import java.util.Collections;
-import java.util.List;
+import
com.aliyun.sdk.service.rocketmq20220801.models.ListInstancesResponseBody;
+import org.junit.jupiter.api.Test;
-@Getter
-public class PageResult<T> {
- private List<T> items;
- private long total;
- private int page;
- private int size;
+import static org.assertj.core.api.Assertions.assertThat;
- private PageResult() {}
+class AliyunConvertersTest {
- public static <T> PageResult<T> of(List<T> items, long total, int page,
int size) {
- PageResult<T> result = new PageResult<>();
- result.items = items;
- result.total = total;
- result.page = page;
- result.size = size;
- return result;
- }
+ @Test
+ void toInstanceOptionShouldClampCountsOutsideTheIntegerRange() {
+ ListInstancesResponseBody.List data =
ListInstancesResponseBody.List.builder()
+ .topicCount(Long.MAX_VALUE)
+ .groupCount(Long.MIN_VALUE)
+ .build();
+
+ var result = AliyunConverters.toInstanceOptionVO(data);
- public static <T> PageResult<T> empty(int page, int size) {
- return of(Collections.emptyList(), 0, page, size);
+ assertThat(result.getTopicCount()).isEqualTo(Integer.MAX_VALUE);
+ assertThat(result.getGroupCount()).isZero();
}
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersTraceElementsTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersTraceElementsTest.java
new file mode 100644
index 00000000..31e07770
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/alibaba/AliyunConvertersTraceElementsTest.java
@@ -0,0 +1,47 @@
+/*
+ * 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.provider.alibaba;
+
+import com.aliyun.sdk.service.rocketmq20220801.models.GetTraceResponseBody;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class AliyunConvertersTraceElementsTest {
+
+ @Test
+ void toTraceRecordShouldSkipNullSdkListElements() {
+ GetTraceResponseBody.ProducerInfo producer =
GetTraceResponseBody.ProducerInfo.builder()
+
.records(Arrays.asList((GetTraceResponseBody.ProducerInfoRecords) null))
+ .build();
+ GetTraceResponseBody.BrokerInfo broker =
GetTraceResponseBody.BrokerInfo.builder()
+ .operations(Arrays.asList((GetTraceResponseBody.Operations)
null))
+ .build();
+ GetTraceResponseBody.ConsumerInfos consumer =
GetTraceResponseBody.ConsumerInfos.builder()
+ .records(Arrays.asList((GetTraceResponseBody.Records) null))
+ .build();
+ GetTraceResponseBody.Data data = GetTraceResponseBody.Data.builder()
+ .producerInfo(producer)
+ .brokerInfo(broker)
+ .consumerInfos(Arrays.asList(null, consumer))
+ .build();
+
+ assertThat(AliyunConverters.toTraceRecord(data).getNodes()).isEmpty();
+ }
+}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index 1fa3d303..6ff3e36e 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -158,6 +158,35 @@ class RocketMQMessageProviderTest {
}
}
+ @Test
+ void queryByTopicReturnsNewestMessagesFirst() throws Exception {
+ MessageQueue queue = new MessageQueue("TopicA", "broker-a", 0);
+ MessageExt older = new MessageExt();
+ older.setMsgId("older");
+ older.setTopic("TopicA");
+ older.setStoreTimestamp(150L);
+ MessageExt newer = new MessageExt();
+ newer.setMsgId("newer");
+ newer.setTopic("TopicA");
+ newer.setStoreTimestamp(250L);
+ PullResult pullResult = new PullResult(PullStatus.FOUND, 11L, 10L, 11L,
+ List.of(older, newer));
+ try (MockedConstruction<DefaultMQPullConsumer> ignored =
+ mockConstruction(DefaultMQPullConsumer.class, (consumer,
context) -> {
+ doNothing().when(consumer).start();
+
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+ when(consumer.searchOffset(eq(queue),
anyLong())).thenReturn(10L);
+ when(consumer.pull(queue, "*", 10L,
32)).thenReturn(pullResult);
+ doNothing().when(consumer).shutdown();
+ })) {
+ List<MessageRecordVO> messages = provider.queryMessages(
+ "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+ assertThat(messages).extracting(MessageRecordVO::getMsgId)
+ .containsExactly("newer", "older");
+ }
+ }
+
@Test
void toRecordVOBoundsMessageBodyAndProperties() {
MessageExt message = new MessageExt();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java
index 1fbcc5a1..c3a89867 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/credential/CloudCredentialServiceTest.java
@@ -182,6 +182,7 @@ class CloudCredentialServiceTest {
stored.setVendor(InstanceVendor.ALIYUN);
when(credentialRepository.findById("cred-1")).thenReturn(Optional.of(stored));
when(instanceRepository.existsByCredentialId("cred-1")).thenReturn(false);
+ when(credentialRepository.deleteById("cred-1")).thenReturn(true);
service.delete("cred-1");
@@ -191,6 +192,22 @@ class CloudCredentialServiceTest {
eq("cred-1"), eq(null), eq("name=null, vendor=ALIYUN"),
eq("SUCCESS"), eq(null));
}
+ @Test
+ void deleteShouldRejectConcurrentRemovalBeforeInvalidatingClientsTest() {
+ CloudCredentialVO stored = new CloudCredentialVO();
+ stored.setId("cred-1");
+ stored.setVendor(InstanceVendor.ALIYUN);
+
when(credentialRepository.findById("cred-1")).thenReturn(Optional.of(stored));
+
when(instanceRepository.existsByCredentialId("cred-1")).thenReturn(false);
+ when(credentialRepository.deleteById("cred-1")).thenReturn(false);
+
+ assertThatThrownBy(() -> service.delete("cred-1"))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Cloud credential not found: cred-1");
+
+ verify(aliyunClientFactory, never()).invalidateCredential(any());
+ }
+
@Test
void revealShouldReturnUnmaskedCredentialTest() {
CloudCredentialVO stored = new CloudCredentialVO();
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepositoryTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepositoryTest.java
index c313a832..8d45a7e9 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepositoryTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/credential/MybatisPlusCloudCredentialRepositoryTest.java
@@ -31,6 +31,7 @@ import java.util.Optional;
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.when;
@ExtendWith(MockitoExtension.class)
@@ -42,6 +43,20 @@ class MybatisPlusCloudCredentialRepositoryTest {
@InjectMocks
private MybatisPlusCloudCredentialRepository repository;
+ @Test
+ void saveShouldReportALostConcurrentUpdate() {
+ CloudCredentialVO credential = new CloudCredentialVO();
+ credential.setId("cred-1");
+ credential.setVendor(InstanceVendor.ALIYUN);
+
when(credentialMapper.selectById("cred-1")).thenReturn(entity("cred-1",
"ALIYUN"));
+
when(credentialMapper.updateById(any(RmqCloudCredential.class))).thenReturn(0);
+
+ assertThatThrownBy(() -> repository.save(credential))
+ .isInstanceOf(BusinessException.class)
+ .hasMessage("Cloud credential update was not applied: cred-1")
+ .satisfies(error -> assertThat(((BusinessException)
error).getCode()).isEqualTo(409));
+ }
+
@Test
void findByIdShouldMapValidPersistedVendor() {
when(credentialMapper.selectById("cred-valid")).thenReturn(entity("cred-valid",
"ALIYUN"));