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 407fe51 feat: add Prometheus range query adapter (#432)
407fe51 is described below
commit 407fe5158f5066241016cbe5e9a8050b1930aa7a
Author: wizcraft_kris <[email protected]>
AuthorDate: Wed Jul 22 17:30:51 2026 +0800
feat: add Prometheus range query adapter (#432)
Add a real Prometheus /api/v1/query_range adapter as the foundation for
observability (#431): configurable base URL, timeouts, Basic/Bearer auth, error
mapping, and tests.
---
docs/api-spec.md | 76 ++++-
.../studio/cluster/metrics/MetricDataVO.java | 35 ++-
.../studio/cluster/metrics/MetricQueryDTO.java | 23 ++
.../studio/cluster/metrics/MetricsController.java | 17 +-
.../studio/cluster/metrics/MetricsService.java | 4 +-
...{MetricDataVO.java => PrometheusException.java} | 25 +-
.../cluster/metrics/PrometheusMetricsSource.java | 260 +++++++++++++++--
...tricsService.java => PrometheusProperties.java} | 30 +-
.../common/exception/GlobalExceptionHandler.java | 28 ++
server/src/main/resources/application.yml | 10 +
.../cluster/metrics/MetricsControllerTest.java | 55 ++++
.../studio/cluster/metrics/MetricsServiceTest.java | 80 +++--
.../metrics/PrometheusMetricsSourceTest.java | 321 +++++++++++++++++++++
web/src/api/metrics.ts | 44 ++-
14 files changed, 912 insertions(+), 96 deletions(-)
diff --git a/docs/api-spec.md b/docs/api-spec.md
index 581e8d0..2b0bd14 100644
--- a/docs/api-spec.md
+++ b/docs/api-spec.md
@@ -1652,17 +1652,83 @@ POST /api/metrics/query
| 字段 | 类型 | 必填 | 说明 |
|------|------|------|------|
-| `metric` | `string` | 是 | 指标名称(如
`tps_in`、`tps_out`、`message_count`、`disk_usage`) |
+| `metric` | `string` | 是 | PromQL 表达式,最大 4096 个字符 |
| `start` | `number` | 是 | 起始时间(Unix 时间戳,秒) |
| `end` | `number` | 是 | 结束时间(Unix 时间戳,秒) |
-| `step` | `string` | 否 | 采样步长(如 `"60s"`、`"5m"`、`"1h"`) |
+| `step` | `string` | 是 | 查询分辨率,可以是持续时间或秒数(如 `"30s"`、`"5m"`、`"1h"`) |
-**Response `data`:** `MetricsResult`
+**Request 示例:**
+
+```json
+{
+ "metric": "sum(rate(rocketmq_messages_in_total[1m])) by (node_id)",
+ "start": 1784112606,
+ "end": 1784114406,
+ "step": "30s"
+}
+```
+
+**Response `data`:** `MetricData`
+
+| 字段 | 类型 | 说明 |
+|------|------|------|
+| `resultType` | `string` | Prometheus 结果类型;范围查询通常为 `matrix` |
+| `series` | `MetricSeries[]` | 查询返回的时间序列 |
+| `warnings` | `string[]` | Prometheus 返回的非致命警告,没有警告时为空数组 |
+
+**MetricSeries:**
+
+| 字段 | 类型 | 说明 |
+|------|------|------|
+| `labels` | `object` | 序列的完整标签集合,包括可能存在的 `__name__` |
+| `values` | `MetricSample[]` | 浮点样本;没有浮点样本时为空数组 |
+| `histograms` | `MetricHistogramSample[]` | Native Histogram 样本;没有 Histogram
样本时为空数组 |
+
+同一序列可能只有 `values`、只有 `histograms`,或同时包含两者。
+
+**MetricSample:**
+
+| 字段 | 类型 | 说明 |
+|------|------|------|
+| `timestamp` | `number` | Unix 时间戳,可能包含小数秒 |
+| `value` | `string` | Prometheus 样本原始字符串,保留小数精度以及 `NaN`、`+Inf`、`-Inf` |
+
+**MetricHistogramSample:**
| 字段 | 类型 | 说明 |
|------|------|------|
-| `metric` | `string` | 指标名称 |
-| `values` | `[number, number][]` | 数据点数组,每项为 `[timestamp, value]` |
+| `timestamp` | `number` | Unix 时间戳,可能包含小数秒 |
+| `histogram` | `object` | Prometheus Native Histogram 原始对象,包含 `count`、`sum` 和
`buckets` |
+
+**Prometheus 配置:**
+
+```yaml
+studio:
+ metrics:
+ prometheus:
+ base-url: ${STUDIO_METRICS_PROMETHEUS_BASE_URL:}
+ connect-timeout: ${STUDIO_METRICS_PROMETHEUS_CONNECT_TIMEOUT:3s}
+ read-timeout: ${STUDIO_METRICS_PROMETHEUS_READ_TIMEOUT:10s}
+ username: ${STUDIO_METRICS_PROMETHEUS_USERNAME:}
+ password: ${STUDIO_METRICS_PROMETHEUS_PASSWORD:}
+ bearer-token: ${STUDIO_METRICS_PROMETHEUS_BEARER_TOKEN:}
+```
+
+- `base-url` 是 Prometheus 或 Prometheus-compatible 服务的 URL 前缀;服务会在其后追加
`/api/v1/query_range`。
+- 未配置 `base-url` 时,查询接口返回 HTTP 503。
+- `connect-timeout` 默认 3 秒,`read-timeout` 默认 10 秒。
+- Bearer Token 的优先级高于 Basic Auth;Basic Auth 必须同时配置 `username` 和 `password`。
+- 密码和 Token 等敏感配置应通过环境变量或其他外部化配置传入,不应提交到代码仓库。
+
+**错误响应:**
+
+| HTTP 状态 | 场景 |
+|-----------|------|
+| `400` | JSON 无法解析、字段类型错误、请求校验失败或 PromQL 参数错误 |
+| `422` | Prometheus 无法执行 PromQL 表达式 |
+| `502` | 无法连接 Prometheus,或 Prometheus 返回非法响应 |
+| `503` | Prometheus 未配置或暂时不可用 |
+| `504` | Prometheus 查询超时 |
---
diff --git
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricDataVO.java
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricDataVO.java
index 06fbf70..29846dd 100644
--- a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricDataVO.java
+++ b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricDataVO.java
@@ -16,11 +16,13 @@
*/
package com.rocketmq.studio.cluster.metrics;
+import com.fasterxml.jackson.databind.JsonNode;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
+import java.util.Map;
import java.util.List;
@Data
@@ -28,6 +30,35 @@ import java.util.List;
@NoArgsConstructor
@AllArgsConstructor
public class MetricDataVO {
- private String metric;
- private List<long[]> values;
+ private String resultType;
+ private List<MetricSeriesVO> series;
+ private List<String> warnings;
+
+ @Data
+ @Builder
+ @NoArgsConstructor
+ @AllArgsConstructor
+ public static class MetricSeriesVO {
+ private Map<String, String> labels;
+ private List<MetricSampleVO> values;
+ private List<MetricHistogramSampleVO> histograms;
+ }
+
+ @Data
+ @Builder
+ @NoArgsConstructor
+ @AllArgsConstructor
+ public static class MetricSampleVO {
+ private double timestamp;
+ private String value;
+ }
+
+ @Data
+ @Builder
+ @NoArgsConstructor
+ @AllArgsConstructor
+ public static class MetricHistogramSampleVO {
+ private double timestamp;
+ private JsonNode histogram;
+ }
}
diff --git
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
index 0884ea4..de032a0 100644
---
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
+++
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricQueryDTO.java
@@ -16,6 +16,10 @@
*/
package com.rocketmq.studio.cluster.metrics;
+import io.swagger.v3.oas.annotations.media.Schema;
+import jakarta.validation.constraints.NotBlank;
+import jakarta.validation.constraints.Positive;
+import jakarta.validation.constraints.Size;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
@@ -25,9 +29,28 @@ import lombok.NoArgsConstructor;
@Builder
@NoArgsConstructor
@AllArgsConstructor
+@Schema(description = "Prometheus range query")
public class MetricQueryDTO {
+ @Schema(description = "PromQL expression evaluated by Prometheus",
+ example = "sum(rate(rocketmq_messages_in_total[1m])) by
(node_id)", minLength = 1,
+ requiredMode = Schema.RequiredMode.REQUIRED)
+ @NotBlank(message = "Metric query is required")
+ @Size(max = 4096, message = "Metric query must not exceed 4096 characters")
private String metric;
+
+ @Schema(description = "Range start as a Unix timestamp in seconds",
example = "1784112606",
+ requiredMode = Schema.RequiredMode.REQUIRED)
+ @Positive(message = "Metric query start must be positive")
private long start;
+
+ @Schema(description = "Range end as a Unix timestamp in seconds", example
= "1784114406",
+ requiredMode = Schema.RequiredMode.REQUIRED)
+ @Positive(message = "Metric query end must be positive")
private long end;
+
+ @Schema(description = "Prometheus query resolution step as a duration or
number of seconds", example = "30s",
+ minLength = 1, requiredMode = Schema.RequiredMode.REQUIRED)
+ @NotBlank(message = "Metric query step is required")
+ @Size(max = 32, message = "Metric query step must not exceed 32
characters")
private String step;
}
diff --git
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsController.java
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsController.java
index 9015c01..6ee8334 100644
---
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsController.java
+++
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsController.java
@@ -17,6 +17,10 @@
package com.rocketmq.studio.cluster.metrics;
import com.rocketmq.studio.common.domain.Result;
+import io.swagger.v3.oas.annotations.Operation;
+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.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
@@ -30,8 +34,19 @@ public class MetricsController {
private final MetricsService metricsService;
+ @Operation(summary = "Query Prometheus range metrics",
+ description = "Executes a PromQL range query against the
configured Prometheus server")
+ @ApiResponses({
+ @ApiResponse(responseCode = "200", description = "Range query
completed successfully",
+ useReturnTypeSchema = true),
+ @ApiResponse(responseCode = "400", description = "Invalid request or
PromQL expression"),
+ @ApiResponse(responseCode = "422", description = "Prometheus could not
execute the expression"),
+ @ApiResponse(responseCode = "502", description = "Prometheus
connection or response failure"),
+ @ApiResponse(responseCode = "503", description = "Prometheus is
unavailable or not configured"),
+ @ApiResponse(responseCode = "504", description = "Prometheus query
timed out")
+ })
@PostMapping("/query")
- public Result<MetricDataVO> query(@RequestBody MetricQueryDTO query) {
+ public Result<MetricDataVO> query(@Valid @RequestBody MetricQueryDTO
query) {
return Result.ok(metricsService.query(query));
}
}
diff --git
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsService.java
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsService.java
index 0a38127..65bbe77 100644
---
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsService.java
+++
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsService.java
@@ -28,8 +28,8 @@ public class MetricsService {
private final MetricsSource metricsSource;
public MetricDataVO query(MetricQueryDTO query) {
- log.info("Querying metrics: metric={}, start={}, end={}, step={}",
- query.getMetric(), query.getStart(), query.getEnd(),
query.getStep());
+ log.debug("Querying metrics: start={}, end={}, step={}",
+ query.getStart(), query.getEnd(), query.getStep());
return metricsSource.query(query);
}
}
diff --git
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricDataVO.java
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusException.java
similarity index 67%
copy from
server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricDataVO.java
copy to
server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusException.java
index 06fbf70..e286c66 100644
--- a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricDataVO.java
+++
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusException.java
@@ -16,18 +16,19 @@
*/
package com.rocketmq.studio.cluster.metrics;
-import lombok.AllArgsConstructor;
-import lombok.Builder;
-import lombok.Data;
-import lombok.NoArgsConstructor;
+import lombok.Getter;
-import java.util.List;
+@Getter
+public class PrometheusException extends RuntimeException {
+ private final int statusCode;
-@Data
-@Builder
-@NoArgsConstructor
-@AllArgsConstructor
-public class MetricDataVO {
- private String metric;
- private List<long[]> values;
+ public PrometheusException(int statusCode, String message) {
+ super(message);
+ this.statusCode = statusCode;
+ }
+
+ public PrometheusException(int statusCode, String message, Throwable
cause) {
+ super(message, cause);
+ this.statusCode = statusCode;
+ }
}
diff --git
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
index cc48fe5..701bab1 100644
---
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
+++
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusMetricsSource.java
@@ -16,53 +16,263 @@
*/
package com.rocketmq.studio.cluster.metrics;
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
+import org.springframework.http.HttpHeaders;
+import org.springframework.http.HttpStatus;
+import org.springframework.http.MediaType;
+import org.springframework.http.client.SimpleClientHttpRequestFactory;
import org.springframework.stereotype.Component;
+import org.springframework.util.LinkedMultiValueMap;
+import org.springframework.util.MultiValueMap;
+import org.springframework.util.StringUtils;
+import org.springframework.web.client.ResourceAccessException;
+import org.springframework.web.client.RestClient;
+import org.springframework.web.client.RestClientException;
+import org.springframework.web.client.RestClientResponseException;
-import java.util.ArrayList;
+import java.io.IOException;
+import java.net.SocketTimeoutException;
+import java.net.URI;
+import java.util.Iterator;
+import java.util.LinkedHashMap;
import java.util.List;
+import java.util.Map;
+import java.util.stream.StreamSupport;
@Slf4j
@Component
public class PrometheusMetricsSource implements MetricsSource {
+ private static final String QUERY_RANGE_PATH = "/api/v1/query_range";
+
+ private final RestClient restClient;
+ private final ObjectMapper objectMapper;
+ private final PrometheusProperties properties;
+
+ public PrometheusMetricsSource(RestClient.Builder restClientBuilder,
ObjectMapper objectMapper,
+ PrometheusProperties properties) {
+ SimpleClientHttpRequestFactory requestFactory = new
SimpleClientHttpRequestFactory();
+ requestFactory.setConnectTimeout(properties.getConnectTimeout());
+ requestFactory.setReadTimeout(properties.getReadTimeout());
+ this.restClient =
restClientBuilder.requestFactory(requestFactory).build();
+ this.objectMapper = objectMapper;
+ this.properties = properties;
+ }
+
@Override
public MetricDataVO query(MetricQueryDTO query) {
- log.info("Querying Prometheus: metric={}, start={}, end={}, step={}",
- query.getMetric(), query.getStart(), query.getEnd(),
query.getStep());
+ validateQuery(query);
+ URI queryRangeUri = queryRangeUri();
+ MultiValueMap<String, String> form = new LinkedMultiValueMap<>();
+ form.add("query", query.getMetric());
+ form.add("start", Long.toString(query.getStart()));
+ form.add("end", Long.toString(query.getEnd()));
+ form.add("step", query.getStep());
+
+ log.debug("Querying Prometheus range: start={}, end={}, step={}",
+ query.getStart(), query.getEnd(), query.getStep());
+
+ try {
+ JsonNode response = restClient.post()
+ .uri(queryRangeUri)
+ .contentType(MediaType.APPLICATION_FORM_URLENCODED)
+ .headers(this::applyAuthentication)
+ .body(form)
+ .retrieve()
+ .body(JsonNode.class);
+ return parseResponse(response);
+ } catch (PrometheusException exception) {
+ throw exception;
+ } catch (RestClientResponseException exception) {
+ throw responseException(exception);
+ } catch (ResourceAccessException exception) {
+ if (hasCause(exception, SocketTimeoutException.class)) {
+ throw new
PrometheusException(HttpStatus.GATEWAY_TIMEOUT.value(),
+ "Prometheus query timed out", exception);
+ }
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ "Failed to connect to Prometheus", exception);
+ } catch (RestClientException exception) {
+ if (hasCause(exception, SocketTimeoutException.class)) {
+ throw new
PrometheusException(HttpStatus.GATEWAY_TIMEOUT.value(),
+ "Prometheus query timed out", exception);
+ }
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ "Prometheus query failed", exception);
+ }
+ }
+
+ private void validateQuery(MetricQueryDTO query) {
+ if (query == null) {
+ throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
"Metric query is required");
+ }
+ if (query.getEnd() < query.getStart()) {
+ throw new PrometheusException(HttpStatus.BAD_REQUEST.value(),
+ "Metric query end must not be earlier than start");
+ }
+ }
- // Stub: generate sample data points
- List<long[]> values = new ArrayList<>();
- long stepSeconds = parseStep(query.getStep());
- long start = query.getStart();
- long end = query.getEnd();
+ private URI queryRangeUri() {
+ if (!StringUtils.hasText(properties.getBaseUrl())) {
+ throw new
PrometheusException(HttpStatus.SERVICE_UNAVAILABLE.value(),
+ "Prometheus base URL is not configured");
+ }
+ try {
+ String baseUrl = properties.getBaseUrl().strip();
+ while (baseUrl.endsWith("/")) {
+ baseUrl = baseUrl.substring(0, baseUrl.length() - 1);
+ }
+ URI uri = URI.create(baseUrl + QUERY_RANGE_PATH);
+ if (!"http".equalsIgnoreCase(uri.getScheme()) &&
!"https".equalsIgnoreCase(uri.getScheme())) {
+ throw new IllegalArgumentException("Unsupported Prometheus URL
scheme");
+ }
+ return uri;
+ } catch (IllegalArgumentException exception) {
+ throw new
PrometheusException(HttpStatus.SERVICE_UNAVAILABLE.value(),
+ "Prometheus base URL is invalid", exception);
+ }
+ }
+
+ private void applyAuthentication(HttpHeaders headers) {
+ if (StringUtils.hasText(properties.getBearerToken())) {
+ headers.setBearerAuth(properties.getBearerToken());
+ return;
+ }
+ boolean hasUsername = StringUtils.hasText(properties.getUsername());
+ boolean hasPassword = StringUtils.hasText(properties.getPassword());
+ if (hasUsername != hasPassword) {
+ throw new
PrometheusException(HttpStatus.SERVICE_UNAVAILABLE.value(),
+ "Prometheus basic authentication is incomplete");
+ }
+ if (hasUsername) {
+ headers.setBasicAuth(properties.getUsername(),
properties.getPassword());
+ }
+ }
- for (long ts = start; ts <= end; ts += stepSeconds) {
- values.add(new long[]{ts, (long) (Math.random() * 100)});
+ private MetricDataVO parseResponse(JsonNode response) {
+ if (response == null ||
!"success".equals(response.path("status").asText())) {
+ throw responseBodyException(response,
HttpStatus.BAD_GATEWAY.value());
}
+ JsonNode data = response.path("data");
+ JsonNode result = data.path("result");
+ if (!data.isObject() || !result.isArray() ||
!StringUtils.hasText(data.path("resultType").asText())) {
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ "Prometheus returned a malformed response");
+ }
+
+ List<MetricDataVO.MetricSeriesVO> series =
StreamSupport.stream(result.spliterator(), false)
+ .map(this::parseSeries)
+ .toList();
+ List<String> warnings = parseWarnings(response.path("warnings"));
+
return MetricDataVO.builder()
- .metric(query.getMetric())
- .values(values)
+ .resultType(data.path("resultType").asText())
+ .series(series)
+ .warnings(warnings)
+ .build();
+ }
+
+ private MetricDataVO.MetricSeriesVO parseSeries(JsonNode seriesNode) {
+ JsonNode metric = seriesNode.path("metric");
+ JsonNode values = seriesNode.path("values");
+ JsonNode histograms = seriesNode.path("histograms");
+ boolean hasValues = values.isArray();
+ boolean hasHistograms = histograms.isArray();
+ boolean hasSamples = hasValues || hasHistograms;
+ boolean invalidValues = !values.isMissingNode() && !hasValues;
+ boolean invalidHistograms = !histograms.isMissingNode() &&
!hasHistograms;
+ if (!metric.isObject() || invalidValues || invalidHistograms ||
!hasSamples) {
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ "Prometheus returned a malformed time series");
+ }
+
+ Map<String, String> labels = new LinkedHashMap<>();
+ Iterator<Map.Entry<String, JsonNode>> fields = metric.fields();
+ fields.forEachRemaining(entry -> labels.put(entry.getKey(),
entry.getValue().asText()));
+
+ List<MetricDataVO.MetricSampleVO> samples = hasValues
+ ? StreamSupport.stream(values.spliterator(),
false).map(this::parseSample).toList()
+ : List.of();
+ List<MetricDataVO.MetricHistogramSampleVO> histogramSamples =
hasHistograms
+ ? StreamSupport.stream(histograms.spliterator(),
false).map(this::parseHistogramSample).toList()
+ : List.of();
+ return MetricDataVO.MetricSeriesVO.builder()
+ .labels(labels)
+ .values(samples)
+ .histograms(histogramSamples)
+ .build();
+ }
+
+ private MetricDataVO.MetricSampleVO parseSample(JsonNode sampleNode) {
+ if (!sampleNode.isArray() || sampleNode.size() != 2 ||
!sampleNode.get(0).isNumber()) {
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ "Prometheus returned a malformed sample");
+ }
+ return MetricDataVO.MetricSampleVO.builder()
+ .timestamp(sampleNode.get(0).asDouble())
+ .value(sampleNode.get(1).asText())
.build();
}
- private long parseStep(String step) {
- if (step == null || step.isEmpty()) {
- return 60;
+ private MetricDataVO.MetricHistogramSampleVO parseHistogramSample(JsonNode
sampleNode) {
+ if (!sampleNode.isArray() || sampleNode.size() != 2
+ || !sampleNode.get(0).isNumber() ||
!sampleNode.get(1).isObject()) {
+ throw new PrometheusException(HttpStatus.BAD_GATEWAY.value(),
+ "Prometheus returned a malformed histogram sample");
}
+ return MetricDataVO.MetricHistogramSampleVO.builder()
+ .timestamp(sampleNode.get(0).asDouble())
+ .histogram(sampleNode.get(1))
+ .build();
+ }
+
+ private List<String> parseWarnings(JsonNode warningsNode) {
+ if (!warningsNode.isArray()) {
+ return List.of();
+ }
+ return StreamSupport.stream(warningsNode.spliterator(), false)
+ .map(JsonNode::asText)
+ .toList();
+ }
+
+ private PrometheusException responseException(RestClientResponseException
exception) {
+ JsonNode response = null;
try {
- if (step.endsWith("s")) {
- return Long.parseLong(step.substring(0, step.length() - 1));
- } else if (step.endsWith("m")) {
- return Long.parseLong(step.substring(0, step.length() - 1)) *
60;
- } else if (step.endsWith("h")) {
- return Long.parseLong(step.substring(0, step.length() - 1)) *
3600;
+ response =
objectMapper.readTree(exception.getResponseBodyAsString());
+ } catch (IOException ignored) {
+ log.debug("Failed to parse Prometheus error response");
+ }
+ int upstreamStatus = exception.getStatusCode().value();
+ int statusCode = switch (upstreamStatus) {
+ case 400, 422, 503 -> upstreamStatus;
+ default -> HttpStatus.BAD_GATEWAY.value();
+ };
+ return responseBodyException(response, statusCode);
+ }
+
+ private PrometheusException responseBodyException(JsonNode response, int
statusCode) {
+ String errorType = response == null ? "" :
response.path("errorType").asText();
+ String error = response == null ? "" : response.path("error").asText();
+ if (StringUtils.hasText(error)) {
+ String message = StringUtils.hasText(errorType)
+ ? "Prometheus query failed (" + errorType + "): " + error
+ : "Prometheus query failed: " + error;
+ return new PrometheusException(statusCode, message);
+ }
+ return new PrometheusException(statusCode, "Prometheus query failed");
+ }
+
+ private boolean hasCause(Throwable throwable, Class<? extends Throwable>
causeType) {
+ Throwable current = throwable;
+ while (current != null) {
+ if (causeType.isInstance(current)) {
+ return true;
}
- return Long.parseLong(step);
- } catch (NumberFormatException e) {
- log.warn("Failed to parse step '{}', defaulting to 60s", step);
- return 60;
+ current = current.getCause();
}
+ return false;
}
}
diff --git
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsService.java
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusProperties.java
similarity index 60%
copy from
server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsService.java
copy to
server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusProperties.java
index 0a38127..e37e582 100644
---
a/server/src/main/java/com/rocketmq/studio/cluster/metrics/MetricsService.java
+++
b/server/src/main/java/com/rocketmq/studio/cluster/metrics/PrometheusProperties.java
@@ -16,20 +16,22 @@
*/
package com.rocketmq.studio.cluster.metrics;
-import lombok.RequiredArgsConstructor;
-import lombok.extern.slf4j.Slf4j;
-import org.springframework.stereotype.Service;
+import lombok.Getter;
+import lombok.Setter;
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.stereotype.Component;
-@Slf4j
-@Service
-@RequiredArgsConstructor
-public class MetricsService {
+import java.time.Duration;
- private final MetricsSource metricsSource;
-
- public MetricDataVO query(MetricQueryDTO query) {
- log.info("Querying metrics: metric={}, start={}, end={}, step={}",
- query.getMetric(), query.getStart(), query.getEnd(),
query.getStep());
- return metricsSource.query(query);
- }
+@Getter
+@Setter
+@Component
+@ConfigurationProperties(prefix = "studio.metrics.prometheus")
+public class PrometheusProperties {
+ private String baseUrl;
+ private Duration connectTimeout = Duration.ofSeconds(3);
+ private Duration readTimeout = Duration.ofSeconds(10);
+ private String username;
+ private String password;
+ private String bearerToken;
}
diff --git
a/server/src/main/java/com/rocketmq/studio/common/exception/GlobalExceptionHandler.java
b/server/src/main/java/com/rocketmq/studio/common/exception/GlobalExceptionHandler.java
index 3bc1c19..26cffb9 100644
---
a/server/src/main/java/com/rocketmq/studio/common/exception/GlobalExceptionHandler.java
+++
b/server/src/main/java/com/rocketmq/studio/common/exception/GlobalExceptionHandler.java
@@ -16,10 +16,14 @@
*/
package com.rocketmq.studio.common.exception;
+import com.rocketmq.studio.cluster.metrics.PrometheusException;
import com.rocketmq.studio.common.domain.Result;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.http.HttpStatus;
+import org.springframework.http.ResponseEntity;
+import org.springframework.http.converter.HttpMessageNotReadableException;
+import org.springframework.web.bind.MethodArgumentNotValidException;
import org.springframework.web.bind.annotation.ExceptionHandler;
import org.springframework.web.bind.annotation.ResponseStatus;
import org.springframework.web.bind.annotation.RestControllerAdvice;
@@ -36,6 +40,30 @@ public class GlobalExceptionHandler {
return Result.error(ex.getCode(), ex.getMessage());
}
+ @ExceptionHandler(PrometheusException.class)
+ public ResponseEntity<Result<?>>
handlePrometheusException(PrometheusException ex) {
+ log.warn("Prometheus exception: status={}, message={}",
ex.getStatusCode(), ex.getMessage());
+ return ResponseEntity.status(ex.getStatusCode())
+ .body(Result.error(ex.getStatusCode(), ex.getMessage()));
+ }
+
+ @ExceptionHandler(MethodArgumentNotValidException.class)
+ @ResponseStatus(HttpStatus.BAD_REQUEST)
+ public Result<?> handleValidationException(MethodArgumentNotValidException
ex) {
+ String message = ex.getBindingResult().getFieldErrors().stream()
+ .findFirst()
+ .map(error -> error.getDefaultMessage() == null ? "Invalid
request" : error.getDefaultMessage())
+ .orElse("Invalid request");
+ return Result.error(HttpStatus.BAD_REQUEST.value(), message);
+ }
+
+ @ExceptionHandler(HttpMessageNotReadableException.class)
+ @ResponseStatus(HttpStatus.BAD_REQUEST)
+ public Result<?>
handleHttpMessageNotReadableException(HttpMessageNotReadableException ex) {
+ log.warn("Invalid request body");
+ return Result.error(HttpStatus.BAD_REQUEST.value(), "Invalid request
body");
+ }
+
@ExceptionHandler(Exception.class)
@ResponseStatus(HttpStatus.INTERNAL_SERVER_ERROR)
public Result<?> handleException(Exception ex) {
diff --git a/server/src/main/resources/application.yml
b/server/src/main/resources/application.yml
index 3c98b4c..a213118 100644
--- a/server/src/main/resources/application.yml
+++ b/server/src/main/resources/application.yml
@@ -12,3 +12,13 @@ springdoc:
path: /api-docs
swagger-ui:
path: /swagger-ui.html
+
+studio:
+ metrics:
+ prometheus:
+ base-url: ${STUDIO_METRICS_PROMETHEUS_BASE_URL:}
+ connect-timeout: ${STUDIO_METRICS_PROMETHEUS_CONNECT_TIMEOUT:3s}
+ read-timeout: ${STUDIO_METRICS_PROMETHEUS_READ_TIMEOUT:10s}
+ username: ${STUDIO_METRICS_PROMETHEUS_USERNAME:}
+ password: ${STUDIO_METRICS_PROMETHEUS_PASSWORD:}
+ bearer-token: ${STUDIO_METRICS_PROMETHEUS_BEARER_TOKEN:}
diff --git
a/server/src/test/java/com/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
b/server/src/test/java/com/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
new file mode 100644
index 0000000..70b64bc
--- /dev/null
+++
b/server/src/test/java/com/rocketmq/studio/cluster/metrics/MetricsControllerTest.java
@@ -0,0 +1,55 @@
+/*
+ * 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 com.rocketmq.studio.cluster.metrics;
+
+import org.junit.jupiter.api.Test;
+import org.springframework.beans.factory.annotation.Autowired;
+import
org.springframework.boot.test.autoconfigure.web.servlet.AutoConfigureMockMvc;
+import org.springframework.boot.test.autoconfigure.web.servlet.WebMvcTest;
+import org.springframework.boot.test.mock.mockito.MockBean;
+import org.springframework.http.MediaType;
+import org.springframework.test.web.servlet.MockMvc;
+
+import static org.mockito.Mockito.verifyNoInteractions;
+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;
+
+@WebMvcTest(MetricsController.class)
+@AutoConfigureMockMvc(addFilters = false)
+class MetricsControllerTest {
+
+ @Autowired
+ private MockMvc mockMvc;
+
+ @MockBean
+ private MetricsService metricsService;
+
+ @Test
+ void queryShouldReturnBadRequestWhenFieldTypeIsInvalid() throws Exception {
+ mockMvc.perform(post("/api/metrics/query")
+ .contentType(MediaType.APPLICATION_JSON)
+ .content("""
+
{"metric":"up","start":"abc","end":123,"step":"30s"}
+ """))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value(400))
+ .andExpect(jsonPath("$.message").value("Invalid request
body"));
+
+ verifyNoInteractions(metricsService);
+ }
+}
diff --git
a/server/src/test/java/com/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
b/server/src/test/java/com/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
index b689775..a8f7118 100644
---
a/server/src/test/java/com/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
+++
b/server/src/test/java/com/rocketmq/studio/cluster/metrics/MetricsServiceTest.java
@@ -22,9 +22,9 @@ import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
-import java.util.Arrays;
import java.util.Collections;
import java.util.List;
+import java.util.Map;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
@@ -48,22 +48,21 @@ class MetricsServiceTest {
.end(1700003600L)
.step("1m")
.build();
- List<long[]> values = Arrays.asList(
- new long[]{1700000000L, 45},
- new long[]{1700000060L, 52},
- new long[]{1700000120L, 48}
+ List<MetricDataVO.MetricSampleVO> values = List.of(
+ sample(1700000000L, "45.5"),
+ sample(1700000060L, "52"),
+ sample(1700000120L, "48")
);
- MetricDataVO data =
MetricDataVO.builder().metric("cpu_usage").values(values).build();
+ MetricDataVO data = metricData("cpu_usage", values);
when(metricsSource.query(query)).thenReturn(data);
MetricDataVO result = metricsService.query(query);
- assertThat(result.getMetric()).isEqualTo("cpu_usage");
- assertThat(result.getValues()).hasSize(3);
- assertThat(result.getValues().get(0)[0]).isEqualTo(1700000000L);
- assertThat(result.getValues().get(0)[1]).isEqualTo(45);
- assertThat(result.getValues().get(1)[1]).isEqualTo(52);
- assertThat(result.getValues().get(2)[1]).isEqualTo(48);
+ assertThat(result.getSeries()).hasSize(1);
+
assertThat(result.getSeries().get(0).getLabels()).containsEntry("__name__",
"cpu_usage");
+ assertThat(result.getSeries().get(0).getValues()).hasSize(3);
+
assertThat(result.getSeries().get(0).getValues().get(0).getTimestamp()).isEqualTo(1700000000D);
+
assertThat(result.getSeries().get(0).getValues().get(0).getValue()).isEqualTo("45.5");
verify(metricsSource).query(query);
}
@@ -75,13 +74,12 @@ class MetricsServiceTest {
.end(1700003600L)
.step("5m")
.build();
- MetricDataVO data =
MetricDataVO.builder().metric("disk_io").values(Collections.emptyList()).build();
+ MetricDataVO data = emptyMetricData();
when(metricsSource.query(query)).thenReturn(data);
MetricDataVO result = metricsService.query(query);
- assertThat(result.getMetric()).isEqualTo("disk_io");
- assertThat(result.getValues()).isEmpty();
+ assertThat(result.getSeries()).isEmpty();
}
@Test
@@ -92,7 +90,7 @@ class MetricsServiceTest {
.end(1700086400L)
.step("1h")
.build();
- MetricDataVO data =
MetricDataVO.builder().metric("tps").values(Collections.emptyList()).build();
+ MetricDataVO data = emptyMetricData();
when(metricsSource.query(any(MetricQueryDTO.class))).thenReturn(data);
metricsService.query(query);
@@ -104,7 +102,7 @@ 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();
- MetricDataVO data =
MetricDataVO.builder().metric("cpu").values(Collections.emptyList()).build();
+ MetricDataVO data = emptyMetricData();
when(metricsSource.query(any(MetricQueryDTO.class))).thenReturn(data);
MetricDataVO result15s = metricsService.query(query15s);
@@ -124,20 +122,19 @@ class MetricsServiceTest {
.end(1700003600L)
.step("1m")
.build();
- List<long[]> values = Arrays.asList(
- new long[]{1700000000L, 72},
- new long[]{1700000060L, 73},
- new long[]{1700000120L, 71},
- new long[]{1700000180L, 74},
- new long[]{1700000240L, 75}
+ List<MetricDataVO.MetricSampleVO> values = List.of(
+ sample(1700000000L, "72"),
+ sample(1700000060L, "73"),
+ sample(1700000120L, "71"),
+ sample(1700000180L, "74"),
+ sample(1700000240L, "75")
);
- MetricDataVO data =
MetricDataVO.builder().metric("memory_usage").values(values).build();
+ MetricDataVO data = metricData("memory_usage", values);
when(metricsSource.query(query)).thenReturn(data);
MetricDataVO result = metricsService.query(query);
- assertThat(result.getMetric()).isEqualTo("memory_usage");
- assertThat(result.getValues()).hasSize(5);
+ assertThat(result.getSeries().get(0).getValues()).hasSize(5);
}
@Test
@@ -145,12 +142,39 @@ class MetricsServiceTest {
String[] metrics = {"rocketmq_tps", "rocketmq_latency_p99",
"broker_disk_usage", "consumer_lag"};
for (String metricName : metrics) {
MetricQueryDTO query =
MetricQueryDTO.builder().metric(metricName).start(1L).end(2L).step("1m").build();
- MetricDataVO data =
MetricDataVO.builder().metric(metricName).values(Collections.emptyList()).build();
+ MetricDataVO data = metricData(metricName,
Collections.emptyList());
when(metricsSource.query(query)).thenReturn(data);
MetricDataVO result = metricsService.query(query);
- assertThat(result.getMetric()).isEqualTo(metricName);
+
assertThat(result.getSeries().get(0).getLabels()).containsEntry("__name__",
metricName);
}
}
+
+ private MetricDataVO emptyMetricData() {
+ return MetricDataVO.builder()
+ .resultType("matrix")
+ .series(Collections.emptyList())
+ .warnings(Collections.emptyList())
+ .build();
+ }
+
+ private MetricDataVO metricData(String metricName,
List<MetricDataVO.MetricSampleVO> values) {
+ MetricDataVO.MetricSeriesVO series =
MetricDataVO.MetricSeriesVO.builder()
+ .labels(Map.of("__name__", metricName))
+ .values(values)
+ .build();
+ return MetricDataVO.builder()
+ .resultType("matrix")
+ .series(List.of(series))
+ .warnings(Collections.emptyList())
+ .build();
+ }
+
+ private MetricDataVO.MetricSampleVO sample(double timestamp, String value)
{
+ return MetricDataVO.MetricSampleVO.builder()
+ .timestamp(timestamp)
+ .value(value)
+ .build();
+ }
}
diff --git
a/server/src/test/java/com/rocketmq/studio/cluster/metrics/PrometheusMetricsSourceTest.java
b/server/src/test/java/com/rocketmq/studio/cluster/metrics/PrometheusMetricsSourceTest.java
new file mode 100644
index 0000000..3d2fbb2
--- /dev/null
+++
b/server/src/test/java/com/rocketmq/studio/cluster/metrics/PrometheusMetricsSourceTest.java
@@ -0,0 +1,321 @@
+/*
+ * 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 com.rocketmq.studio.cluster.metrics;
+
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.sun.net.httpserver.HttpExchange;
+import com.sun.net.httpserver.HttpServer;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.springframework.http.HttpStatus;
+import org.springframework.web.client.RestClient;
+
+import java.io.IOException;
+import java.net.InetSocketAddress;
+import java.net.URLDecoder;
+import java.nio.charset.StandardCharsets;
+import java.time.Duration;
+import java.util.Base64;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+class PrometheusMetricsSourceTest {
+
+ private HttpServer server;
+ private String baseUrl;
+
+ @BeforeEach
+ void setUp() throws IOException {
+ server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
+ baseUrl = "http://127.0.0.1:" + server.getAddress().getPort();
+ server.start();
+ }
+
+ @AfterEach
+ void tearDown() {
+ server.stop(0);
+ }
+
+ @Test
+ void queryShouldPreserveSeriesLabelsDecimalValuesAndWarnings() {
+ AtomicReference<String> requestBody = new AtomicReference<>();
+ AtomicReference<String> requestMethod = new AtomicReference<>();
+ AtomicReference<String> contentType = new AtomicReference<>();
+ server.createContext("/api/v1/query_range", exchange -> {
+ requestMethod.set(exchange.getRequestMethod());
+
contentType.set(exchange.getRequestHeaders().getFirst("Content-Type"));
+ requestBody.set(new
String(exchange.getRequestBody().readAllBytes(), StandardCharsets.UTF_8));
+ respond(exchange, 200, """
+ {
+ "status": "success",
+ "data": {
+ "resultType": "matrix",
+ "result": [
+ {
+ "metric": {"node_id": "broker-a", "cluster":
"cluster-a"},
+ "values": [[1784107658, "0.30000000000000004"],
[1784107688, "NaN"]]
+ },
+ {
+ "metric": {"node_id": "broker-b", "cluster":
"cluster-a"},
+ "values": [[1784107658.5, "1.25"]]
+ }
+ ]
+ },
+ "warnings": ["partial response"]
+ }
+ """);
+ });
+
+ PrometheusMetricsSource source = source(Duration.ofSeconds(2));
+ MetricDataVO result = source.query(query());
+
+ assertThat(result.getResultType()).isEqualTo("matrix");
+ assertThat(result.getSeries()).hasSize(2);
+ assertThat(result.getSeries().get(0).getLabels())
+ .containsEntry("node_id", "broker-a")
+ .containsEntry("cluster", "cluster-a");
+ assertThat(result.getSeries().get(0).getValues().get(0).getValue())
+ .isEqualTo("0.30000000000000004");
+
assertThat(result.getSeries().get(0).getValues().get(1).getValue()).isEqualTo("NaN");
+
assertThat(result.getSeries().get(1).getValues().get(0).getTimestamp()).isEqualTo(1784107658.5D);
+ assertThat(result.getWarnings()).containsExactly("partial response");
+
+ assertThat(requestMethod.get()).isEqualTo("POST");
+
assertThat(contentType.get()).startsWith("application/x-www-form-urlencoded");
+ String decodedBody = URLDecoder.decode(requestBody.get(),
StandardCharsets.UTF_8);
+
assertThat(decodedBody).contains("query=sum(rate(rocketmq_messages_in_total[1m]))");
+ assertThat(decodedBody).contains("start=1784107658");
+ assertThat(decodedBody).contains("end=1784108558");
+ assertThat(decodedBody).contains("step=30s");
+ }
+
+ @Test
+ void queryShouldPreserveHistogramOnlySeries() {
+ server.createContext("/api/v1/query_range", exchange ->
respond(exchange, 200, """
+ {
+ "status": "success",
+ "data": {
+ "resultType": "matrix",
+ "result": [{
+ "metric": {"__name__": "rocketmq_rpc_latency"},
+ "histograms": [[1784107658, {
+ "count": "12",
+ "sum": "3.5",
+ "buckets": [[3, "-0.5", "0.5", "4"], [0, "0.5",
"+Inf", "8"]]
+ }]]
+ }]
+ }
+ }
+ """));
+
+ MetricDataVO.MetricSeriesVO series =
source(Duration.ofSeconds(2)).query(query()).getSeries().get(0);
+
+ assertThat(series.getValues()).isEmpty();
+ assertThat(series.getHistograms()).hasSize(1);
+
assertThat(series.getHistograms().get(0).getTimestamp()).isEqualTo(1784107658D);
+
assertThat(series.getHistograms().get(0).getHistogram().path("count").asText()).isEqualTo("12");
+
assertThat(series.getHistograms().get(0).getHistogram().path("sum").asText()).isEqualTo("3.5");
+
assertThat(series.getHistograms().get(0).getHistogram().path("buckets")).hasSize(2);
+ }
+
+ @Test
+ void queryShouldPreserveFloatAndHistogramSamplesInSameSeries() {
+ server.createContext("/api/v1/query_range", exchange ->
respond(exchange, 200, """
+ {
+ "status": "success",
+ "data": {
+ "resultType": "matrix",
+ "result": [{
+ "metric": {"__name__": "request_duration_seconds"},
+ "values": [[1784107658, "1.25"]],
+ "histograms": [[1784107688, {
+ "count": "2",
+ "sum": "1.5",
+ "buckets": [[3, "-0.5", "0.5", "2"]]
+ }]]
+ }]
+ }
+ }
+ """));
+
+ MetricDataVO.MetricSeriesVO series =
source(Duration.ofSeconds(2)).query(query()).getSeries().get(0);
+
+ assertThat(series.getValues()).hasSize(1);
+ assertThat(series.getValues().get(0).getValue()).isEqualTo("1.25");
+ assertThat(series.getHistograms()).hasSize(1);
+
assertThat(series.getHistograms().get(0).getHistogram().path("count").asText()).isEqualTo("2");
+ }
+
+ @Test
+ void queryShouldApplyBasicAuthentication() {
+ AtomicReference<String> authorization = new AtomicReference<>();
+ server.createContext("/api/v1/query_range", exchange -> {
+
authorization.set(exchange.getRequestHeaders().getFirst("Authorization"));
+ respond(exchange, 200, successResponse());
+ });
+ PrometheusProperties properties = properties(Duration.ofSeconds(2));
+ properties.setUsername("studio");
+ properties.setPassword("secret");
+
+ source(properties).query(query());
+
+ String credentials =
Base64.getEncoder().encodeToString("studio:secret".getBytes(StandardCharsets.UTF_8));
+ assertThat(authorization.get()).isEqualTo("Basic " + credentials);
+ }
+
+ @Test
+ void queryShouldPreferBearerAuthentication() {
+ AtomicReference<String> authorization = new AtomicReference<>();
+ server.createContext("/api/v1/query_range", exchange -> {
+
authorization.set(exchange.getRequestHeaders().getFirst("Authorization"));
+ respond(exchange, 200, successResponse());
+ });
+ PrometheusProperties properties = properties(Duration.ofSeconds(2));
+ properties.setBearerToken("test-token");
+ properties.setUsername("ignored-user");
+ properties.setPassword("ignored-password");
+
+ source(properties).query(query());
+
+ assertThat(authorization.get()).isEqualTo("Bearer test-token");
+ }
+
+ @Test
+ void queryShouldExposePrometheusErrorDetails() {
+ server.createContext("/api/v1/query_range", exchange ->
respond(exchange, 422, """
+ {"status":"error","errorType":"execution","error":"invalid
expression"}
+ """));
+
+ PrometheusMetricsSource source = source(Duration.ofSeconds(2));
+
+ assertThatThrownBy(() -> source.query(query()))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> {
+ PrometheusException prometheusException =
(PrometheusException) exception;
+
assertThat(prometheusException.getStatusCode()).isEqualTo(422);
+ assertThat(prometheusException.getMessage())
+ .isEqualTo("Prometheus query failed (execution):
invalid expression");
+ });
+ }
+
+ @Test
+ void queryShouldRejectMalformedPrometheusResponse() {
+ server.createContext("/api/v1/query_range", exchange ->
respond(exchange, 200, """
+ {"status":"success","data":{"resultType":"matrix","result":{}}}
+ """));
+
+ assertThatThrownBy(() -> source(Duration.ofSeconds(2)).query(query()))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> assertThat(((PrometheusException)
exception).getStatusCode())
+ .isEqualTo(HttpStatus.BAD_GATEWAY.value()))
+ .hasMessage("Prometheus returned a malformed response");
+ }
+
+ @Test
+ void queryShouldRejectEndEarlierThanStart() {
+ MetricQueryDTO invalidQuery = MetricQueryDTO.builder()
+ .metric("up")
+ .start(2L)
+ .end(1L)
+ .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 not be earlier than start");
+ }
+
+ @Test
+ void queryShouldFailLoudWhenPrometheusIsNotConfigured() {
+ PrometheusProperties properties = new PrometheusProperties();
+ PrometheusMetricsSource source = new PrometheusMetricsSource(
+ RestClient.builder(), new ObjectMapper(), properties);
+
+ assertThatThrownBy(() -> source.query(query()))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> assertThat(((PrometheusException)
exception).getStatusCode())
+ .isEqualTo(HttpStatus.SERVICE_UNAVAILABLE.value()))
+ .hasMessage("Prometheus base URL is not configured");
+ }
+
+ @Test
+ void queryShouldReportReadTimeout() {
+ server.createContext("/api/v1/query_range", exchange -> {
+ try {
+ Thread.sleep(300);
+ respond(exchange, 200,
"{\"status\":\"success\",\"data\":{\"resultType\":\"matrix\",\"result\":[]}}");
+ } catch (InterruptedException exception) {
+ Thread.currentThread().interrupt();
+ } catch (IOException ignored) {
+ // The client closes the exchange after the expected timeout.
+ }
+ });
+
+ PrometheusMetricsSource source = source(Duration.ofMillis(50));
+
+ assertThatThrownBy(() -> source.query(query()))
+ .isInstanceOf(PrometheusException.class)
+ .satisfies(exception -> {
+ PrometheusException prometheusException =
(PrometheusException) exception;
+ assertThat(prometheusException.getStatusCode())
+ .isEqualTo(HttpStatus.GATEWAY_TIMEOUT.value());
+ })
+ .hasMessage("Prometheus query timed out");
+ }
+
+ private PrometheusMetricsSource source(Duration readTimeout) {
+ return source(properties(readTimeout));
+ }
+
+ private PrometheusMetricsSource source(PrometheusProperties properties) {
+ return new PrometheusMetricsSource(RestClient.builder(), new
ObjectMapper(), properties);
+ }
+
+ private PrometheusProperties properties(Duration readTimeout) {
+ PrometheusProperties properties = new PrometheusProperties();
+ properties.setBaseUrl(baseUrl);
+ properties.setConnectTimeout(Duration.ofSeconds(1));
+ properties.setReadTimeout(readTimeout);
+ return properties;
+ }
+
+ private MetricQueryDTO query() {
+ return MetricQueryDTO.builder()
+ .metric("sum(rate(rocketmq_messages_in_total[1m]))")
+ .start(1784107658L)
+ .end(1784108558L)
+ .step("30s")
+ .build();
+ }
+
+ private void respond(HttpExchange exchange, int statusCode, String body)
throws IOException {
+ byte[] response = body.getBytes(StandardCharsets.UTF_8);
+ exchange.getResponseHeaders().set("Content-Type", "application/json");
+ exchange.sendResponseHeaders(statusCode, response.length);
+ exchange.getResponseBody().write(response);
+ exchange.close();
+ }
+
+ private String successResponse() {
+ return
"{\"status\":\"success\",\"data\":{\"resultType\":\"matrix\",\"result\":[]}}";
+ }
+}
diff --git a/web/src/api/metrics.ts b/web/src/api/metrics.ts
index ddc0cc1..d0b71a9 100644
--- a/web/src/api/metrics.ts
+++ b/web/src/api/metrics.ts
@@ -40,6 +40,41 @@ export interface DashboardData {
clusters: ClusterOverview[];
}
+export interface MetricSample {
+ timestamp: number;
+ value: string;
+}
+
+export interface MetricHistogram {
+ count: string;
+ sum: string;
+ buckets: [number, string, string, string][];
+}
+
+export interface MetricHistogramSample {
+ timestamp: number;
+ histogram: MetricHistogram;
+}
+
+export interface MetricSeries {
+ labels: Record<string, string>;
+ values: MetricSample[];
+ histograms: MetricHistogramSample[];
+}
+
+export interface MetricData {
+ resultType: string;
+ series: MetricSeries[];
+ warnings: string[];
+}
+
+export interface MetricQuery {
+ metric: string;
+ start: number;
+ end: number;
+ step: string;
+}
+
// ─── Dashboard ──────────────────────────────────────────────────
export async function getDashboard() {
const res = await client.get<{ data: DashboardData }>('/dashboard');
@@ -47,12 +82,7 @@ export async function getDashboard() {
}
// ─── Metrics ────────────────────────────────────────────────────
-export async function queryMetrics(query: {
- metric: string;
- start: number;
- end: number;
- step?: string;
-}) {
- const res = await client.post<{ data: unknown }>('/metrics/query', query);
+export async function queryMetrics(query: MetricQuery) {
+ const res = await client.post<{ data: MetricData }>('/metrics/query', query);
return res.data.data;
}