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 bd346021 feat: validate metrics and Prometheus queries (#637)
bd346021 is described below
commit bd346021c2f60d94c1b3276627bdf1b821024b61
Author: aias00 <[email protected]>
AuthorDate: Tue Jul 28 07:14:51 2026 -0700
feat: validate metrics and Prometheus queries (#637)
* [Studio] Validate metrics and Prometheus queries
* feat: add RocketMQ metric profiles
* fix: apply data source test authentication
* chore: keep metrics PR scoped to metrics
---
.../{MetricsService.java => MetricProfile.java} | 39 +++++---
.../cluster/metrics/MetricProfileService.java | 106 ++++++++++++++++++++
.../{MetricsService.java => MetricProfileVO.java} | 37 ++++---
.../studio/cluster/metrics/MetricsController.java | 13 +++
.../studio/cluster/metrics/MetricsService.java | 75 ++++++++++++++
.../cluster/metrics/PrometheusMetricsSource.java | 12 +++
.../studio/cluster/metrics/SemanticMetric.java | 49 +++++++++
.../cluster/metrics/MetricProfileServiceTest.java | 110 +++++++++++++++++++++
.../cluster/metrics/MetricsControllerTest.java | 72 ++++++++++++++
.../studio/cluster/metrics/MetricsServiceTest.java | 73 ++++++++++++++
.../metrics/PrometheusMetricsSourceTest.java | 64 ++++++++++++
web/src/api/metrics.test.ts | 19 +++-
12 files changed, 643 insertions(+), 26 deletions(-)
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/MetricProfile.java
similarity index 50%
copy from
server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
copy to
server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfile.java
index c0d6ed45..18749dde 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/MetricProfile.java
@@ -16,20 +16,35 @@
*/
package org.apache.rocketmq.studio.cluster.metrics;
-import lombok.RequiredArgsConstructor;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.stereotype.Service;
+public enum MetricProfile {
+ ROCKETMQ_4_EXPORTER(
+ "rocketmq4-exporter",
+ "RocketMQ 4.x Exporter",
+ "RocketMQ 4.x clusters scraped through the standalone
rocketmq-exporter"),
+ ROCKETMQ_5_NATIVE(
+ "rocketmq5-native",
+ "RocketMQ 5.x Native",
+ "RocketMQ 5.1+ native OpenTelemetry and Prometheus metrics");
-@Slf4j
-@Service
-@RequiredArgsConstructor
-public class MetricsService {
+ private final String id;
+ private final String displayName;
+ private final String description;
- private final MetricsSource metricsSource;
+ MetricProfile(String id, String displayName, String description) {
+ this.id = id;
+ this.displayName = displayName;
+ this.description = description;
+ }
+
+ public String getId() {
+ return id;
+ }
+
+ public String getDisplayName() {
+ return displayName;
+ }
- public MetricDataVO query(MetricQueryDTO query) {
- log.debug("Querying metrics: start={}, end={}, step={}",
- query.getStart(), query.getEnd(), query.getStep());
- return metricsSource.query(query);
+ public String getDescription() {
+ return description;
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileService.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileService.java
new file mode 100644
index 00000000..ba7f3146
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileService.java
@@ -0,0 +1,106 @@
+/*
+ * 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.cluster.metrics;
+
+import org.springframework.stereotype.Service;
+
+import java.util.List;
+
+@Service
+public class MetricProfileService {
+
+ public List<MetricProfileVO> listProfiles() {
+ return List.of(
+ profile(MetricProfile.ROCKETMQ_4_EXPORTER,
rocketmq4ExporterMetrics()),
+ profile(MetricProfile.ROCKETMQ_5_NATIVE,
rocketmq5NativeMetrics())
+ );
+ }
+
+ private MetricProfileVO profile(MetricProfile profile,
+ List<MetricProfileVO.MetricMappingVO>
metrics) {
+ return MetricProfileVO.builder()
+ .id(profile.getId())
+ .name(profile.getDisplayName())
+ .description(profile.getDescription())
+ .metrics(metrics)
+ .build();
+ }
+
+ private List<MetricProfileVO.MetricMappingVO> rocketmq4ExporterMetrics() {
+ return List.of(
+ mapping(SemanticMetric.MESSAGE_IN_TPS, "rocketmq_broker_tps",
+ "sum(rocketmq_broker_tps) by (cluster, broker)",
+ "cluster", "broker"),
+ mapping(SemanticMetric.MESSAGE_OUT_TPS,
"rocketmq_consumer_tps",
+ "sum(rocketmq_consumer_tps) by (cluster, group,
topic)",
+ "cluster", "group", "topic"),
+ mapping(SemanticMetric.THROUGHPUT_IN,
"rocketmq_producer_message_size",
+ "sum(rocketmq_producer_message_size) by (cluster,
topic)",
+ "cluster", "topic"),
+ mapping(SemanticMetric.THROUGHPUT_OUT,
"rocketmq_consumer_message_size",
+ "sum(rocketmq_consumer_message_size) by (cluster,
group, topic)",
+ "cluster", "group", "topic"),
+ mapping(SemanticMetric.CONSUMER_LAG_MESSAGES,
"rocketmq_message_accumulation",
+ "sum(rocketmq_message_accumulation) by (cluster,
group, topic)",
+ "cluster", "group", "topic"),
+ mapping(SemanticMetric.CONSUMER_LAG_LATENCY,
"rocketmq_group_get_latency_by_storetime",
+ "max(rocketmq_group_get_latency_by_storetime) by
(cluster, group, topic)",
+ "cluster", "group", "topic"),
+ mapping(SemanticMetric.BROKER_HEALTH, "up",
+ "min(up{job=~\".*rocketmq.*\"}) by (job, instance)",
+ "job", "instance")
+ );
+ }
+
+ private List<MetricProfileVO.MetricMappingVO> rocketmq5NativeMetrics() {
+ return List.of(
+ mapping(SemanticMetric.MESSAGE_IN_TPS,
"rocketmq_messages_in_total",
+ "sum(rate(rocketmq_messages_in_total[1m])) by
(cluster, node_id)",
+ "cluster", "node_id", "topic", "message_type"),
+ mapping(SemanticMetric.MESSAGE_OUT_TPS,
"rocketmq_messages_out_total",
+ "sum(rate(rocketmq_messages_out_total[1m])) by
(cluster, node_id, consumer_group)",
+ "cluster", "node_id", "topic", "consumer_group"),
+ mapping(SemanticMetric.THROUGHPUT_IN,
"rocketmq_throughput_in_total",
+ "sum(rate(rocketmq_throughput_in_total[1m])) by
(cluster, node_id)",
+ "cluster", "node_id", "topic", "message_type"),
+ mapping(SemanticMetric.THROUGHPUT_OUT,
"rocketmq_throughput_out_total",
+ "sum(rate(rocketmq_throughput_out_total[1m])) by
(cluster, node_id, consumer_group)",
+ "cluster", "node_id", "topic", "consumer_group"),
+ mapping(SemanticMetric.CONSUMER_LAG_MESSAGES,
"rocketmq_consumer_lag_messages",
+ "sum(rocketmq_consumer_lag_messages) by (cluster,
topic, consumer_group)",
+ "cluster", "topic", "consumer_group"),
+ mapping(SemanticMetric.CONSUMER_LAG_LATENCY,
"rocketmq_consumer_lag_latency",
+ "max(rocketmq_consumer_lag_latency) by (cluster,
topic, consumer_group)",
+ "cluster", "topic", "consumer_group"),
+ mapping(SemanticMetric.BROKER_HEALTH,
"rocketmq_processor_watermark",
+ "max(rocketmq_processor_watermark) by (cluster,
node_id, processor)",
+ "cluster", "node_id", "processor")
+ );
+ }
+
+ private MetricProfileVO.MetricMappingVO mapping(SemanticMetric
semanticMetric, String prometheusMetric,
+ String promql, String...
labels) {
+ return MetricProfileVO.MetricMappingVO.builder()
+ .semanticMetric(semanticMetric.getKey())
+ .name(semanticMetric.getDisplayName())
+ .unit(semanticMetric.getUnit())
+ .prometheusMetric(prometheusMetric)
+ .promql(promql)
+ .labels(List.of(labels))
+ .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/MetricProfileVO.java
similarity index 56%
copy from
server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsService.java
copy to
server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileVO.java
index c0d6ed45..7c7f65d9 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/MetricProfileVO.java
@@ -16,20 +16,33 @@
*/
package org.apache.rocketmq.studio.cluster.metrics;
-import lombok.RequiredArgsConstructor;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.stereotype.Service;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
-@Slf4j
-@Service
-@RequiredArgsConstructor
-public class MetricsService {
+import java.util.List;
- private final MetricsSource metricsSource;
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class MetricProfileVO {
+ private String id;
+ private String name;
+ private String description;
+ private List<MetricMappingVO> metrics;
- public MetricDataVO query(MetricQueryDTO query) {
- log.debug("Querying metrics: start={}, end={}, step={}",
- query.getStart(), query.getEnd(), query.getStep());
- return metricsSource.query(query);
+ @Data
+ @Builder
+ @NoArgsConstructor
+ @AllArgsConstructor
+ public static class MetricMappingVO {
+ private String semanticMetric;
+ private String name;
+ private String unit;
+ private String prometheusMetric;
+ private String promql;
+ private List<String> labels;
}
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsController.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsController.java
index 9399ef53..02b6d855 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/MetricsController.java
@@ -22,17 +22,30 @@ import io.swagger.v3.oas.annotations.responses.ApiResponse;
import io.swagger.v3.oas.annotations.responses.ApiResponses;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
+import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
+import java.util.List;
+
@RestController
@RequestMapping("/api/metrics")
@RequiredArgsConstructor
public class MetricsController {
private final MetricsService metricsService;
+ private final MetricProfileService metricProfileService;
+
+ @Operation(summary = "List RocketMQ metric profiles",
+ description = "Returns version-aware semantic PromQL mappings for
RocketMQ metrics")
+ @ApiResponse(responseCode = "200", description = "Metric profiles listed
successfully",
+ useReturnTypeSchema = true)
+ @GetMapping("/profiles")
+ public Result<List<MetricProfileVO>> listProfiles() {
+ return Result.ok(metricProfileService.listProfiles());
+ }
@Operation(summary = "Query Prometheus range metrics",
description = "Executes a PromQL range query against the
configured Prometheus server")
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 c0d6ed45..a68eba17 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
@@ -18,18 +18,93 @@ package org.apache.rocketmq.studio.cluster.metrics;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
+import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Service;
+import org.springframework.util.StringUtils;
+
+import java.math.BigDecimal;
+import java.util.Map;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
@Slf4j
@Service
@RequiredArgsConstructor
public class MetricsService {
+ private static final long MAX_RANGE_SECONDS = 31L * 24 * 60 * 60;
+ private static final long MAX_SAMPLE_POINTS = 11_000L;
+ private static final Pattern NUMBER_PATTERN =
Pattern.compile("\\d+(?:\\.\\d+)?");
+ private static final Pattern DURATION_PART_PATTERN =
Pattern.compile("(\\d+(?:\\.\\d+)?)(ms|s|m|h|d|w|y)");
+ private static final Map<String, BigDecimal> UNIT_TO_MILLIS = Map.of(
+ "ms", BigDecimal.ONE,
+ "s", BigDecimal.valueOf(1_000L),
+ "m", BigDecimal.valueOf(60_000L),
+ "h", BigDecimal.valueOf(3_600_000L),
+ "d", BigDecimal.valueOf(86_400_000L),
+ "w", BigDecimal.valueOf(604_800_000L),
+ "y", BigDecimal.valueOf(31_536_000_000L)
+ );
private final MetricsSource metricsSource;
public MetricDataVO query(MetricQueryDTO query) {
+ validateQueryWindow(query);
log.debug("Querying metrics: start={}, end={}, step={}",
query.getStart(), query.getEnd(), query.getStep());
return metricsSource.query(query);
}
+
+ private void validateQueryWindow(MetricQueryDTO query) {
+ if (query == null) {
+ throw badRequest("Metric query is required");
+ }
+ long rangeSeconds = query.getEnd() - query.getStart();
+ if (rangeSeconds <= 0) {
+ throw badRequest("Metric query end must be later than start");
+ }
+ if (rangeSeconds > MAX_RANGE_SECONDS) {
+ throw badRequest("Metric query range must not exceed 31 days");
+ }
+ BigDecimal stepMillis = parseStepMillis(query.getStep());
+ if (stepMillis.signum() <= 0) {
+ throw badRequest("Metric query step must be positive");
+ }
+ BigDecimal samplePoints = BigDecimal.valueOf(rangeSeconds)
+ .multiply(BigDecimal.valueOf(1_000L))
+ .divideToIntegralValue(stepMillis)
+ .add(BigDecimal.ONE);
+ if (samplePoints.compareTo(BigDecimal.valueOf(MAX_SAMPLE_POINTS)) > 0)
{
+ throw badRequest("Metric query returns too many samples; increase
step or reduce range");
+ }
+ }
+
+ private BigDecimal parseStepMillis(String step) {
+ if (!StringUtils.hasText(step)) {
+ throw badRequest("Metric query step is required");
+ }
+ String value = step.strip();
+ if (NUMBER_PATTERN.matcher(value).matches()) {
+ return new BigDecimal(value).multiply(BigDecimal.valueOf(1_000L));
+ }
+
+ Matcher matcher = DURATION_PART_PATTERN.matcher(value);
+ BigDecimal millis = BigDecimal.ZERO;
+ int position = 0;
+ while (matcher.find()) {
+ if (matcher.start() != position) {
+ throw badRequest("Metric query step is invalid");
+ }
+ BigDecimal amount = new BigDecimal(matcher.group(1));
+ millis =
millis.add(amount.multiply(UNIT_TO_MILLIS.get(matcher.group(2))));
+ position = matcher.end();
+ }
+ if (position != value.length()) {
+ throw badRequest("Metric query step is invalid");
+ }
+ return millis;
+ }
+
+ private PrometheusException badRequest(String message) {
+ return new PrometheusException(HttpStatus.BAD_REQUEST.value(),
message);
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
index 44f7ced0..913cd56c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
@@ -108,10 +108,22 @@ public class PrometheusMetricsSource implements
MetricsSource {
if (query == null) {
throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
"Metric query is required");
}
+ if (!StringUtils.hasText(query.getMetric())) {
+ throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
"Metric query is required");
+ }
+ if (query.getStart() <= 0) {
+ throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
"Metric query start must be positive");
+ }
+ if (query.getEnd() <= 0) {
+ throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
"Metric query end must be positive");
+ }
if (query.getEnd() < query.getStart()) {
throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
"Metric query end must not be earlier than start");
}
+ if (!StringUtils.hasText(query.getStep())) {
+ throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
"Metric query step is required");
+ }
}
private URI queryRangeUri() {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/SemanticMetric.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/SemanticMetric.java
new file mode 100644
index 00000000..c3c66f24
--- /dev/null
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/SemanticMetric.java
@@ -0,0 +1,49 @@
+/*
+ * 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.cluster.metrics;
+
+public enum SemanticMetric {
+ MESSAGE_IN_TPS("message_in_tps", "Message In TPS", "messages/s"),
+ MESSAGE_OUT_TPS("message_out_tps", "Message Out TPS", "messages/s"),
+ THROUGHPUT_IN("throughput_in", "Throughput In", "bytes/s"),
+ THROUGHPUT_OUT("throughput_out", "Throughput Out", "bytes/s"),
+ CONSUMER_LAG_MESSAGES("consumer_lag_messages", "Consumer Lag", "messages"),
+ CONSUMER_LAG_LATENCY("consumer_lag_latency", "Consumer Lag Latency", "ms"),
+ BROKER_HEALTH("broker_health", "Broker Health", "up");
+
+ private final String key;
+ private final String displayName;
+ private final String unit;
+
+ SemanticMetric(String key, String displayName, String unit) {
+ this.key = key;
+ this.displayName = displayName;
+ this.unit = unit;
+ }
+
+ public String getKey() {
+ return key;
+ }
+
+ public String getDisplayName() {
+ return displayName;
+ }
+
+ public String getUnit() {
+ return unit;
+ }
+}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileServiceTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileServiceTest.java
new file mode 100644
index 00000000..3cd93129
--- /dev/null
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricProfileServiceTest.java
@@ -0,0 +1,110 @@
+/*
+ * 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.cluster.metrics;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+import java.util.Map;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+class MetricProfileServiceTest {
+
+ private final MetricProfileService service = new MetricProfileService();
+
+ @Test
+ void listProfilesShouldExposeRocketmq4And5Mappings() {
+ Map<String, MetricProfileVO> profiles = service.listProfiles().stream()
+ .collect(Collectors.toMap(MetricProfileVO::getId,
Function.identity()));
+
+
assertThat(profiles.keySet()).containsExactlyInAnyOrder("rocketmq4-exporter",
"rocketmq5-native");
+ assertThat(semanticMetrics(profiles.get("rocketmq4-exporter")))
+ .containsExactlyInAnyOrderElementsOf(allSemanticMetricKeys());
+ assertThat(semanticMetrics(profiles.get("rocketmq5-native")))
+ .containsExactlyInAnyOrderElementsOf(allSemanticMetricKeys());
+ }
+
+ @Test
+ void rocketmq5ProfileShouldUseNativeMetricNames() {
+ MetricProfileVO profile = findProfile("rocketmq5-native");
+
+ assertThat(mapping(profile, SemanticMetric.MESSAGE_IN_TPS))
+
.extracting(MetricProfileVO.MetricMappingVO::getPrometheusMetric,
+ MetricProfileVO.MetricMappingVO::getPromql)
+ .containsExactly("rocketmq_messages_in_total",
+ "sum(rate(rocketmq_messages_in_total[1m])) by
(cluster, node_id)");
+ assertThat(mapping(profile,
SemanticMetric.CONSUMER_LAG_MESSAGES).getPrometheusMetric())
+ .isEqualTo("rocketmq_consumer_lag_messages");
+ assertThat(mapping(profile,
SemanticMetric.BROKER_HEALTH).getPrometheusMetric())
+ .isEqualTo("rocketmq_processor_watermark");
+ }
+
+ @Test
+ void rocketmq4ProfileShouldUseExporterMetricNames() {
+ MetricProfileVO profile = findProfile("rocketmq4-exporter");
+
+ assertThat(mapping(profile,
SemanticMetric.MESSAGE_IN_TPS).getPrometheusMetric())
+ .isEqualTo("rocketmq_broker_tps");
+ assertThat(mapping(profile,
SemanticMetric.THROUGHPUT_IN).getPrometheusMetric())
+ .isEqualTo("rocketmq_producer_message_size");
+ assertThat(mapping(profile,
SemanticMetric.THROUGHPUT_OUT).getPrometheusMetric())
+ .isEqualTo("rocketmq_consumer_message_size");
+ assertThat(mapping(profile,
SemanticMetric.CONSUMER_LAG_MESSAGES).getPrometheusMetric())
+ .isEqualTo("rocketmq_message_accumulation");
+ assertThat(mapping(profile,
SemanticMetric.CONSUMER_LAG_LATENCY).getPrometheusMetric())
+ .isEqualTo("rocketmq_group_get_latency_by_storetime");
+ }
+
+ @Test
+ void mappingsShouldExposeLabelsAndUnitsForDashboardRendering() {
+ MetricProfileVO profile = findProfile("rocketmq5-native");
+ MetricProfileVO.MetricMappingVO messageOut = mapping(profile,
SemanticMetric.MESSAGE_OUT_TPS);
+
+ assertThat(messageOut.getLabels()).contains("cluster", "node_id",
"topic", "consumer_group");
+ assertThat(messageOut.getUnit()).isEqualTo("messages/s");
+ assertThat(messageOut.getName()).isEqualTo("Message Out TPS");
+ }
+
+ private MetricProfileVO findProfile(String id) {
+ return service.listProfiles().stream()
+ .filter(profile -> profile.getId().equals(id))
+ .findFirst()
+ .orElseThrow();
+ }
+
+ private MetricProfileVO.MetricMappingVO mapping(MetricProfileVO profile,
SemanticMetric semanticMetric) {
+ return profile.getMetrics().stream()
+ .filter(metric ->
metric.getSemanticMetric().equals(semanticMetric.getKey()))
+ .findFirst()
+ .orElseThrow();
+ }
+
+ private List<String> semanticMetrics(MetricProfileVO profile) {
+ return profile.getMetrics().stream()
+ .map(MetricProfileVO.MetricMappingVO::getSemanticMetric)
+ .toList();
+ }
+
+ private List<String> allSemanticMetricKeys() {
+ return List.of(SemanticMetric.values()).stream()
+ .map(SemanticMetric::getKey)
+ .toList();
+ }
+}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
index 659827e3..3b1af182 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
@@ -24,7 +24,11 @@ import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.http.MediaType;
import org.springframework.test.web.servlet.MockMvc;
+import java.util.List;
+
+import static org.mockito.Mockito.when;
import static org.mockito.Mockito.verifyNoInteractions;
+import static
org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get;
import static
org.springframework.test.web.servlet.request.MockMvcRequestBuilders.post;
import static
org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath;
import static
org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@@ -39,6 +43,74 @@ class MetricsControllerTest {
@MockBean
private MetricsService metricsService;
+ @MockBean
+ private MetricProfileService metricProfileService;
+
+ @Test
+ void listProfilesShouldReturnVersionAwareMappings() throws Exception {
+ MetricProfileVO profile = MetricProfileVO.builder()
+ .id("rocketmq5-native")
+ .name("RocketMQ 5.x Native")
+ .metrics(List.of(MetricProfileVO.MetricMappingVO.builder()
+ .semanticMetric("message_in_tps")
+ .prometheusMetric("rocketmq_messages_in_total")
+ .promql("sum(rate(rocketmq_messages_in_total[1m])) by
(cluster, node_id)")
+ .labels(List.of("cluster", "node_id"))
+ .build()))
+ .build();
+ when(metricProfileService.listProfiles()).thenReturn(List.of(profile));
+
+ mockMvc.perform(get("/api/metrics/profiles"))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.code").value(200))
+ .andExpect(jsonPath("$.data[0].id").value("rocketmq5-native"))
+
.andExpect(jsonPath("$.data[0].metrics[0].semanticMetric").value("message_in_tps"))
+ .andExpect(jsonPath("$.data[0].metrics[0].prometheusMetric")
+ .value("rocketmq_messages_in_total"));
+ }
+
+ @Test
+ void queryShouldReturnBadRequestWhenMetricIsBlank() throws Exception {
+ mockMvc.perform(post("/api/metrics/query")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+ {"metric":"
","start":1784107658,"end":1784108558,"step":"30s"}
+ """))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("Metric query is
required"));
+
+ verifyNoInteractions(metricsService);
+ }
+
+ @Test
+ void queryShouldReturnBadRequestWhenStartIsNotPositive() throws Exception {
+ mockMvc.perform(post("/api/metrics/query")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+
{"metric":"up","start":0,"end":1784108558,"step":"30s"}
+ """))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("Metric query start
must be positive"));
+
+ verifyNoInteractions(metricsService);
+ }
+
+ @Test
+ void queryShouldReturnBadRequestWhenStepIsBlank() throws Exception {
+ mockMvc.perform(post("/api/metrics/query")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+
{"metric":"up","start":1784107658,"end":1784108558,"step":""}
+ """))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("Metric query step is
required"));
+
+ verifyNoInteractions(metricsService);
+ }
+
@Test
void queryShouldReturnBadRequestWhenFieldTypeIsInvalid() throws Exception {
mockMvc.perform(post("/api/metrics/query")
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 7e2929ad..b67e4adb 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
@@ -26,9 +26,11 @@ import java.util.Collections;
import java.util.List;
import java.util.Map;
+import static org.assertj.core.api.Assertions.assertThatExceptionOfType;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
@@ -102,16 +104,26 @@ class MetricsServiceTest {
void queryShouldHandleVariousStepSizes() {
MetricQueryDTO query15s =
MetricQueryDTO.builder().metric("cpu").start(1L).end(2L).step("15s").build();
MetricQueryDTO query1h =
MetricQueryDTO.builder().metric("cpu").start(1L).end(2L).step("1h").build();
+ MetricQueryDTO queryCombined =
MetricQueryDTO.builder().metric("cpu").start(1L).end(7_201L)
+ .step("1h30m").build();
+ MetricQueryDTO queryNumeric =
MetricQueryDTO.builder().metric("cpu").start(1L).end(2L)
+ .step("0.5").build();
MetricDataVO data = emptyMetricData();
when(metricsSource.query(any(MetricQueryDTO.class))).thenReturn(data);
MetricDataVO result15s = metricsService.query(query15s);
MetricDataVO result1h = metricsService.query(query1h);
+ MetricDataVO resultCombined = metricsService.query(queryCombined);
+ MetricDataVO resultNumeric = metricsService.query(queryNumeric);
assertThat(result15s).isNotNull();
assertThat(result1h).isNotNull();
+ assertThat(resultCombined).isNotNull();
+ assertThat(resultNumeric).isNotNull();
verify(metricsSource).query(query15s);
verify(metricsSource).query(query1h);
+ verify(metricsSource).query(queryCombined);
+ verify(metricsSource).query(queryNumeric);
}
@Test
@@ -151,6 +163,58 @@ class MetricsServiceTest {
}
}
+ @Test
+ void queryShouldRejectInvalidWindow() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("rocketmq_messages_in_total")
+ .start(1700003600L)
+ .end(1700000000L)
+ .step("1m")
+ .build();
+
+ assertBadRequest(query, "Metric query end must be later than start");
+ verifyNoInteractions(metricsSource);
+ }
+
+ @Test
+ void queryShouldRejectOversizedWindow() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("rocketmq_messages_in_total")
+ .start(1700000000L)
+ .end(1702678401L)
+ .step("1h")
+ .build();
+
+ assertBadRequest(query, "Metric query range must not exceed 31 days");
+ verifyNoInteractions(metricsSource);
+ }
+
+ @Test
+ void queryShouldRejectInvalidStep() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("rocketmq_messages_in_total")
+ .start(1700000000L)
+ .end(1700003600L)
+ .step("five minutes")
+ .build();
+
+ assertBadRequest(query, "Metric query step is invalid");
+ verifyNoInteractions(metricsSource);
+ }
+
+ @Test
+ void queryShouldRejectTooManySamples() {
+ MetricQueryDTO query = MetricQueryDTO.builder()
+ .metric("rocketmq_messages_in_total")
+ .start(1700000000L)
+ .end(1700011001L)
+ .step("1s")
+ .build();
+
+ assertBadRequest(query, "Metric query returns too many samples;
increase step or reduce range");
+ verifyNoInteractions(metricsSource);
+ }
+
private MetricDataVO emptyMetricData() {
return MetricDataVO.builder()
.resultType("matrix")
@@ -177,4 +241,13 @@ class MetricsServiceTest {
.value(value)
.build();
}
+
+ private void assertBadRequest(MetricQueryDTO query, String message) {
+ assertThatExceptionOfType(PrometheusException.class)
+ .isThrownBy(() -> metricsService.query(query))
+ .satisfies(exception -> {
+ assertThat(exception.getStatusCode()).isEqualTo(400);
+ assertThat(exception.getMessage()).isEqualTo(message);
+ });
+ }
}
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 116ed9a1..32d0e502 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
@@ -244,6 +244,70 @@ class PrometheusMetricsSourceTest {
.hasMessage("Metric query end must not be earlier than start");
}
+ @Test
+ void queryShouldRejectBlankMetric() {
+ MetricQueryDTO invalidQuery = MetricQueryDTO.builder()
+ .metric(" ")
+ .start(1L)
+ .end(2L)
+ .step("30s")
+ .build();
+
+ assertThatThrownBy(() ->
source(Duration.ofSeconds(2)).query(invalidQuery))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> assertThat(((PrometheusException)
exception).getStatusCode())
+ .isEqualTo(HttpStatus.BAD_REQUEST.value()))
+ .hasMessage("Metric query is required");
+ }
+
+ @Test
+ void queryShouldRejectNonPositiveStart() {
+ MetricQueryDTO invalidQuery = MetricQueryDTO.builder()
+ .metric("up")
+ .start(0L)
+ .end(2L)
+ .step("30s")
+ .build();
+
+ assertThatThrownBy(() ->
source(Duration.ofSeconds(2)).query(invalidQuery))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> assertThat(((PrometheusException)
exception).getStatusCode())
+ .isEqualTo(HttpStatus.BAD_REQUEST.value()))
+ .hasMessage("Metric query start must be positive");
+ }
+
+ @Test
+ void queryShouldRejectNonPositiveEnd() {
+ MetricQueryDTO invalidQuery = MetricQueryDTO.builder()
+ .metric("up")
+ .start(1L)
+ .end(0L)
+ .step("30s")
+ .build();
+
+ assertThatThrownBy(() ->
source(Duration.ofSeconds(2)).query(invalidQuery))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> assertThat(((PrometheusException)
exception).getStatusCode())
+ .isEqualTo(HttpStatus.BAD_REQUEST.value()))
+ .hasMessage("Metric query end must be positive");
+ }
+
+ @Test
+ void queryShouldRejectBlankStep() {
+ MetricQueryDTO invalidQuery = MetricQueryDTO.builder()
+ .metric("up")
+ .start(1L)
+ .end(2L)
+ .step(" ")
+ .build();
+
+ assertThatThrownBy(() ->
source(Duration.ofSeconds(2)).query(invalidQuery))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> assertThat(((PrometheusException)
exception).getStatusCode())
+ .isEqualTo(HttpStatus.BAD_REQUEST.value()))
+ .hasMessage("Metric query step is required");
+ }
+
@Test
void queryShouldFailLoudWhenPrometheusIsNotConfigured() {
PrometheusProperties properties = new PrometheusProperties();
diff --git a/web/src/api/metrics.test.ts b/web/src/api/metrics.test.ts
index 9db303e5..092fea5c 100644
--- a/web/src/api/metrics.test.ts
+++ b/web/src/api/metrics.test.ts
@@ -57,11 +57,26 @@ describe('metrics API', () => {
it('posts a metrics query and returns its result', async () => {
const query = { metric: 'TPS_IN', start: 1, end: 2, step: '1m' };
+ const result = {
+ resultType: 'matrix',
+ series: [
+ {
+ labels: {
+ __name__: 'rocketmq_messages_in_total',
+ cluster: 'rmq-cn-v5-prod-01',
+ },
+ values: [{ timestamp: 1, value: '8' }],
+ histograms: [],
+ },
+ ],
+ warnings: [],
+ };
+
mock.onPost('/metrics/query').reply((config) => {
expect(JSON.parse(config.data)).toEqual(query);
- return [200, { code: 200, data: [{ timestamp: 1, value: 8 }] }];
+ return [200, { code: 200, data: result }];
});
- await expect(queryMetrics(query)).resolves.toEqual([{ timestamp: 1, value:
8 }]);
+ await expect(queryMetrics(query)).resolves.toEqual(result);
});
});