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 165a72e5f fix(alert): keep cloud availability incidents stable across
status changes (#5140)
165a72e5f is described below
commit 165a72e5fce41cc2adb2eaab90dbb3fa0503a8be
Author: coder999o <[email protected]>
AuthorDate: Thu Oct 1 18:24:53 2026 +0800
fix(alert): keep cloud availability incidents stable across status changes
(#5140)
---
docs/studio-native-alerting-design.md | 4 +-
.../studio/ops/alert/AlertFingerprint.java | 5 +-
.../studio/ops/alert/AlertFingerprintTest.java | 13 ++++
.../studio/ops/alert/NativeAlertProcessorTest.java | 70 ++++++++++++++++++++++
4 files changed, 90 insertions(+), 2 deletions(-)
diff --git a/docs/studio-native-alerting-design.md
b/docs/studio-native-alerting-design.md
index 7caf50d4f..7b1f75d53 100644
--- a/docs/studio-native-alerting-design.md
+++ b/docs/studio-native-alerting-design.md
@@ -198,7 +198,9 @@ FIRING -- user acknowledges --> ACKED
PENDING/FIRING -- silence matches --> state unchanged, notification suppressed
```
-The fingerprint is `sha256(ruleId + instanceId + sorted(labels))`. A single
rule therefore creates independent events for different Brokers, Topics,
queues, or consumer groups.
+The fingerprint is `sha256(ruleId + instanceId + sorted(identity labels))`.
Descriptive labels that change while the same incident persists — currently
`cloudStatus`, the cloud control-plane instance status behind
`cloud.instance.availability` — are excluded from identity but still carried on
the sample and on every emitted event, so one unavailable cloud instance keeps
a single incident across STOPPED → STARTING → RUNNING. A single rule therefore
still creates independent events for dif [...]
+
+Because existing `rmq_alert_state` rows carry fingerprints computed with
`cloudStatus` included, the first collection cycle after an upgrade emits one
RESOLVED + FIRING pair per active cloud-availability incident; the migration is
one-shot and self-healing.
Current delivery emits on `FIRING`, periodic `REMINDER` transitions controlled
by the rule's `reminderInterval`, and `RESOLVED`. A separate `cooldownSeconds`
policy remains future work. A value recovery always emits a `RESOLVED` event.
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 56494e547..c3e4da83d 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
@@ -28,7 +28,10 @@ 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(escape(instanceId)).append('\n');
- new TreeMap<>(labels == null ? Map.of() : labels).forEach((key, value)
-> input.append(escape(key))
+ Map<String, String> identityLabels = new TreeMap<>(labels == null ?
Map.of() : labels);
+ // Cloud status describes an incident but changes while the same
instance remains unavailable.
+ identityLabels.remove("cloudStatus");
+ identityLabels.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));
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 9956d2c98..039088594 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
@@ -46,4 +46,17 @@ class AlertFingerprintTest {
assertThat(AlertFingerprint.of(7L, "local", embeddedLabel))
.isNotEqualTo(AlertFingerprint.of(7L, "local",
separateLabels));
}
+
+ @Test
+ void cloudStatusDoesNotChangeTheIdentityOfACloudInstanceTest() {
+ String stopped = AlertFingerprint.of(7L, "cloud-local",
+ Map.of("cloudInstanceId", "rmq-cloud", "cloudStatus",
"STOPPED"));
+
+ assertThat(stopped).isEqualTo(AlertFingerprint.of(7L, "cloud-local",
+ Map.of("cloudInstanceId", "rmq-cloud", "cloudStatus",
"STARTING")))
+ .isEqualTo(AlertFingerprint.of(7L, "cloud-local",
+ Map.of("cloudInstanceId", "rmq-cloud")));
+ assertThat(stopped).isNotEqualTo(AlertFingerprint.of(7L, "cloud-local",
+ Map.of("cloudInstanceId", "another-cloud")));
+ }
}
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 82c321e6e..86860162c 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
@@ -176,6 +176,76 @@ class NativeAlertProcessorTest {
.isEqualTo(AlertStateStatus.FIRING);
}
+ @Test
+ void cloudAvailabilityStatusChangesKeepOneIncidentTest() {
+ AlertService service = mock(AlertService.class);
+ AlertRuleVO rule =
AlertRuleVO.builder().id(7L).domain(AlertDomain.CLUSTER)
+ .name("Cloud instance
unavailable").metric("cloud.instance.availability")
+
.operator("UNAVAILABLE").enabled(true).instanceId("cloud-local")
+ .consecutiveSamples(1).build();
+ when(service.listRules(AlertDomain.CLUSTER)).thenReturn(List.of(rule));
+ Map<AlertStateKey, AlertRuleState> saved = new HashMap<>();
+ AlertStateRepository states = mock(AlertStateRepository.class);
+ when(states.find(any(AlertStateKey.class)))
+ .thenAnswer(invocation ->
Optional.ofNullable(saved.get(invocation.getArgument(0))));
+ when(states.save(any(AlertStateKey.class), any(AlertRuleState.class)))
+ .thenAnswer(invocation -> {
+ saved.put(invocation.getArgument(0),
invocation.getArgument(1));
+ return true;
+ });
+ AlertRepository alerts = mock(AlertRepository.class);
+ NativeAlertProcessor processor = processor(service, states, alerts);
+ Instant collectedAt = Instant.parse("2026-09-29T00:00:00Z");
+
+ processor.process(List.of(cloudAvailabilitySample("STOPPED",
collectedAt)));
+ processor.process(List.of(cloudAvailabilitySample("STARTING",
collectedAt.plusSeconds(60))));
+
+ assertThat(saved).hasSize(1);
+ verify(alerts, times(1)).saveAlert(any(SystemAlertVO.class));
+
+ processor.process(List.of(new
MetricSample("cloud.instance.availability", AlertDomain.CLUSTER,
+ "cloud-local", null, Map.of("cloudInstanceId", "rmq-cloud",
"cloudStatus", "RUNNING"),
+ 1D, MetricAvailability.AVAILABLE,
collectedAt.plusSeconds(120))));
+
+
assertThat(saved.values()).singleElement().extracting(AlertRuleState::status)
+ .isEqualTo(AlertStateStatus.RESOLVED);
+ verify(alerts, times(2)).saveAlert(any(SystemAlertVO.class));
+ }
+
+ @Test
+ void cloudStatusChangeDoesNotResolveTheExistingIncidentTest() {
+ AlertService service = mock(AlertService.class);
+ AlertRuleVO rule =
AlertRuleVO.builder().id(7L).domain(AlertDomain.CLUSTER)
+ .name("Cloud instance
unavailable").metric("cloud.instance.availability")
+
.operator("UNAVAILABLE").enabled(true).instanceId("cloud-local")
+ .consecutiveSamples(1).build();
+ when(service.listRules(AlertDomain.CLUSTER)).thenReturn(List.of(rule));
+ AlertStateKey key = new AlertStateKey(7L, AlertFingerprint.of(7L,
"cloud-local",
+ Map.of("cloudInstanceId", "rmq-cloud")));
+ Instant firedAt = Instant.parse("2026-09-29T00:00:00Z");
+ AlertRuleState firing = new AlertRuleState(AlertStateStatus.FIRING, 1,
null,
+ firedAt, firedAt, firedAt, null);
+ AlertStateRepository states = mock(AlertStateRepository.class);
+ when(states.findActive(any(MetricCollectionScope.class), any()))
+ .thenReturn(List.of(new ActiveAlertState(key, firing,
"cloud-local",
+ Map.of("cloudInstanceId", "rmq-cloud", "cloudStatus",
"STOPPED"))));
+
+ NativeAlertProcessor processor = new NativeAlertProcessor(service,
+ mock(NativeAlertEvaluationService.class), new
AlertStateMachine(), states,
+ mock(AlertRepository.class),
mock(NotificationOutboxService.class), suppression(), mockTxManager());
+ processor.processSuccessfulCollection(new
MetricCollectionScope(AlertDomain.CLUSTER, "cloud-local",
+ java.util.Set.of("cloud.instance.availability")),
+ List.of(cloudAvailabilitySample("STARTING",
firedAt.plusSeconds(60))));
+
+ verify(states, never()).save(eq(key), any(AlertRuleState.class));
+ }
+
+ private static MetricSample cloudAvailabilitySample(String status, Instant
collectedAt) {
+ return new MetricSample("cloud.instance.availability",
AlertDomain.CLUSTER, "cloud-local", null,
+ Map.of("cloudInstanceId", "rmq-cloud", "cloudStatus", status),
null,
+ MetricAvailability.UNAVAILABLE, collectedAt);
+ }
+
@Test
void loadsRulesOncePerDomainForABatchOfSamplesTest() {
AlertService service = mock(AlertService.class);