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 22dbbbfd3 fix(alert): consolidate locale normalization and alerting
robustness fixes (#2802)
22dbbbfd3 is described below
commit 22dbbbfd3be022a1bcbbc4e31fc28be72baf94de
Author: yyqdbngt <[email protected]>
AuthorDate: Mon Aug 31 19:23:48 2026 +0800
fix(alert): consolidate locale normalization and alerting robustness fixes
(#2802)
* fix(alert): normalize aggregation with root locale
* fix(alert): normalize delivery filters with root locale
* fix(alert): validate channels with root locale
* fix(nameserver): normalize hosts with root locale
* fix(alert): stabilize alert pagination order
* fix(alert): make fingerprint input unambiguous
* fix(alert): isolate rule test collector failures
* fix(alert): render notification templates in one pass
* fix(alert): reject null rules in imports
* fix(alert): scan all suppression candidates
* fix(alert): reject invalid reminder intervals
---------
Co-authored-by: Yue Wang <[email protected]>
---
.../cluster/nameserver/NamesrvAddrParser.java | 5 ++--
.../studio/ops/alert/AlertFingerprint.java | 15 ++++++++--
.../alert/AlertNotificationSuppressionService.java | 31 +++++++++++++++------
.../ops/alert/AlertNotificationTemplate.java | 20 +++++++++++---
.../studio/ops/alert/AlertRuleTransferDTO.java | 2 +-
.../studio/ops/alert/AlertStateMachine.java | 3 ++
.../ops/alert/MybatisPlusAlertRepository.java | 2 +-
.../studio/ops/alert/NativeAlertProcessor.java | 3 +-
.../studio/ops/alert/NativeAlertRulePolicy.java | 4 ++-
.../ops/alert/NativeAlertRuleTestService.java | 26 +++++++++++++++---
.../ops/alert/NotificationOutboxService.java | 7 +++--
.../cluster/nameserver/NamesrvAddrParserTest.java | 14 ++++++++++
.../studio/ops/alert/AlertFingerprintTest.java | 9 ++++++
.../AlertNotificationSuppressionServiceTest.java | 26 ++++++++++++++++++
.../ops/alert/AlertNotificationTemplateTest.java | 11 ++++++++
.../studio/ops/alert/AlertRuleTransferDTOTest.java | 32 ++++++++++++++++++++++
.../studio/ops/alert/AlertStateMachineTest.java | 12 ++++++++
.../ops/alert/MybatisPlusAlertRepositoryTest.java | 10 ++++++-
.../studio/ops/alert/NativeAlertProcessorTest.java | 31 +++++++++++++++++++++
.../ops/alert/NativeAlertRulePolicyTest.java | 14 ++++++++++
.../ops/alert/NativeAlertRuleTestServiceTest.java | 21 ++++++++++++++
.../ops/alert/NotificationOutboxServiceTest.java | 19 +++++++++++++
22 files changed, 287 insertions(+), 30 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParser.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParser.java
index 45bbca3b4..e0fa94de8 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParser.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParser.java
@@ -20,6 +20,7 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import java.util.ArrayList;
import java.util.List;
+import java.util.Locale;
/**
* Parses the NameServer address list stored in the registry. Accepts comma or
semicolon
@@ -62,7 +63,7 @@ public final class NamesrvAddrParser {
if (!isValidIpv6Literal(ipv6)) {
throw new BusinessException(400, "namesrvAddr segment has a
malformed IPv6 literal: " + segment);
}
- normalizedHost = "[" + ipv6.toLowerCase() + "]";
+ normalizedHost = "[" + ipv6.toLowerCase(Locale.ROOT) + "]";
} else {
if (host.isEmpty()) {
throw new BusinessException(400, "namesrvAddr segment is
missing a host: " + segment);
@@ -70,7 +71,7 @@ public final class NamesrvAddrParser {
if (host.chars().anyMatch(ch -> ch == ':' ||
Character.isWhitespace(ch))) {
throw new BusinessException(400, "namesrvAddr segment has an
unexpected character: " + segment);
}
- normalizedHost = host.toLowerCase();
+ normalizedHost = host.toLowerCase(Locale.ROOT);
}
int port;
try {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprint.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprint.java
index c8e0f05d9..56494e547 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprint.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprint.java
@@ -27,9 +27,9 @@ public final class AlertFingerprint {
}
public static String of(long ruleId, String instanceId, Map<String,
String> labels) {
- StringBuilder input = new
StringBuilder().append(ruleId).append('\n').append(instanceId).append('\n');
- new TreeMap<>(labels == null ? Map.of() : labels).forEach((key, value)
-> input.append(key)
- .append('=').append(value).append('\n'));
+ StringBuilder input = new
StringBuilder().append(ruleId).append('\n').append(escape(instanceId)).append('\n');
+ new TreeMap<>(labels == null ? Map.of() : labels).forEach((key, value)
-> input.append(escape(key))
+ .append('=').append(escape(value)).append('\n'));
try {
byte[] bytes =
MessageDigest.getInstance("SHA-256").digest(input.toString().getBytes(StandardCharsets.UTF_8));
StringBuilder fingerprint = new StringBuilder(bytes.length * 2);
@@ -41,4 +41,13 @@ public final class AlertFingerprint {
throw new IllegalStateException("SHA-256 is not available",
impossible);
}
}
+
+ private static String escape(String value) {
+ if (value == null) {
+ return "null";
+ }
+ return value.replace("\\", "\\\\")
+ .replace("\n", "\\n")
+ .replace("=", "\\=");
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionService.java
index 06297f926..727921623 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionService.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.ops.alert;
import lombok.RequiredArgsConstructor;
+import org.apache.rocketmq.studio.common.domain.PageResult;
import org.springframework.stereotype.Service;
import java.time.Duration;
@@ -32,6 +33,7 @@ import java.util.Optional;
@RequiredArgsConstructor
public class AlertNotificationSuppressionService {
private static final Duration CORRELATION_WINDOW = Duration.ofMinutes(30);
+ private static final int CANDIDATE_PAGE_SIZE = 100;
private final AlertRepository alertRepository;
@@ -40,16 +42,27 @@ public class AlertNotificationSuppressionService {
return Optional.empty();
}
LocalDateTime windowStart = event.getTime().minus(CORRELATION_WINDOW);
- List<SystemAlertVO> candidates = alertRepository.findAlertsPage(new
SystemAlertQuery(null, AlertDomain.CLUSTER,
- event.getInstanceId(), null, null, null, windowStart,
event.getTime(), 1, 100))
- .getItems().stream()
- .filter(candidate -> AlertCorrelationScope.matches(event,
candidate))
- .toList();
Map<String, SystemAlertVO> latestByIncident = new HashMap<>();
- for (SystemAlertVO candidate : candidates) {
- String incident = candidate.getFingerprint() == null ?
String.valueOf(candidate.getId())
- : candidate.getFingerprint();
- latestByIncident.merge(incident, candidate, (left, right) ->
later(left, right) ? left : right);
+ int page = 1;
+ long fetched = 0;
+ while (true) {
+ PageResult<SystemAlertVO> result =
alertRepository.findAlertsPage(new SystemAlertQuery(
+ null, AlertDomain.CLUSTER, event.getInstanceId(), null,
null, null,
+ windowStart, event.getTime(), page, CANDIDATE_PAGE_SIZE));
+ List<SystemAlertVO> candidates = result.getItems();
+ for (SystemAlertVO candidate : candidates) {
+ if (!AlertCorrelationScope.matches(event, candidate)) {
+ continue;
+ }
+ String incident = candidate.getFingerprint() == null ?
String.valueOf(candidate.getId())
+ : candidate.getFingerprint();
+ latestByIncident.merge(incident, candidate, (left, right) ->
later(left, right) ? left : right);
+ }
+ fetched += candidates.size();
+ if (candidates.isEmpty() || fetched >= result.getTotal()) {
+ break;
+ }
+ page++;
}
return latestByIncident.values().stream()
.filter(candidate ->
"FIRING".equalsIgnoreCase(candidate.getTransition()))
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplate.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplate.java
index 1db1c9d18..4c8d63247 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplate.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplate.java
@@ -9,6 +9,8 @@ package org.apache.rocketmq.studio.ops.alert;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.Set;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
/** Renders the fixed, documented placeholders allowed in an alert
notification. */
final class AlertNotificationTemplate {
@@ -17,16 +19,26 @@ final class AlertNotificationTemplate {
"broker.disk.usage_ratio",
"broker.jvm.heap.usage_ratio",
"broker.send_queue.usage_ratio");
+ private static final Pattern PLACEHOLDER =
Pattern.compile("\\$\\{([A-Za-z][A-Za-z0-9]*)}");
private AlertNotificationTemplate() {
}
static String render(String template, SystemAlertVO alert, AlertRuleVO
rule) {
- String result = hasText(template) ? template.trim() : DEFAULT_TEMPLATE;
- for (Map.Entry<String, String> entry : values(alert, rule).entrySet())
{
- result = result.replace("${" + entry.getKey() + "}",
entry.getValue());
+ String source = hasText(template) ? template.trim() : DEFAULT_TEMPLATE;
+ Map<String, String> replacements = values(alert, rule);
+ Matcher matcher = PLACEHOLDER.matcher(source);
+ StringBuilder result = new StringBuilder(source.length());
+ while (matcher.find()) {
+ String replacement = replacements.get(matcher.group(1));
+ if (replacement == null) {
+ matcher.appendReplacement(result,
Matcher.quoteReplacement(matcher.group()));
+ } else {
+ matcher.appendReplacement(result,
Matcher.quoteReplacement(replacement));
+ }
}
- return result;
+ matcher.appendTail(result);
+ return result.toString();
}
private static Map<String, String> values(SystemAlertVO alert, AlertRuleVO
rule) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleTransferDTO.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleTransferDTO.java
index 0faa53229..3c65450e5 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleTransferDTO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleTransferDTO.java
@@ -25,5 +25,5 @@ public class AlertRuleTransferDTO {
private AlertDomain domain;
@NotEmpty(message = "rules must not be empty")
- private List<@Valid AlertRuleRequestDTO> rules;
+ private List<@NotNull(message = "rule must not be null") @Valid
AlertRuleRequestDTO> rules;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachine.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachine.java
index d8f3ac261..a4ad996f7 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachine.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachine.java
@@ -43,6 +43,9 @@ public class AlertStateMachine {
if (requiredDuration == null || requiredDuration.isNegative()) {
throw new IllegalArgumentException("requiredDuration must not be
negative");
}
+ if (reminderInterval == null || reminderInterval.isNegative()) {
+ throw new IllegalArgumentException("reminderInterval must not be
negative");
+ }
AlertRuleState state = previous == null ? AlertRuleState.initial() :
previous;
if (!evaluation.matches()) {
return new AlertStateUpdate(state, AlertStateTransition.NONE);
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 11c34248e..6b77d01b2 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
@@ -195,7 +195,7 @@ public class MybatisPlusAlertRepository implements
AlertRepository {
"JSON_CONTAINS(labels_json, JSON_OBJECT({0}, {1}))",
query.labelKey(), query.labelValue())
.ge(query.from() != null, "time", query.from())
.le(query.to() != null, "time", query.to())
- .orderByDesc("time");
+ .orderByDesc("time", "id");
Page<RmqSystemAlert> result = alertMapper.selectPage(new
Page<>(query.page(), query.pageSize()), conditions);
return
PageResult.of(result.getRecords().stream().map(MybatisPlusAlertRepository::toAlertVO).toList(),
result.getTotal(), query.page(), query.pageSize());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessor.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessor.java
index 2d5b31103..b4559a533 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessor.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessor.java
@@ -28,6 +28,7 @@ import java.time.LocalDateTime;
import java.time.Duration;
import java.time.ZoneOffset;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Optional;
import java.util.TreeMap;
@@ -122,7 +123,7 @@ public class NativeAlertProcessor {
if (window.isEmpty()) {
return sample;
}
- double value = switch (rule.getAggregation() == null ? "LAST" :
rule.getAggregation().toUpperCase()) {
+ double value = switch (rule.getAggregation() == null ? "LAST" :
rule.getAggregation().toUpperCase(Locale.ROOT)) {
case "MAX" -> window.stream().mapToDouble(item ->
item.value()).max().orElse(sample.value());
case "MIN" -> window.stream().mapToDouble(item ->
item.value()).min().orElse(sample.value());
case "AVG" -> window.stream().mapToDouble(item ->
item.value()).average().orElse(sample.value());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicy.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicy.java
index cce2e7c58..18d75d261 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicy.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicy.java
@@ -19,6 +19,7 @@ package org.apache.rocketmq.studio.ops.alert;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.springframework.util.StringUtils;
+import java.util.Locale;
import java.util.Map;
import java.util.Set;
@@ -83,7 +84,8 @@ final class NativeAlertRulePolicy {
return;
}
for (String channel : rule.getChannels()) {
- if (!StringUtils.hasText(channel) ||
!NOTIFICATION_CHANNELS.contains(channel.trim().toLowerCase())) {
+ if (!StringUtils.hasText(channel)
+ ||
!NOTIFICATION_CHANNELS.contains(channel.trim().toLowerCase(Locale.ROOT))) {
throw new BusinessException(400, "Unsupported notification
channel: " + channel);
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestService.java
index 786a06da4..c4cc20933 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestService.java
@@ -17,6 +17,7 @@
package org.apache.rocketmq.studio.ops.alert;
import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.studio.cluster.metrics.BusinessMetricsCollector;
import org.apache.rocketmq.studio.cluster.metrics.ClusterMetricsCollector;
import org.apache.rocketmq.studio.cluster.metrics.MetricSample;
@@ -31,6 +32,7 @@ import java.util.List;
/** Executes a native rule once without persisting a snapshot, state, or
event. */
@Service
@RequiredArgsConstructor
+@Slf4j
public class NativeAlertRuleTestService {
private final InstanceRepository instanceRepository;
private final List<ClusterMetricsCollector> clusterCollectors;
@@ -43,11 +45,27 @@ public class NativeAlertRuleTestService {
.orElseThrow(() -> new BusinessException(404, "Instance not
found: " + rule.getInstanceId()));
List<MetricSample> samples = new ArrayList<>();
if (rule.getDomain() == AlertDomain.CLUSTER) {
- clusterCollectors.stream().filter(collector ->
collector.supports(instance))
- .forEach(collector ->
samples.addAll(collector.collect(instance)));
+ for (ClusterMetricsCollector collector : clusterCollectors) {
+ try {
+ if (collector.supports(instance)) {
+ samples.addAll(collector.collect(instance));
+ }
+ } catch (RuntimeException error) {
+ log.warn("Native cluster metric test collector failed for
instance {}: {}",
+ instance.getName(), error.getMessage());
+ }
+ }
} else {
- businessCollectors.stream().filter(collector ->
collector.supports(instance))
- .forEach(collector ->
samples.addAll(collector.collect(instance)));
+ for (BusinessMetricsCollector collector : businessCollectors) {
+ try {
+ if (collector.supports(instance)) {
+ samples.addAll(collector.collect(instance));
+ }
+ } catch (RuntimeException error) {
+ log.warn("Native business metric test collector failed for
instance {}: {}",
+ instance.getName(), error.getMessage());
+ }
+ }
}
return AlertRuleTestResultVO.builder().samples(samples.stream()
.filter(sample -> rule.getMetric().equals(sample.metricKey()))
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
index 26acbe785..8a76f8e08 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxService.java
@@ -35,6 +35,7 @@ import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.util.LinkedHashSet;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Set;
import java.util.ArrayList;
@@ -131,7 +132,7 @@ public class NotificationOutboxService {
Set<String> channels = new LinkedHashSet<>();
if (rule.getChannels() != null) {
rule.getChannels().stream().filter(StringUtils::hasText)
- .map(value ->
value.trim().toLowerCase()).forEach(channels::add);
+ .map(value ->
value.trim().toLowerCase(Locale.ROOT)).forEach(channels::add);
}
for (String channel : channels) {
if (!"dingtalk".equals(channel) && !"sms".equals(channel) &&
!"email".equals(channel)) {
@@ -218,7 +219,7 @@ public class NotificationOutboxService {
private static String normalizeFilter(String value) {
String normalized = normalizeTrim(value);
- return normalized == null ? null : normalized.toLowerCase();
+ return normalized == null ? null : normalized.toLowerCase(Locale.ROOT);
}
private static String normalizeTrim(String value) {
@@ -231,7 +232,7 @@ public class NotificationOutboxService {
return null;
}
try {
- return
NotificationOutboxStatus.valueOf(value.toUpperCase()).name();
+ return
NotificationOutboxStatus.valueOf(value.toUpperCase(Locale.ROOT)).name();
} catch (IllegalArgumentException error) {
throw new
org.apache.rocketmq.studio.common.exception.BusinessException(400,
"Unknown notification delivery status: " + status);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParserTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParserTest.java
index 71f1dd6fa..2b30bb5fc 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParserTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/nameserver/NamesrvAddrParserTest.java
@@ -19,6 +19,8 @@ package org.apache.rocketmq.studio.cluster.nameserver;
import org.apache.rocketmq.studio.common.exception.BusinessException;
import org.junit.jupiter.api.Test;
+import java.util.Locale;
+
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -34,6 +36,18 @@ class NamesrvAddrParserTest {
assertThat(NamesrvAddrParser.normalize("NS1.Example.COM:9876")).isEqualTo("ns1.example.com:9876");
}
+ @Test
+ void lowercasesHostsIndependentlyOfTheDefaultLocaleTest() {
+ Locale previous = Locale.getDefault();
+ try {
+ Locale.setDefault(Locale.forLanguageTag("tr-TR"));
+
assertThat(NamesrvAddrParser.normalize("INTERNAL.Example.COM:9876"))
+ .isEqualTo("internal.example.com:9876");
+ } finally {
+ Locale.setDefault(previous);
+ }
+ }
+
@Test
void normalizesSeparatorsAndWhitespaceTest() {
assertThat(NamesrvAddrParser.normalize(" ns1:9876 ; ns2:9876 ,ns3:9876
"))
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprintTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprintTest.java
index 7eec5cb52..9956d2c98 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprintTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertFingerprintTest.java
@@ -37,4 +37,13 @@ class AlertFingerprintTest {
.isEqualTo(AlertFingerprint.of(7L, "local", second))
.hasSize(64);
}
+
+ @Test
+ void
separatorCharactersCannotCreateTheSameFingerprintForDifferentLabelsTest() {
+ Map<String, String> embeddedLabel = Map.of("a", "b\nc=d");
+ Map<String, String> separateLabels = Map.of("a", "b", "c", "d");
+
+ assertThat(AlertFingerprint.of(7L, "local", embeddedLabel))
+ .isNotEqualTo(AlertFingerprint.of(7L, "local",
separateLabels));
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionServiceTest.java
index 9a036fc4b..99e250f30 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationSuppressionServiceTest.java
@@ -22,10 +22,13 @@ import org.junit.jupiter.api.Test;
import java.time.LocalDateTime;
import java.util.List;
import java.util.Map;
+import java.util.Optional;
+import java.util.stream.IntStream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -87,6 +90,29 @@ class AlertNotificationSuppressionServiceTest {
.isEmpty();
}
+ @Test
+ void searchesBeyondTheFirstCandidatePageTest() {
+ AlertRepository repository = mock(AlertRepository.class);
+ LocalDateTime now = LocalDateTime.now();
+ List<SystemAlertVO> unrelated = IntStream.range(0, 100)
+ .mapToObj(index -> event((long) index, AlertDomain.CLUSTER,
"FIRING",
+ "broker-" + index, now.minusMinutes(1)))
+ .toList();
+ SystemAlertVO cause = event(101L, AlertDomain.CLUSTER, "FIRING",
"target-broker",
+ now.minusMinutes(2));
+ when(repository.findAlertsPage(any())).thenReturn(
+ PageResult.of(unrelated, 101, 1, 100),
+ PageResult.of(List.of(cause), 101, 2, 100));
+
+ Optional<SystemAlertVO> result = new
AlertNotificationSuppressionService(repository)
+ .findSuppressingClusterAlert(event(102L, AlertDomain.BUSINESS,
"FIRING",
+ "target-broker", now));
+
+ assertThat(result).contains(cause);
+ verify(repository, times(2)).findAlertsPage(any());
+
verify(repository).findAlertsPage(org.mockito.ArgumentMatchers.argThat(query ->
query.page() == 2));
+ }
+
private static SystemAlertVO event(Long id, AlertDomain domain, String
transition, String brokerName,
LocalDateTime time) {
return
SystemAlertVO.builder().id(id).domain(domain).transition(transition).instanceId("local")
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplateTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplateTest.java
index 61eb8abdf..08897ad8f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplateTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertNotificationTemplateTest.java
@@ -40,4 +40,15 @@ class AlertNotificationTemplateTest {
assertThat(AlertNotificationTemplate.render(null, alert, null))
.isEqualTo("[info] Test - connection works\nLabels: ");
}
+
+ @Test
+ void doesNotExpandPlaceholderSyntaxIntroducedByAlertValuesTest() {
+ SystemAlertVO alert = SystemAlertVO.builder()
+ .title("${description}")
+ .description("internal detail")
+ .build();
+
+ assertThat(AlertNotificationTemplate.render("${title}", alert, null))
+ .isEqualTo("${description}");
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleTransferDTOTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleTransferDTOTest.java
new file mode 100644
index 000000000..49b2ca3c0
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertRuleTransferDTOTest.java
@@ -0,0 +1,32 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0.
+ */
+package org.apache.rocketmq.studio.ops.alert;
+
+import jakarta.validation.Validation;
+import jakarta.validation.Validator;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class AlertRuleTransferDTOTest {
+
+ private final Validator validator =
Validation.buildDefaultValidatorFactory().getValidator();
+
+ @Test
+ void rejectsNullRuleEntriesBeforeImportTest() {
+ AlertRuleTransferDTO transfer = new AlertRuleTransferDTO();
+ transfer.setVersion(AlertRuleTransferDTO.VERSION);
+ transfer.setDomain(AlertDomain.CLUSTER);
+ transfer.setRules(Collections.singletonList(null));
+
+ assertThat(validator.validate(transfer))
+ .extracting(violation -> violation.getMessage())
+ .contains("rule must not be null");
+ }
+}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachineTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachineTest.java
index 36bc5ef3e..baacb9565 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachineTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/AlertStateMachineTest.java
@@ -23,6 +23,7 @@ import java.time.Instant;
import java.time.Duration;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
class AlertStateMachineTest {
private final AlertStateMachine stateMachine = new AlertStateMachine();
@@ -103,6 +104,17 @@ class AlertStateMachineTest {
assertThat(noReminder.transition()).isEqualTo(AlertStateTransition.NONE);
}
+ @Test
+ void rejectsInvalidReminderIntervalsTest() {
+ assertThatThrownBy(() -> stateMachine.advance(null, met(), 1,
Duration.ZERO, null, now))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessage("reminderInterval must not be negative");
+ assertThatThrownBy(() -> stateMachine.advance(null, met(), 1,
Duration.ZERO,
+ Duration.ofSeconds(-1), now))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessage("reminderInterval must not be negative");
+ }
+
private static AlertEvaluationResult met() {
return new AlertEvaluationResult(true, true, 0.9,
MetricAvailability.AVAILABLE);
}
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 7e3844e65..852003124 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
@@ -291,7 +291,8 @@ class MybatisPlusAlertRepositoryTest {
repository.findAlertsPage(new SystemAlertQuery(null, null, null, null,
null, null, null, null, 1, 20));
- verify(alertMapper).selectPage(any(Page.class), any());
+ verify(alertMapper).selectPage(any(Page.class),
+
argThat(MybatisPlusAlertRepositoryTest::hasStableAlertOrdering));
}
private static boolean hasScopeLabelAndTimeFilters(Wrapper<RmqSystemAlert>
query) {
@@ -314,6 +315,13 @@ class MybatisPlusAlertRepositoryTest {
return queryWrapper.getParamNameValuePairs().containsValue("info");
}
+ private static boolean hasStableAlertOrdering(Wrapper<RmqSystemAlert>
query) {
+ if (!(query instanceof QueryWrapper<?> queryWrapper)) {
+ return false;
+ }
+ return queryWrapper.getCustomSqlSegment().contains("ORDER BY time
DESC,id DESC");
+ }
+
private static boolean hasBusinessRulePageFilters(Wrapper<RmqAlertRule>
query) {
if (!(query instanceof QueryWrapper<?> queryWrapper)) {
return false;
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessorTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessorTest.java
index 56f3c65cb..f93c82ac9 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessorTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertProcessorTest.java
@@ -24,6 +24,7 @@ import org.junit.jupiter.api.Test;
import java.time.Instant;
import java.util.HashMap;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Optional;
@@ -204,6 +205,36 @@ class NativeAlertProcessorTest {
});
}
+ @Test
+ void evaluatesAggregationIndependentlyOfTheDefaultLocaleTest() {
+ Locale previous = Locale.getDefault();
+ try {
+ Locale.setDefault(Locale.forLanguageTag("tr-TR"));
+ AlertService service = mock(AlertService.class);
+ AlertRuleVO rule =
AlertRuleVO.builder().id(1L).domain(AlertDomain.BUSINESS).name("Orders lag")
+
.metric("consumer.lag.total").operator(">").threshold(5).enabled(true).instanceId("local")
+
.consumerGroup("orders").aggregation("min").windowSeconds(300).build();
+
when(service.listRules(AlertDomain.BUSINESS)).thenReturn(List.of(rule));
+ MetricSnapshotRepository snapshots =
mock(MetricSnapshotRepository.class);
+ MetricSample current = sample("orders", 20D);
+ when(snapshots.findRecent(any(MetricSample.class),
any(Instant.class)))
+ .thenReturn(List.of(sample("orders", 10D),
sample("orders", 30D), current));
+ AlertStateRepository states = mock(AlertStateRepository.class);
+
when(states.find(any(AlertStateKey.class))).thenReturn(Optional.empty());
+ when(states.save(any(AlertStateKey.class),
any(AlertRuleState.class))).thenReturn(true);
+
+ new NativeAlertProcessor(service, new AlertRuleEvaluator(), new
AlertStateMachine(), states, snapshots,
+ mock(AlertRepository.class),
mock(NotificationOutboxService.class), suppression())
+ .process(List.of(current));
+
+ org.mockito.ArgumentCaptor<AlertRuleState> saved =
org.mockito.ArgumentCaptor.forClass(AlertRuleState.class);
+ verify(states).save(any(AlertStateKey.class), saved.capture());
+ assertThat(saved.getValue().currentValue()).isEqualTo(10D);
+ } finally {
+ Locale.setDefault(previous);
+ }
+ }
+
@Test
void recordsTheRuleTriggerTimeWhenEmittingAFiringEventTest() {
AlertService service = mock(AlertService.class);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicyTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicyTest.java
index a00980ae8..88e0c3c79 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicyTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRulePolicyTest.java
@@ -20,6 +20,7 @@ import
org.apache.rocketmq.studio.common.exception.BusinessException;
import org.junit.jupiter.api.Test;
import java.util.List;
+import java.util.Locale;
import static org.assertj.core.api.Assertions.assertThatCode;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -93,6 +94,19 @@ class NativeAlertRulePolicyTest {
.isInstanceOf(BusinessException.class).hasMessageContaining("Unsupported
notification channel");
}
+ @Test
+ void acceptsNotificationChannelsIndependentlyOfTheDefaultLocaleTest() {
+ Locale previous = Locale.getDefault();
+ try {
+ Locale.setDefault(Locale.forLanguageTag("tr-TR"));
+ assertThatCode(() ->
NativeAlertRulePolicy.validate(rule(AlertDomain.BUSINESS,
+ "rocketmq_consumer_lag_messages").channels(List.of("
DINGTALK ")).build()))
+ .doesNotThrowAnyException();
+ } finally {
+ Locale.setDefault(previous);
+ }
+ }
+
private static AlertRuleVO.AlertRuleVOBuilder rule(AlertDomain domain,
String metric) {
return AlertRuleVO.builder().domain(domain).name("Test
rule").metric(metric).consecutiveSamples(1);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestServiceTest.java
index b42d2ea36..1edf22952 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NativeAlertRuleTestServiceTest.java
@@ -78,6 +78,27 @@ class NativeAlertRuleTestServiceTest {
});
}
+ @Test
+ void keepsResultsFromHealthyCollectorsWhenAnotherCollectorFailsTest() {
+ InstanceRepository instances = mock(InstanceRepository.class);
+ BusinessMetricsCollector failing =
mock(BusinessMetricsCollector.class);
+ BusinessMetricsCollector healthy =
mock(BusinessMetricsCollector.class);
+ InstanceVO instance = InstanceVO.builder().name("local").build();
+
when(instances.findByIdentifier("local")).thenReturn(Optional.of(instance));
+ when(failing.supports(instance)).thenReturn(true);
+ when(failing.collect(instance)).thenThrow(new
IllegalStateException("provider unavailable"));
+ when(healthy.supports(instance)).thenReturn(true);
+ when(healthy.collect(instance)).thenReturn(List.of(sample("orders",
20)));
+ AlertRuleVO rule =
AlertRuleVO.builder().domain(AlertDomain.BUSINESS).metric("consumer.lag.total")
+
.instanceId("local").consumerGroup("orders").operator(">").threshold(10).build();
+
+ AlertRuleTestResultVO result = new
NativeAlertRuleTestService(instances, List.of(),
+ List.of(failing, healthy), new
AlertRuleEvaluator()).test(rule);
+
+ assertThat(result.samples()).singleElement()
+ .satisfies(sample ->
assertThat(sample.currentValue()).isEqualTo(20));
+ }
+
private static MetricSample sample(String group, double value) {
return sample("consumer.lag.total", group, value);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxServiceTest.java
index f2818f6d1..87954e859 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxServiceTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/ops/alert/NotificationOutboxServiceTest.java
@@ -27,6 +27,7 @@ import java.time.Duration;
import java.util.TimeZone;
import java.util.Arrays;
import java.util.List;
+import java.util.Locale;
import java.util.Optional;
import static org.mockito.ArgumentMatchers.any;
@@ -340,6 +341,24 @@ class NotificationOutboxServiceTest {
assertThat(result.getItems()).containsExactly(delivery);
}
+ @Test
+ void normalizesDeliveryFiltersIndependentlyOfTheDefaultLocaleTest() {
+ Locale previous = Locale.getDefault();
+ try {
+ Locale.setDefault(Locale.forLanguageTag("tr-TR"));
+ RmqAlertNotificationOutboxMapper mapper =
mock(RmqAlertNotificationOutboxMapper.class);
+ when(mapper.countPage("dingtalk", "PENDING",
"Local")).thenReturn(0L);
+
+ new NotificationOutboxService(mapper,
mock(SettingsRepository.class), mock(AlertSilenceService.class),
+ mock(AlertRepository.class),
mock(OperationAuditService.class))
+ .listDeliveries(" DINGTALK ", "pending", "Local", 1, 20);
+
+ verify(mapper).countPage("dingtalk", "PENDING", "Local");
+ } finally {
+ Locale.setDefault(previous);
+ }
+ }
+
@Test
void retriesOnlyFailedDeliveryAndResetsItsDispatchStateTest() {
RmqAlertNotificationOutboxMapper mapper =
mock(RmqAlertNotificationOutboxMapper.class);