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 9d46e0dd feat: support HA replication lag alert rules in alert module 
(#658)
9d46e0dd is described below

commit 9d46e0dd0fd799a691dfff68efa8104a872ecf40
Author: zhaohai <[email protected]>
AuthorDate: Fri Jul 31 15:50:51 2026 +0800

    feat: support HA replication lag alert rules in alert module (#658)
    
    Rework of the original branch-compensate proposal: instead of adding a
    parallel per-metric CRUD stack, extend the existing alert module so
    replication lag rules reuse /api/alert-rules:
    
    - Add brokerName/clusterName scope and severity fields to AlertRuleVO
    - Render broker/cluster label selectors in exported Prometheus expressions
    - Use rule severity (critical/warning/info, default warning) in export
    - Route replication/slave metrics to the broker rule group
---
 .../rocketmq/studio/ops/alert/AlertRuleVO.java     |  6 ++
 .../rocketmq/studio/ops/alert/AlertService.java    | 35 ++++++++++-
 .../studio/ops/alert/AlertServiceTest.java         | 67 ++++++++++++++++++++++
 3 files changed, 106 insertions(+), 2 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java 
b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java
index ec467237..2b8557a2 100644
--- a/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java
+++ b/server/src/main/java/org/apache/rocketmq/studio/ops/alert/AlertRuleVO.java
@@ -39,4 +39,10 @@ public class AlertRuleVO {
     private boolean enabled;
     private String lastTriggered;
     private String description;
+    /** Target broker name pattern (e.g. "broker-a", or "*" for all) */
+    private String brokerName;
+    /** Target cluster name pattern (e.g. "DefaultCluster", or "*" for all) */
+    private String clusterName;
+    /** Alert severity: critical, warning, info */
+    private String severity;
 }
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 5fb04726..59df2845 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
@@ -168,7 +168,7 @@ public class AlertService {
                 alertName(rule),
                 expression(rule),
                 duration(rule),
-                "warning",
+                severity(rule),
                 team,
                 summary(rule),
                 description(rule));
@@ -194,7 +194,35 @@ public class AlertService {
     private String expression(AlertRuleVO rule) {
         String metric = hasText(rule.getMetric()) ? rule.getMetric() : 
"rocketmq_consumer_lag_messages";
         String operator = hasText(rule.getOperator()) ? rule.getOperator() : 
">";
-        return metric + " " + operator + " " + 
formatThreshold(rule.getThreshold());
+        return metric + labelSelector(rule) + " " + operator + " " + 
formatThreshold(rule.getThreshold());
+    }
+
+    private String labelSelector(AlertRuleVO rule) {
+        StringBuilder selector = new StringBuilder();
+        appendLabel(selector, "cluster", rule.getClusterName());
+        appendLabel(selector, "broker", rule.getBrokerName());
+        return selector.isEmpty() ? "" : "{" + selector + "}";
+    }
+
+    private void appendLabel(StringBuilder selector, String label, String 
value) {
+        if (!hasText(value) || "*".equals(value.trim())) {
+            return;
+        }
+        if (!selector.isEmpty()) {
+            selector.append(',');
+        }
+        selector.append(label).append("=\"").append(value.trim()).append('"');
+    }
+
+    private String severity(AlertRuleVO rule) {
+        String severity = rule.getSeverity();
+        if (hasText(severity)) {
+            String normalized = severity.trim().toLowerCase();
+            if ("critical".equals(normalized) || "warning".equals(normalized) 
|| "info".equals(normalized)) {
+                return normalized;
+            }
+        }
+        return "warning";
     }
 
     private String formatThreshold(double threshold) {
@@ -212,6 +240,9 @@ public class AlertService {
         if (!hasText(metric)) {
             return "broker";
         }
+        if (metric.contains("replication") || metric.contains("fall_behind") 
|| metric.contains("slave")) {
+            return "broker";
+        }
         if (metric.contains("consumer") || metric.contains("lag")) {
             return "consumer";
         }
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 318d34f2..4c110c51 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
@@ -108,6 +108,73 @@ class AlertServiceTest {
                 .contains("description: \"Lag too high\"");
     }
 
+    @Test
+    void 
exportPrometheusRulesYamlShouldRenderReplicationLagRuleWithScopeAndSeverity() {
+        AlertRuleVO rule = AlertRuleVO.builder()
+                .name("Replication Lag High")
+                .metric("rocketmq_broker_replication_lag_bytes")
+                .operator(">")
+                .threshold(104857600)
+                .duration("5m")
+                .brokerName("broker-a")
+                .clusterName("DefaultCluster")
+                .severity("critical")
+                .description("Slave falls behind master")
+                .build();
+        when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+
+        String result = alertService.exportPrometheusRulesYaml();
+
+        assertThat(result)
+                .contains("rocketmq-broker.rules")
+                .contains("expr: 
rocketmq_broker_replication_lag_bytes{cluster=\"DefaultCluster\",broker=\"broker-a\"}
 > 104857600")
+                .contains("for: 5m")
+                .contains("severity: critical");
+    }
+
+    @Test
+    void 
exportPrometheusRulesYamlShouldIgnoreWildcardScopeAndInvalidSeverity() {
+        AlertRuleVO rule = AlertRuleVO.builder()
+                .name("Replication Lag Any Broker")
+                .metric("rocketmq_broker_replication_lag_bytes")
+                .operator(">")
+                .threshold(1024)
+                .brokerName("*")
+                .clusterName(" ")
+                .severity("fatal")
+                .build();
+        when(alertRepository.findAllRules()).thenReturn(List.of(rule));
+
+        String result = alertService.exportPrometheusRulesYaml();
+
+        assertThat(result)
+                .contains("expr: rocketmq_broker_replication_lag_bytes > 1024")
+                .contains("severity: warning");
+    }
+
+    @Test
+    void createRuleShouldPreserveReplicationScopeFields() {
+        AlertRuleVO input = AlertRuleVO.builder()
+                .name("Replication Lag High")
+                .metric("rocketmq_broker_replication_lag_bytes")
+                .operator(">")
+                .threshold(104857600)
+                .thresholdUnit("B")
+                .brokerName("broker-a")
+                .clusterName("DefaultCluster")
+                .severity("critical")
+                .build();
+        
when(alertRepository.saveRule(any(AlertRuleVO.class))).thenAnswer(invocation -> 
invocation.getArgument(0));
+
+        AlertRuleVO result = alertService.createRule(input);
+
+        assertThat(result.getId()).isNotNull().isNotEmpty();
+        assertThat(result.getBrokerName()).isEqualTo("broker-a");
+        assertThat(result.getClusterName()).isEqualTo("DefaultCluster");
+        assertThat(result.getSeverity()).isEqualTo("critical");
+        verify(alertRepository).saveRule(result);
+    }
+
     @Test
     void createRuleShouldAssignId() {
         AlertRuleVO input = AlertRuleVO.builder().name("New 
Rule").metric("tps")

Reply via email to