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);

Reply via email to