This is an automated email from the ASF dual-hosted git repository. zqr10159 pushed a commit to branch 2.0.0 in repository https://gitbox.apache.org/repos/asf/hertzbeat.git
commit be0d6e5c4e0dcc51035c84aa1c8f2e3932ae0b96 Author: Logic <[email protected]> AuthorDate: Fri Aug 28 15:08:43 2026 +0800 Add secure Perses telemetry query boundary --- .../common/entity/dto/query/DatasourceQuery.java | 3 + ...ticatedGreptimeThreeSignalPublicApiE2eTest.java | 4 +- .../GreptimeThreeSignalInstrumentationE2eTest.java | 4 +- .../PrometheusActiveSourcePublicApiE2eTest.java | 2 +- .../manager/support/GlobalExceptionHandler.java | 9 + .../support/GlobalExceptionHandlerTest.java | 14 + .../controller/OtlpIngestionController.java | 6 +- .../service/OtlpIngestionWorkspaceService.java | 22 + .../impl/OtlpIngestionWorkspaceServiceImpl.java | 82 +++- .../CollectorScopedMetricsQueryServiceImpl.java | 130 +++++- .../controller/OtlpIngestionControllerTest.java | 28 +- .../OtlpIngestionWorkspaceServiceImplTest.java | 21 + ...CollectorScopedMetricsQueryServiceImplTest.java | 178 +++++++- .../warehouse/db/PromqlQueryExecutor.java | 5 + .../repository/MetricQueryRepository.java | 19 +- .../repository/PromqlMetricQueryRepository.java | 4 +- .../db/GreptimePromqlQueryExecutorTest.java | 5 + .../PromqlMetricQueryRepositoryTest.java | 7 +- script/ci/test_official_otel_demo_metrics_poll.py | 47 +++ script/dev/run-official-otel-demo.sh | 17 +- script/dev/verify-otlp-three-signal-demo.sh | 4 +- web-app/scripts/perses-boundary-contract.test.mjs | 4 + .../datasource/hertzbeat-query-client.test.ts | 465 +++++++++++++++++++++ .../perses/datasource/hertzbeat-query-client.ts | 212 ++++++++++ .../perses/datasource/hertzbeat-query-contract.ts | 155 +++++++ .../perses/datasource/hertzbeat-query-schema.ts | 316 ++++++++++++++ web-app/src/platform/perses/index.ts | 23 + 27 files changed, 1723 insertions(+), 63 deletions(-) diff --git a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/dto/query/DatasourceQuery.java b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/dto/query/DatasourceQuery.java index 1f97f46af6..3b7647c518 100644 --- a/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/dto/query/DatasourceQuery.java +++ b/hertzbeat-common-core/src/main/java/org/apache/hertzbeat/common/entity/dto/query/DatasourceQuery.java @@ -56,4 +56,7 @@ public class DatasourceQuery { @Schema(title = "query time step, like 5m or 1h") private String step; + + @Schema(title = "maximum number of series returned by the datasource") + private Integer limit; } diff --git a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/AuthenticatedGreptimeThreeSignalPublicApiE2eTest.java b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/AuthenticatedGreptimeThreeSignalPublicApiE2eTest.java index ce681addda..926ad6cfdd 100644 --- a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/AuthenticatedGreptimeThreeSignalPublicApiE2eTest.java +++ b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/AuthenticatedGreptimeThreeSignalPublicApiE2eTest.java @@ -231,7 +231,7 @@ class AuthenticatedGreptimeThreeSignalPublicApiE2eTest extends GreptimeThreeSign Map<String, String> metricParameters = new LinkedHashMap<>(parameters); metricParameters.put("query", METRIC_QUERY); - metricParameters.put("step", "1s"); + metricParameters.put("step", "1"); metricParameters.put("limit", "20"); JsonNode metrics = authenticatedGet("/api/ingestion/otlp/metrics/console", metricParameters, token); assertThat(metrics.path("stats").path("nonEmptySeries").asInt()).isZero(); @@ -249,7 +249,7 @@ class AuthenticatedGreptimeThreeSignalPublicApiE2eTest extends GreptimeThreeSign private void assertMetricsQuery(JsonNode context, String token) throws Exception { Map<String, String> parameters = commonQueryParameters(context); parameters.put("query", METRIC_QUERY); - parameters.put("step", "1s"); + parameters.put("step", "1"); parameters.put("limit", "20"); JsonNode data = authenticatedGet("/api/ingestion/otlp/metrics/console", parameters, token); assertThat(data.path("context").path("collectorId").asText()).isEqualTo(COLLECTOR_ID); diff --git a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeThreeSignalInstrumentationE2eTest.java b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeThreeSignalInstrumentationE2eTest.java index c0a3367c1d..ece7c87c76 100644 --- a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeThreeSignalInstrumentationE2eTest.java +++ b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/GreptimeThreeSignalInstrumentationE2eTest.java @@ -316,7 +316,7 @@ class GreptimeThreeSignalInstrumentationE2eTest extends GreptimeThreeSignalE2eSu "default", null, null, metricsContext.startedAt(), end, metricsContext.serviceName(), metricsContext.serviceNamespace(), metricsContext.environment(), metricsContext.collectorId(), INSTANCE_ID, ENDPOINT, "hertzbeat_e2e_requests", null, null, - null, null, "1s", "20", null)); + null, null, "1", "20", null)); assertThat(metrics.getContext().getCollectorId()).isEqualTo(metricsContext.collectorId()); assertThat(metrics.getContext().getInstance()).isEqualTo(INSTANCE_ID); assertThat(metrics.getContext().getEndpoint()).isEqualTo(ENDPOINT); @@ -335,7 +335,7 @@ class GreptimeThreeSignalInstrumentationE2eTest extends GreptimeThreeSignalE2eSu "default", null, null, metricsContext.startedAt(), end, metricsContext.serviceName(), metricsContext.serviceNamespace(), metricsContext.environment(), metricsContext.collectorId(), "other-instance", ENDPOINT, "hertzbeat_e2e_requests", null, null, - null, null, "1s", "20", null)); + null, null, "1", "20", null)); assertThat(missingInstanceMetrics.getStats().getNonEmptySeries()).isZero(); org.springframework.data.domain.Page<LogEntry> logs = logQueryService.list( diff --git a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/PrometheusActiveSourcePublicApiE2eTest.java b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/PrometheusActiveSourcePublicApiE2eTest.java index 8c22ee9b3c..efedf343a5 100644 --- a/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/PrometheusActiveSourcePublicApiE2eTest.java +++ b/hertzbeat-e2e/hertzbeat-observability-e2e/src/test/java/org/apache/hertzbeat/observability/storage/PrometheusActiveSourcePublicApiE2eTest.java @@ -136,7 +136,7 @@ class PrometheusActiveSourcePublicApiE2eTest extends GreptimeThreeSignalE2eSuppo parameters.put("instance", context.path("serviceInstanceId").asText()); parameters.put("endpoint", context.path("endpoint").asText()); parameters.put("query", METRIC_QUERY); - parameters.put("step", "1s"); + parameters.put("step", "1"); parameters.put("limit", "20"); await().atMost(Duration.ofSeconds(30)).pollInterval(Duration.ofSeconds(1)).untilAsserted(() -> { JsonNode data = successfulJson(send(get( diff --git a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandler.java b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandler.java index 1dafb467cf..f99f84f857 100644 --- a/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandler.java +++ b/hertzbeat-manager/src/main/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandler.java @@ -29,6 +29,7 @@ import org.apache.hertzbeat.common.transaction.MetadataWriteAdmissionException; import org.apache.hertzbeat.common.entity.dto.Message; import org.apache.hertzbeat.common.support.exception.CommonException; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; +import org.apache.hertzbeat.observability.shared.query.ObservabilityQueryRequestException; import org.apache.hertzbeat.alert.notice.AlertNoticeException; import org.apache.hertzbeat.manager.support.exception.MonitorDatabaseException; import org.apache.hertzbeat.manager.support.exception.MonitorDetectException; @@ -86,6 +87,14 @@ public class GlobalExceptionHandler { .body(Message.fail(FAIL_CODE, TELEMETRY_STORAGE_UNAVAILABLE_MESSAGE)); } + /** Return a stable HTTP error without echoing rejected query content. */ + @ExceptionHandler(ObservabilityQueryRequestException.class) + @ResponseBody + ResponseEntity<Message<Void>> handleObservabilityQueryRequestException() { + return ResponseEntity.badRequest() + .body(Message.fail(PARAM_INVALID_CODE, ObservabilityQueryRequestException.ERROR_CODE)); + } + /** Return a retryable, cache-safe maintenance response without logging private state. */ @ExceptionHandler(MetadataWriteAdmissionException.class) @ResponseBody diff --git a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandlerTest.java b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandlerTest.java index e28ed1ae23..b9ee24236f 100644 --- a/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandlerTest.java +++ b/hertzbeat-manager/src/test/java/org/apache/hertzbeat/manager/support/GlobalExceptionHandlerTest.java @@ -24,6 +24,7 @@ import ch.qos.logback.core.read.ListAppender; import jakarta.servlet.http.HttpServletResponse; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; import org.apache.hertzbeat.common.transaction.MetadataWriteAdmissionException; +import org.apache.hertzbeat.observability.shared.query.ObservabilityQueryRequestException; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.slf4j.LoggerFactory; @@ -139,6 +140,14 @@ class GlobalExceptionHandlerTest { .andExpect(content().string(not(containsString(PRIVATE_SERIALIZATION_DETAIL)))); } + @Test + void invalidObservabilityQueryReturnsStableBadRequest() throws Exception { + mockMvc.perform(MockMvcRequestBuilders.get("/invalid-observability-query")) + .andExpect(status().isBadRequest()) + .andExpect(content().string(containsString(ObservabilityQueryRequestException.ERROR_CODE))) + .andExpect(content().string(not(containsString(PRIVATE_SERIALIZATION_DETAIL)))); + } + @Test void metadataWriteAdmissionFailureIsStableNoStoreServiceUnavailable() throws Exception { mockMvc.perform(MockMvcRequestBuilders.post("/metadata-write-maintenance")) @@ -179,6 +188,11 @@ class GlobalExceptionHandlerTest { throw new TelemetryStorageUnavailableException(); } + @GetMapping("/invalid-observability-query") + void invalidObservabilityQuery() { + throw new ObservabilityQueryRequestException(); + } + @org.springframework.web.bind.annotation.PostMapping("/metadata-write-maintenance") void metadataWriteMaintenance() { throw MetadataWriteAdmissionException.metadataWritesPaused(); diff --git a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionController.java b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionController.java index c52cd033dd..c74d0fb2a6 100644 --- a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionController.java +++ b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionController.java @@ -88,15 +88,15 @@ public class OtlpIngestionController { public ResponseEntity<Message<OtlpMetricsConsoleDto>> metricsConsole( @RequestParam(value = "entityId", required = false) Long entityId, @RequestParam(value = "entityType", required = false) String entityType, - @RequestParam(value = "start", required = false) Long start, - @RequestParam(value = "end", required = false) Long end, + @RequestParam("start") Long start, + @RequestParam("end") Long end, @RequestParam(value = "serviceName", required = false) String serviceName, @RequestParam(value = "serviceNamespace", required = false) String serviceNamespace, @RequestParam(value = "environment", required = false) String environment, @RequestParam(value = "collectorId", required = false) String collectorId, @RequestParam(value = "instance", required = false) String instance, @RequestParam(value = "endpoint", required = false) String endpoint, - @RequestParam(value = "query", required = false) String query, + @RequestParam("query") String query, @RequestParam(value = "filter", required = false) String filter, @RequestParam(value = "groupBy", required = false) String groupBy, @RequestParam(value = "aggregation", required = false) String aggregation, diff --git a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/OtlpIngestionWorkspaceService.java b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/OtlpIngestionWorkspaceService.java index 6f1b71d971..e160785268 100644 --- a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/OtlpIngestionWorkspaceService.java +++ b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/OtlpIngestionWorkspaceService.java @@ -67,6 +67,28 @@ public interface OtlpIngestionWorkspaceService { environment, query, filter, groupBy, aggregation, temporalAggregation, step, limit, operationName); } + /** Execute the public metrics console contract with a datasource-enforced series limit. */ + OtlpMetricsConsoleDto getBoundedMetricsConsole( + String workspaceId, + Long entityId, + String entityType, + Long start, + Long end, + String serviceName, + String serviceNamespace, + String environment, + String collectorId, + String instance, + String endpoint, + String query, + String filter, + String groupBy, + String aggregation, + String temporalAggregation, + String step, + String limit, + String operationName); + OtlpMetricsInventoryDto getMetricsInventory(String workspaceId, Long entityId, String entityType, Long start, Long end, String serviceName, String serviceNamespace, String environment, String limit); diff --git a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImpl.java b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImpl.java index 3dcd7e8c7d..ee5e7d6b96 100644 --- a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImpl.java +++ b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImpl.java @@ -563,6 +563,60 @@ public class OtlpIngestionWorkspaceServiceImpl implements OtlpIngestionWorkspace String step, String limit, String operationName) { + return getMetricsConsole( + workspaceId, entityId, entityType, start, end, serviceName, serviceNamespace, environment, + collectorId, instance, endpoint, query, filter, groupBy, aggregation, temporalAggregation, step, limit, + operationName, false); + } + + @Override + public OtlpMetricsConsoleDto getBoundedMetricsConsole( + String workspaceId, + Long entityId, + String entityType, + Long start, + Long end, + String serviceName, + String serviceNamespace, + String environment, + String collectorId, + String instance, + String endpoint, + String query, + String filter, + String groupBy, + String aggregation, + String temporalAggregation, + String step, + String limit, + String operationName) { + return getMetricsConsole( + workspaceId, entityId, entityType, start, end, serviceName, serviceNamespace, environment, + collectorId, instance, endpoint, query, filter, groupBy, aggregation, temporalAggregation, step, limit, + operationName, true); + } + + private OtlpMetricsConsoleDto getMetricsConsole( + String workspaceId, + Long entityId, + String entityType, + Long start, + Long end, + String serviceName, + String serviceNamespace, + String environment, + String collectorId, + String instance, + String endpoint, + String query, + String filter, + String groupBy, + String aggregation, + String temporalAggregation, + String step, + String limit, + String operationName, + boolean sourceSeriesLimit) { String trustedWorkspaceId = requireMetricsWorkspace(workspaceId); long resolvedEnd = end == null || end <= 0 ? System.currentTimeMillis() : end; long resolvedStart = start == null || start <= 0 || start >= resolvedEnd @@ -643,7 +697,7 @@ public class OtlpIngestionWorkspaceServiceImpl implements OtlpIngestionWorkspace String lastErrorMessage = null; for (String candidateQuery : resolvedQueries) { MetricsQueryExecution execution = executeMetricsConsoleQuery(candidateQuery, resolvedStart, resolvedEnd, - resolvedStep, resolvedSeriesLimit); + resolvedStep, resolvedSeriesLimit, sourceSeriesLimit); if (execution.errorMessage() != null) { lastErrorMessage = execution.errorMessage(); continue; @@ -2418,13 +2472,21 @@ public class OtlpIngestionWorkspaceServiceImpl implements OtlpIngestionWorkspace long resolvedEnd, String resolvedStep, int resolvedSeriesLimit) { - MetricQueryRepository.PromqlRangeQueryResult queryResult = metricQueryRepository.queryPromqlRange( - METRICS_CONSOLE_REF_ID, - query, - resolvedStart, - resolvedEnd, - resolvedStep - ); + return executeMetricsConsoleQuery( + query, resolvedStart, resolvedEnd, resolvedStep, resolvedSeriesLimit, false); + } + + private MetricsQueryExecution executeMetricsConsoleQuery(String query, + long resolvedStart, + long resolvedEnd, + String resolvedStep, + int resolvedSeriesLimit, + boolean sourceSeriesLimit) { + MetricQueryRepository.PromqlRangeQueryResult queryResult = sourceSeriesLimit + ? metricQueryRepository.queryPromqlRange( + METRICS_CONSOLE_REF_ID, query, resolvedStart, resolvedEnd, resolvedStep, resolvedSeriesLimit) + : metricQueryRepository.queryPromqlRange( + METRICS_CONSOLE_REF_ID, query, resolvedStart, resolvedEnd, resolvedStep); if (queryResult == null) { log.warn(MetricQueryRepository.PROMQL_QUERY_FAILED); return new MetricsQueryExecution( @@ -2443,7 +2505,9 @@ public class OtlpIngestionWorkspaceServiceImpl implements OtlpIngestionWorkspace queryResult.errorMessage() ); } - DatasourceQueryData results = limitMetricsConsoleResults(queryResult.results(), resolvedSeriesLimit); + DatasourceQueryData results = sourceSeriesLimit + ? queryResult.results() + : limitMetricsConsoleResults(queryResult.results(), resolvedSeriesLimit); return new MetricsQueryExecution( queryResult.datasource(), results, diff --git a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImpl.java b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImpl.java index 8f7aeb3820..4156f84ae6 100644 --- a/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImpl.java +++ b/hertzbeat-observability/src/main/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImpl.java @@ -17,8 +17,13 @@ package org.apache.hertzbeat.observability.metrics.service.impl; +import java.time.Duration; +import java.util.List; +import java.util.Locale; +import java.util.Set; import java.util.regex.Pattern; import lombok.RequiredArgsConstructor; +import org.apache.hertzbeat.common.entity.dto.query.DatasourceQueryData; import org.apache.hertzbeat.common.observability.dto.metrics.OtlpMetricsConsoleDto; import org.apache.hertzbeat.common.observability.dto.metrics.OtlpMetricsInventoryDto; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; @@ -26,6 +31,7 @@ import org.apache.hertzbeat.observability.ingestion.semantic.OtlpMetricSemanticL import org.apache.hertzbeat.observability.ingestion.semantic.OtlpResourceSemanticAttributes; import org.apache.hertzbeat.observability.ingestion.service.OtlpIngestionWorkspaceService; import org.apache.hertzbeat.observability.metrics.service.CollectorScopedMetricsQueryService; +import org.apache.hertzbeat.observability.shared.query.ObservabilityQueryRequestException; import org.apache.hertzbeat.observability.shared.query.TelemetryQueryContextScope; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; @@ -39,26 +45,39 @@ public class CollectorScopedMetricsQueryServiceImpl implements CollectorScopedMe private static final Pattern SIMPLE_METRIC_NAME = Pattern.compile("[A-Za-z_:][A-Za-z0-9_:]*"); private static final Pattern COLLECTOR_ID = Pattern.compile("[A-Za-z0-9][A-Za-z0-9._:-]{0,127}"); + private static final Duration MAX_TIME_RANGE = Duration.ofDays(1); + private static final int MAX_SERIES = 32; + private static final int MAX_POINTS_PER_SERIES = 1_200; + private static final Set<String> AGGREGATIONS = Set.of("avg", "sum", "min", "max", "count"); + private static final Set<String> TEMPORAL_AGGREGATIONS = Set.of("raw", "rate", "increase", "delta"); private final OtlpIngestionWorkspaceService workspaceService; @Override public OtlpMetricsConsoleDto query(Request request) { String workspaceId = requireWorkspaceId(request.workspaceId()); + long start = requireExactTimeWindow(request.start(), request.end()); + long end = request.end(); String collectorId = normalizeCollectorId(request.collectorId()); TelemetryQueryContextScope queryContextScope = new TelemetryQueryContextScope( request.instance(), request.endpoint()); String query = StringUtils.trimWhitespace(request.query()); - if (StringUtils.hasText(query) && !SIMPLE_METRIC_NAME.matcher(query).matches()) { - return unsupportedQuery(request, collectorId, queryContextScope); + if (!StringUtils.hasText(query) || !SIMPLE_METRIC_NAME.matcher(query).matches()) { + throw new ObservabilityQueryRequestException(); } + String aggregation = normalizeAllowlistedControl(request.aggregation(), "sum", AGGREGATIONS); + String temporalAggregation = normalizeAllowlistedControl( + request.temporalAggregation(), "raw", TEMPORAL_AGGREGATIONS); + String step = resolveEffectiveStep(start, end, request.step()); + String limit = resolveSeriesLimit(request.limit()); String scopedFilter = applyCollectorFilter(request.filter(), collectorId); queryContextScope.validateMetricFilter(scopedFilter); - OtlpMetricsConsoleDto result = workspaceService.getMetricsConsole( - workspaceId, request.entityId(), request.entityType(), request.start(), request.end(), request.serviceName(), + OtlpMetricsConsoleDto result = workspaceService.getBoundedMetricsConsole( + workspaceId, request.entityId(), request.entityType(), start, end, request.serviceName(), request.serviceNamespace(), request.environment(), collectorId, queryContextScope.instance(), - queryContextScope.endpoint(), request.query(), scopedFilter, request.groupBy(), request.aggregation(), - request.temporalAggregation(), request.step(), request.limit(), request.operationName()); + queryContextScope.endpoint(), query, scopedFilter, request.groupBy(), aggregation, + temporalAggregation, step, limit, request.operationName()); + sanitizeAndBoundResponse(result); if (result != null && result.getContext() != null) { result.getContext().setCollectorId(collectorId); result.getContext().setInstance(queryContextScope.instance()); @@ -91,7 +110,7 @@ public class CollectorScopedMetricsQueryServiceImpl implements CollectorScopedMe return null; } if (!COLLECTOR_ID.matcher(normalized).matches()) { - throw new IllegalArgumentException("Collector ID contains unsupported characters"); + throw new ObservabilityQueryRequestException(); } return normalized; } @@ -127,14 +146,93 @@ public class CollectorScopedMetricsQueryServiceImpl implements CollectorScopedMe return normalized; } - private OtlpMetricsConsoleDto unsupportedQuery(Request request, String collectorId, - TelemetryQueryContextScope queryContextScope) { - OtlpMetricsConsoleDto.Context context = new OtlpMetricsConsoleDto.Context( - request.entityId(), request.entityType(), null, request.serviceName(), request.serviceNamespace(), - request.environment(), collectorId, queryContextScope.instance(), queryContextScope.endpoint(), - request.operationName(), request.start(), request.end()); - return new OtlpMetricsConsoleDto( - context, request.query(), null, "promql", null, - new OtlpMetricsConsoleDto.Stats(0, 0, null), "unsupported_query", null); + private long requireExactTimeWindow(Long start, Long end) { + if (start == null || end == null || start <= 0 || end <= start + || end - start > MAX_TIME_RANGE.toMillis()) { + throw new ObservabilityQueryRequestException(); + } + return start; + } + + private String normalizeAllowlistedControl(String value, String defaultValue, Set<String> allowedValues) { + String normalized = StringUtils.trimWhitespace(value); + if (!StringUtils.hasText(normalized)) { + return defaultValue; + } + normalized = normalized.toLowerCase(Locale.ROOT); + if (!allowedValues.contains(normalized)) { + throw new ObservabilityQueryRequestException(); + } + return normalized; + } + + private String resolveEffectiveStep(long start, long end, String requestedStep) { + long minimumStepSeconds = Math.max( + 1L, + Math.ceilDiv(end - start, 1_000L * (MAX_POINTS_PER_SERIES - 1))); + String normalized = StringUtils.trimWhitespace(requestedStep); + long requestedSeconds = defaultStepSeconds(end - start); + if (StringUtils.hasText(normalized)) { + if (!normalized.matches("[1-9]\\d*")) { + throw new ObservabilityQueryRequestException(); + } + try { + requestedSeconds = Long.parseLong(normalized); + } catch (NumberFormatException exception) { + throw new ObservabilityQueryRequestException(); + } + if (requestedSeconds > Duration.ofDays(1).toSeconds()) { + throw new ObservabilityQueryRequestException(); + } + } + return Long.toString(Math.max(requestedSeconds, minimumStepSeconds)); + } + + private long defaultStepSeconds(long rangeMillis) { + if (rangeMillis <= Duration.ofHours(1).toMillis()) { + return 30L; + } + if (rangeMillis <= Duration.ofHours(6).toMillis()) { + return 60L; + } + return 300L; + } + + private String resolveSeriesLimit(String requestedLimit) { + String normalized = StringUtils.trimWhitespace(requestedLimit); + if (!StringUtils.hasText(normalized)) { + return Integer.toString(MAX_SERIES); + } + if (!normalized.matches("[1-9]\\d*")) { + throw new ObservabilityQueryRequestException(); + } + try { + return Integer.toString(Math.min(Integer.parseInt(normalized), MAX_SERIES)); + } catch (NumberFormatException exception) { + throw new ObservabilityQueryRequestException(); + } + } + + private void sanitizeAndBoundResponse(OtlpMetricsConsoleDto result) { + if (result == null) { + return; + } + result.setErrorMessage(null); + if (result.getResults() == null) { + return; + } + DatasourceQueryData results = result.getResults(); + results.setMsg(null); + List<DatasourceQueryData.SchemaData> frames = results.getFrames(); + if (frames == null) { + return; + } + if (frames.size() > MAX_SERIES || frames.stream().anyMatch(this::exceedsPointBudget)) { + throw new TelemetryStorageUnavailableException(); + } + } + + private boolean exceedsPointBudget(DatasourceQueryData.SchemaData frame) { + return frame != null && frame.getData() != null && frame.getData().size() > MAX_POINTS_PER_SERIES; } } diff --git a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionControllerTest.java b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionControllerTest.java index 1f7df02e61..c326c68276 100644 --- a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionControllerTest.java +++ b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/controller/OtlpIngestionControllerTest.java @@ -19,6 +19,7 @@ package org.apache.hertzbeat.observability.ingestion.controller; import static org.mockito.ArgumentMatchers.argThat; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; import java.time.Duration; @@ -252,7 +253,8 @@ class OtlpIngestionControllerTest { null ); when(collectorScopedMetricsQueryService.query(argThat(request -> - "collector-a".equals(request.collectorId()) && "checkout".equals(request.serviceName()) + "team-a".equals(request.workspaceId()) && "collector-a".equals(request.collectorId()) + && "checkout".equals(request.serviceName()) && "checkout-7d9".equals(request.instance()) && "/checkout".equals(request.endpoint())))) .thenReturn(console); @@ -261,6 +263,8 @@ class OtlpIngestionControllerTest { .param("entityType", "service") .param("start", "1000") .param("end", "2000") + .param("query", "http_server_request_duration_count") + .param("workspaceId", "team-b") .param("serviceName", "checkout") .param("serviceNamespace", "commerce") .param("environment", "prod") @@ -284,13 +288,33 @@ class OtlpIngestionControllerTest { .andExpect(jsonPath("$.data.results.frames[0].schema.labels.__name__").value("http_server_requests_seconds_count")); verify(collectorScopedMetricsQueryService).query(argThat(request -> - "collector-a".equals(request.collectorId()) + "team-a".equals(request.workspaceId()) + && "http_server_request_duration_count".equals(request.query()) + && "collector-a".equals(request.collectorId()) && "checkout-7d9".equals(request.instance()) && "/checkout".equals(request.endpoint()) && "span.kind=\"server\"".equals(request.filter()) && "POST /checkout".equals(request.operationName()))); } + @Test + void metricsConsoleRequiresTheExactTypedRangeAndMetricContract() throws Exception { + mockMvc.perform(get("/api/ingestion/otlp/metrics/console") + .param("end", "2000") + .param("query", "http_server_request_duration_count")) + .andExpect(status().isBadRequest()); + mockMvc.perform(get("/api/ingestion/otlp/metrics/console") + .param("start", "1000") + .param("query", "http_server_request_duration_count")) + .andExpect(status().isBadRequest()); + mockMvc.perform(get("/api/ingestion/otlp/metrics/console") + .param("start", "1000") + .param("end", "2000")) + .andExpect(status().isBadRequest()); + + verifyNoInteractions(collectorScopedMetricsQueryService); + } + @Test void shouldReturnWrappedMetricsInventoryPayload() throws Exception { OtlpMetricsInventoryDto inventory = new OtlpMetricsInventoryDto( diff --git a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImplTest.java b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImplTest.java index 9cbe81d56d..080ac6b745 100644 --- a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImplTest.java +++ b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/ingestion/service/impl/OtlpIngestionWorkspaceServiceImplTest.java @@ -1131,6 +1131,27 @@ class OtlpIngestionWorkspaceServiceImplTest { ); } + @Test + void boundedMetricsConsolePushesTheSeriesLimitToTheQueryRepository() { + DatasourceQueryData oversizedQueryData = new DatasourceQueryData( + "otlp-metrics-console", 200, null, metricFrames(33)); + when(metricQueryRepository.hasPromqlExecutor()).thenReturn(true); + when(metricQueryRepository.queryPromqlRange( + eq("otlp-metrics-console"), anyString(), eq(1_000L), eq(2_000L), eq("30s"), eq(32))) + .thenReturn(promqlSuccess(oversizedQueryData)); + + OtlpMetricsConsoleDto console = otlpIngestionWorkspaceService.getBoundedMetricsConsole( + AuthTokenScopes.DEFAULT_WORKSPACE_ID, null, null, 1_000L, 2_000L, + "checkout", "commerce", "prod", null, null, null, + "http_server_request_duration_count", null, null, "sum", "raw", "30", "32", null); + + assertNotNull(console); + assertEquals(33, console.getResults().getFrames().size()); + assertEquals(33, console.getStats().getTotalSeries()); + verify(metricQueryRepository).queryPromqlRange( + eq("otlp-metrics-console"), anyString(), eq(1_000L), eq(2_000L), eq("30s"), eq(32)); + } + @Test void metricsConsoleUsesOperationNameWithHttpRouteFallback() { observabilitySignalIntakeGateway.recordOtlpMetricIntake( diff --git a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImplTest.java b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImplTest.java index 29b9a80d90..c7b1dc3d8f 100644 --- a/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImplTest.java +++ b/hertzbeat-observability/src/test/java/org/apache/hertzbeat/observability/metrics/service/impl/CollectorScopedMetricsQueryServiceImplTest.java @@ -20,14 +20,19 @@ package org.apache.hertzbeat.observability.metrics.service.impl; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; import static org.mockito.Mockito.when; +import java.time.Duration; +import java.util.List; +import org.apache.hertzbeat.common.entity.dto.query.DatasourceQueryData; import org.apache.hertzbeat.common.observability.dto.metrics.OtlpMetricsConsoleDto; import org.apache.hertzbeat.common.observability.dto.metrics.OtlpMetricsInventoryDto; import org.apache.hertzbeat.common.support.exception.TelemetryStorageUnavailableException; import org.apache.hertzbeat.observability.ingestion.service.OtlpIngestionWorkspaceService; import org.apache.hertzbeat.observability.metrics.service.CollectorScopedMetricsQueryService; +import org.apache.hertzbeat.observability.shared.query.ObservabilityQueryRequestException; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -51,10 +56,10 @@ class CollectorScopedMetricsQueryServiceImplTest { void scopesGeneratedMetricQueryThroughCanonicalCollectorLabel() { OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto(); result.setContext(new OtlpMetricsConsoleDto.Context()); - when(workspaceService.getMetricsConsole( + when(workspaceService.getBoundedMetricsConsole( "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", "collector-a", null, null, "http_server_duration", "span_kind=server", - null, null, null, "60s", null, null)).thenReturn(result); + null, "sum", "raw", "60", "32", null)).thenReturn(result); OtlpMetricsConsoleDto actual = service.query(request("collector-a", "http_server_duration")); @@ -85,31 +90,24 @@ class CollectorScopedMetricsQueryServiceImplTest { @Test void scopesDefaultQueryAndPreservesExistingFilter() { - OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto(); - result.setContext(new OtlpMetricsConsoleDto.Context()); - when(workspaceService.getMetricsConsole( - "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", - "collector-east", null, null, null, "span_kind=server", - null, null, null, "60s", null, null)).thenReturn(result); - - OtlpMetricsConsoleDto actual = service.query(request("collector-east", null)); - - assertEquals("collector-east", actual.getContext().getCollectorId()); + assertThrows(ObservabilityQueryRequestException.class, + () -> service.query(request("collector-east", null))); + verifyNoInteractions(workspaceService); } @Test void scopesInstanceAndHttpRouteThroughCanonicalMetricLabels() { OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto(); result.setContext(new OtlpMetricsConsoleDto.Context()); - when(workspaceService.getMetricsConsole( + when(workspaceService.getBoundedMetricsConsole( "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", "collector-a", "checkout-7d9", "/checkout", "http_server_duration", "span_kind=server", - null, null, null, "60s", null, null)).thenReturn(result); + null, "sum", "raw", "60", "32", null)).thenReturn(result); OtlpMetricsConsoleDto actual = service.query(new CollectorScopedMetricsQueryService.Request( "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", "collector-a", "checkout-7d9", "/checkout", "http_server_duration", "span_kind=server", - null, null, null, "60s", null, null)); + null, null, null, "60", null, null)); assertEquals("checkout-7d9", actual.getContext().getInstance()); assertEquals("/checkout", actual.getContext().getEndpoint()); @@ -130,10 +128,10 @@ class CollectorScopedMetricsQueryServiceImplTest { void blankCollectorKeepsLegacyRequestUnchanged() { OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto(); result.setContext(new OtlpMetricsConsoleDto.Context()); - when(workspaceService.getMetricsConsole( + when(workspaceService.getBoundedMetricsConsole( "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", null, null, null, "http_server_duration", "span_kind=server", - null, null, null, "60s", null, null)).thenReturn(result); + null, "sum", "raw", "60", "32", null)).thenReturn(result); OtlpMetricsConsoleDto actual = service.query(request(" ", "http_server_duration")); @@ -142,14 +140,131 @@ class CollectorScopedMetricsQueryServiceImplTest { @Test void rejectsArbitraryPromqlInsteadOfDroppingCollectorScope() { - OtlpMetricsConsoleDto result = service.query(request("collector-a", "sum(rate(http_requests_total[5m]))")); + ObservabilityQueryRequestException failure = assertThrows( + ObservabilityQueryRequestException.class, + () -> service.query(request("collector-a", "sum(rate(http_requests_total[5m]))"))); - assertEquals("unsupported_query", result.getEmptyStateReason()); - assertEquals("collector-a", result.getContext().getCollectorId()); - assertNull(result.getResults()); + assertEquals(ObservabilityQueryRequestException.ERROR_CODE, failure.getMessage()); + verifyNoInteractions(workspaceService); + } + + @Test + void requiresAnExactBoundedTimeWindowBeforeMetricsRead() { + CollectorScopedMetricsQueryService.Request baseline = request("collector-a", "http_server_duration"); + for (CollectorScopedMetricsQueryService.Request invalid : List.of( + withWindow(baseline, null, 200L), + withWindow(baseline, 100L, null), + withWindow(baseline, 0L, 100L), + withWindow(baseline, 200L, 100L), + withWindow(baseline, 100L, 100L), + withWindow(baseline, 100L, 100L + Duration.ofDays(1).toMillis() + 1))) { + ObservabilityQueryRequestException failure = assertThrows( + ObservabilityQueryRequestException.class, () -> service.query(invalid)); + assertEquals(ObservabilityQueryRequestException.ERROR_CODE, failure.getMessage()); + } + verifyNoInteractions(workspaceService); + } + + @Test + void appliesServerOwnedSeriesAndPointBudgets() { + OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto(); + result.setContext(new OtlpMetricsConsoleDto.Context()); + when(workspaceService.getBoundedMetricsConsole( + "team-a", null, null, 1_000L, 86_401_000L, "checkout", "commerce", "prod", + "collector-a", null, null, "http_server_duration", "span_kind=server", + null, "sum", "raw", "73", "32", null)).thenReturn(result); + + service.query(new CollectorScopedMetricsQueryService.Request( + "team-a", null, null, 1_000L, 86_401_000L, "checkout", "commerce", "prod", + "collector-a", null, null, "http_server_duration", "span_kind=server", + null, "SUM", "RAW", "1", "999", null)); + + verify(workspaceService).getBoundedMetricsConsole( + "team-a", null, null, 1_000L, 86_401_000L, "checkout", "commerce", "prod", + "collector-a", null, null, "http_server_duration", "span_kind=server", + null, "sum", "raw", "73", "32", null); + } + + @Test + void rejectsUnallowlistedMetricQueryControls() { + CollectorScopedMetricsQueryService.Request baseline = request("collector-a", "http_server_duration"); + for (CollectorScopedMetricsQueryService.Request invalid : List.of( + withControls(baseline, "sum) by (password) (", null, "60", "20"), + withControls(baseline, "sum", "predict_linear", "60", "20"), + withControls(baseline, "sum", "raw", "60ms", "20"), + withControls(baseline, "sum", "raw", "60", "not-a-number"))) { + assertThrows(ObservabilityQueryRequestException.class, () -> service.query(invalid)); + } verifyNoInteractions(workspaceService); } + @Test + void redactsSuccessfulBackendMessages() { + Object[] row = {1_000L, 1.0}; + DatasourceQueryData.SchemaData frame = new DatasourceQueryData.SchemaData( + new DatasourceQueryData.MetricSchema(List.of(), java.util.Map.of(), java.util.Map.of()), + List.<Object[]>of(row)); + OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto( + new OtlpMetricsConsoleDto.Context(), null, "Greptime-promql", "promql", + new DatasourceQueryData("A", 200, "jdbc:greptime://private?password=secret", List.of(frame)), + new OtlpMetricsConsoleDto.Stats(1, 1, 1_000L), null, "private successful diagnostic"); + when(workspaceService.getBoundedMetricsConsole( + "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", "collector-a", null, null, + "http_server_duration", "span_kind=server", null, "sum", "raw", "60", "32", null)) + .thenReturn(result); + + OtlpMetricsConsoleDto actual = service.query(request("collector-a", "http_server_duration")); + + assertNull(actual.getResults().getMsg()); + assertNull(actual.getErrorMessage()); + } + + @Test + void redactsStorageDiagnosticsWhenTheResultPayloadIsAbsent() { + OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto( + new OtlpMetricsConsoleDto.Context(), null, null, "promql", null, + new OtlpMetricsConsoleDto.Stats(0, 0, null), "load_failed", + "jdbc:greptime://private?password=secret"); + when(workspaceService.getBoundedMetricsConsole( + "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", "collector-a", null, null, + "http_server_duration", "span_kind=server", null, "sum", "raw", "60", "32", null)) + .thenReturn(result); + + OtlpMetricsConsoleDto actual = service.query(request("collector-a", "http_server_duration")); + + assertNull(actual.getErrorMessage()); + assertNull(actual.getResults()); + } + + @Test + void failsClosedWhenTheDatasourceViolatesTheSeriesOrPointBudget() { + Object[] row = {1_000L, 1.0}; + DatasourceQueryData.SchemaData oversizedFrame = new DatasourceQueryData.SchemaData( + new DatasourceQueryData.MetricSchema(List.of(), java.util.Map.of(), java.util.Map.of()), + java.util.stream.IntStream.range(0, 1_201).mapToObj(ignored -> row).toList()); + OtlpMetricsConsoleDto result = new OtlpMetricsConsoleDto( + new OtlpMetricsConsoleDto.Context(), null, "Greptime-promql", "promql", + new DatasourceQueryData("A", 200, null, List.of(oversizedFrame)), + new OtlpMetricsConsoleDto.Stats(1, 1, 1_000L), null, null); + DatasourceQueryData.SchemaData boundedFrame = new DatasourceQueryData.SchemaData( + oversizedFrame.getSchema(), List.<Object[]>of(row)); + OtlpMetricsConsoleDto oversizedSeries = new OtlpMetricsConsoleDto( + new OtlpMetricsConsoleDto.Context(), null, "Greptime-promql", "promql", + new DatasourceQueryData( + "A", 200, null, + java.util.stream.IntStream.range(0, 33).mapToObj(ignored -> boundedFrame).toList()), + new OtlpMetricsConsoleDto.Stats(33, 33, 1_000L), null, null); + when(workspaceService.getBoundedMetricsConsole( + "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", "collector-a", null, null, + "http_server_duration", "span_kind=server", null, "sum", "raw", "60", "32", null)) + .thenReturn(result, oversizedSeries); + + assertThrows(TelemetryStorageUnavailableException.class, + () -> service.query(request("collector-a", "http_server_duration"))); + assertThrows(TelemetryStorageUnavailableException.class, + () -> service.query(request("collector-a", "http_server_duration"))); + } + @Test void rejectsInvalidOrDuplicateCollectorScope() { assertThrows(IllegalArgumentException.class, () -> @@ -189,6 +304,25 @@ class CollectorScopedMetricsQueryServiceImplTest { private CollectorScopedMetricsQueryService.Request request(String collectorId, String query) { return new CollectorScopedMetricsQueryService.Request( "team-a", null, null, 100L, 200L, "checkout", "commerce", "prod", collectorId, null, null, query, - "span_kind=server", null, null, null, "60s", null, null); + "span_kind=server", null, null, null, "60", null, null); + } + + private CollectorScopedMetricsQueryService.Request withWindow( + CollectorScopedMetricsQueryService.Request request, Long start, Long end) { + return new CollectorScopedMetricsQueryService.Request( + request.workspaceId(), request.entityId(), request.entityType(), start, end, request.serviceName(), + request.serviceNamespace(), request.environment(), request.collectorId(), request.instance(), + request.endpoint(), request.query(), request.filter(), request.groupBy(), request.aggregation(), + request.temporalAggregation(), request.step(), request.limit(), request.operationName()); + } + + private CollectorScopedMetricsQueryService.Request withControls( + CollectorScopedMetricsQueryService.Request request, String aggregation, String temporalAggregation, + String step, String limit) { + return new CollectorScopedMetricsQueryService.Request( + request.workspaceId(), request.entityId(), request.entityType(), request.start(), request.end(), + request.serviceName(), request.serviceNamespace(), request.environment(), request.collectorId(), + request.instance(), request.endpoint(), request.query(), request.filter(), request.groupBy(), aggregation, + temporalAggregation, step, limit, request.operationName()); } } diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/db/PromqlQueryExecutor.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/db/PromqlQueryExecutor.java index b970da724c..47c558b31e 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/db/PromqlQueryExecutor.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/db/PromqlQueryExecutor.java @@ -189,6 +189,11 @@ public abstract class PromqlQueryExecutor implements QueryExecutor { } else { throw new IllegalArgumentException(String.format("no such time type for query id %s.", datasourceQuery.getRefId())); } + if (datasourceQuery.getLimit() != null && datasourceQuery.getLimit() > 0) { + uri = UriComponentsBuilder.fromUri(uri) + .queryParam(HTTP_LIMIT_PARAM, datasourceQuery.getLimit()) + .build().toUri(); + } ResponseEntity<PromQlQueryContent> responseEntity = restTemplate.exchange(uri, HttpMethod.GET, httpEntity, PromQlQueryContent.class); if (responseEntity.getStatusCode().is2xxSuccessful()) { diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/MetricQueryRepository.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/MetricQueryRepository.java index 3e136d4e54..30780c26ea 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/MetricQueryRepository.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/MetricQueryRepository.java @@ -46,7 +46,24 @@ public interface MetricQueryRepository { * @param step range step * @return datasource and query result */ - PromqlRangeQueryResult queryPromqlRange(String refId, String query, long start, long end, String step); + default PromqlRangeQueryResult queryPromqlRange( + String refId, String query, long start, long end, String step) { + return queryPromqlRange(refId, query, start, end, step, null); + } + + /** + * Execute a promql range query with a datasource-enforced series limit. + * + * @param refId query ref id + * @param query promql expression + * @param start range start millis + * @param end range end millis + * @param step range step + * @param maxSeries maximum datasource series count, or {@code null} for the repository default + * @return datasource and query result + */ + PromqlRangeQueryResult queryPromqlRange( + String refId, String query, long start, long end, String step, Integer maxSeries); /** * Promql query result wrapper. diff --git a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepository.java b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepository.java index 516de39cd6..ac479d06bf 100644 --- a/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepository.java +++ b/hertzbeat-warehouse/src/main/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepository.java @@ -43,7 +43,8 @@ public class PromqlMetricQueryRepository implements MetricQueryRepository { } @Override - public PromqlRangeQueryResult queryPromqlRange(String refId, String query, long start, long end, String step) { + public PromqlRangeQueryResult queryPromqlRange( + String refId, String query, long start, long end, String step, Integer maxSeries) { QueryExecutor queryExecutor = resolvePromqlExecutor(); if (queryExecutor == null) { return new PromqlRangeQueryResult(null, null, PROMQL_EXECUTOR_UNAVAILABLE); @@ -57,6 +58,7 @@ public class PromqlMetricQueryRepository implements MetricQueryRepository { .start(start) .end(end) .step(step) + .limit(maxSeries) .build(); try { DatasourceQueryData results = queryExecutor.query(datasourceQuery); diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/db/GreptimePromqlQueryExecutorTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/db/GreptimePromqlQueryExecutorTest.java index dbc3a9d5d2..432e09cbfd 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/db/GreptimePromqlQueryExecutorTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/db/GreptimePromqlQueryExecutorTest.java @@ -113,12 +113,17 @@ class GreptimePromqlQueryExecutorTest { .start(1_775_034_288_092L) .end(1_775_037_888_092L) .step("30s") + .limit(32) .build(); DatasourceQueryData result = greptimePromqlQueryExecutor.query(query); assertEquals(200, result.getStatus()); assertEquals(1, result.getFrames().size()); + ArgumentCaptor<URI> uriCaptor = ArgumentCaptor.forClass(URI.class); + verify(restTemplate).exchange( + uriCaptor.capture(), eq(HttpMethod.GET), any(HttpEntity.class), eq(PromQlQueryContent.class)); + assertTrue(uriCaptor.getValue().getQuery().contains("limit=32")); } @Test diff --git a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepositoryTest.java b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepositoryTest.java index 143488b551..35a2b04922 100644 --- a/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepositoryTest.java +++ b/hertzbeat-warehouse/src/test/java/org/apache/hertzbeat/warehouse/repository/PromqlMetricQueryRepositoryTest.java @@ -34,6 +34,7 @@ import org.apache.hertzbeat.common.entity.dto.query.DatasourceQueryData; import org.apache.hertzbeat.warehouse.db.QueryExecutor; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; @@ -63,14 +64,16 @@ class PromqlMetricQueryRepositoryTest { MetricQueryRepository repository = new PromqlMetricQueryRepository(List.of(promqlQueryExecutor)); MetricQueryRepository.PromqlRangeQueryResult result = - repository.queryPromqlRange("ref", "sum(rate(test_total[5m]))", 1000L, 2000L, "30s"); + repository.queryPromqlRange("ref", "sum(rate(test_total[5m]))", 1000L, 2000L, "30s", 32); assertTrue(repository.hasPromqlExecutor()); assertNotNull(result); assertEquals("Greptime-promql", result.datasource()); assertEquals(queryData, result.results()); assertEquals(null, result.errorMessage()); - verify(promqlQueryExecutor).query(any(DatasourceQuery.class)); + ArgumentCaptor<DatasourceQuery> queryCaptor = ArgumentCaptor.forClass(DatasourceQuery.class); + verify(promqlQueryExecutor).query(queryCaptor.capture()); + assertEquals(32, queryCaptor.getValue().getLimit()); } @Test diff --git a/script/ci/test_official_otel_demo_metrics_poll.py b/script/ci/test_official_otel_demo_metrics_poll.py new file mode 100644 index 0000000000..80ff5030eb --- /dev/null +++ b/script/ci/test_official_otel_demo_metrics_poll.py @@ -0,0 +1,47 @@ +#!/usr/bin/env python3 + +# 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. + +"""Contracts for the official OTEL demo metrics verification window.""" + +from __future__ import annotations + +import unittest +from pathlib import Path + + +ROOT = Path(__file__).resolve().parents[2] +DEMO_SCRIPT = ROOT / "script/dev/run-official-otel-demo.sh" + + +class OfficialOtelDemoMetricsPollTest(unittest.TestCase): + + def test_each_metrics_poll_builds_a_fresh_bounded_exact_window(self) -> None: + content = DEMO_SCRIPT.read_text(encoding="utf-8") + self.assertIn("build_metrics_console_path() {", content) + path_builder = content.split("build_metrics_console_path() {", 1)[1].split("\n}", 1)[0] + verify_demo = content.split("verify_demo() {", 1)[1].split("\n}", 1)[0] + metrics_poll = verify_demo.rsplit("poll_until ", 1)[1] + + self.assertIn("date +%s", path_builder) + self.assertIn('metrics_start="$((metrics_end - 3600000))"', path_builder) + self.assertIn("?start=%s&end=%s&query=%s", path_builder) + self.assertIn('"${metrics_start}" "${metrics_end}" "${metrics_query}"', path_builder) + self.assertIn(r'\$(build_metrics_console_path', metrics_poll) + + +if __name__ == "__main__": + unittest.main() diff --git a/script/dev/run-official-otel-demo.sh b/script/dev/run-official-otel-demo.sh index c2b4133c68..e8db2d16ab 100755 --- a/script/dev/run-official-otel-demo.sh +++ b/script/dev/run-official-otel-demo.sh @@ -53,6 +53,7 @@ FLAGD_UI_PORT="${OTEL_DEMO_FLAGD_UI_PORT:-18082}" POLL_INTERVAL_SECONDS="${POLL_INTERVAL_SECONDS:-5}" POLL_ATTEMPTS="${POLL_ATTEMPTS:-30}" CURL_MAX_TIME_SECONDS="${CURL_MAX_TIME_SECONDS:-15}" +OTEL_DEMO_METRIC_QUERY="${OTEL_DEMO_METRIC_QUERY:-http_server_request_duration_seconds_count}" login_token="" @@ -250,10 +251,22 @@ stop_demo_projects() { compose_minimal down --remove-orphans >/dev/null 2>&1 || true } +build_metrics_console_path() { + local metrics_query="$1" + local metrics_end metrics_start + metrics_end="$(($(date +%s) * 1000 + 60000))" + metrics_start="$((metrics_end - 3600000))" + printf '/api/ingestion/otlp/metrics/console?start=%s&end=%s&query=%s' \ + "${metrics_start}" "${metrics_end}" "${metrics_query}" +} + verify_demo() { + local metrics_query if [[ -z "${login_token}" ]]; then login_hertzbeat fi + metrics_query="$(python3 -c 'import sys, urllib.parse; print(urllib.parse.quote(sys.argv[1], safe=""))' \ + "${OTEL_DEMO_METRIC_QUERY}")" log_step "验证 OTLP 概览" poll_until "OTLP 三大信号已激活" \ @@ -274,7 +287,7 @@ verify_demo() { jq -e '.code == 0 and ((.data.content // []) | map(select(.serviceName == \"frontend\" or .serviceName == \"checkout\" or .serviceName == \"cart\" or .serviceName == \"product-catalog\" or .serviceName == \"image-provider\" or .serviceName == \"flagd\")) | length) >= 1' <<<\"\$response\" >/dev/null" poll_until "指标工作台已解析到 demo 服务上下文" \ - "response=\$(api_get '/api/ingestion/otlp/metrics/console' '${login_token}'); \ + "response=\$(api_get \"\$(build_metrics_console_path '${metrics_query}')\" '${login_token}'); \ jq -e '.code == 0 and .data.emptyStateReason != \"no_context\" and ((.data.context.serviceName // \"\") | length) > 0 and ((.data.query // \"\") | length) > 0' <<<\"\$response\" >/dev/null" cat <<EOF @@ -338,6 +351,7 @@ EOF cmd_up() { require_bin curl require_bin jq + require_bin python3 require_bin git require_bin docker @@ -378,6 +392,7 @@ cmd_logs() { cmd_verify() { require_bin curl require_bin jq + require_bin python3 verify_demo } diff --git a/script/dev/verify-otlp-three-signal-demo.sh b/script/dev/verify-otlp-three-signal-demo.sh index 178c427cff..bf957269f6 100755 --- a/script/dev/verify-otlp-three-signal-demo.sh +++ b/script/dev/verify-otlp-three-signal-demo.sh @@ -523,12 +523,14 @@ service_namespace_q="$(url_encode "${SERVICE_NAMESPACE}")" environment_q="$(url_encode "${DEPLOYMENT_ENVIRONMENT}")" entity_type_q="$(url_encode "${HERTZBEAT_ENTITY_TYPE}")" metric_query_q="$(url_encode "${METRIC_QUERY}")" +metrics_end_ms="$(($(date +%s) * 1000 + 60000))" +metrics_start_ms="$((metrics_end_ms - 3600000))" service_version_group_q="$(url_encode "resource:service.version")" host_group_q="$(url_encode "resource:host.name")" k8s_pod_group_q="$(url_encode "resource:k8s.pod.name")" trace_id_q="$(url_encode "${TRACE_ID}")" root_span_id_q="$(url_encode "${ROOT_SPAN_ID}")" -metrics_console_path="/api/ingestion/otlp/metrics/console?entityId=${HERTZBEAT_ENTITY_ID}&entityType=${entity_type_q}&serviceName=${service_name_q}&serviceNamespace=${service_namespace_q}&environment=${environment_q}&query=${metric_query_q}" +metrics_console_path="/api/ingestion/otlp/metrics/console?start=${metrics_start_ms}&end=${metrics_end_ms}&entityId=${HERTZBEAT_ENTITY_ID}&entityType=${entity_type_q}&serviceName=${service_name_q}&serviceNamespace=${service_namespace_q}&environment=${environment_q}&query=${metric_query_q}" metrics_breakout_path="${metrics_console_path}&groupBy=service.version" metrics_resource_breakout_path="${metrics_console_path}&groupBy=service.version,host.name,k8s.pod.name" related_metrics_filter_q="$(url_encode "host.name=\"${HOST_NAME}\" and k8s.namespace.name=\"${K8S_NAMESPACE_NAME}\" and k8s.pod.name=\"${K8S_POD_NAME}\" and container.name=\"${CONTAINER_NAME}\"")" diff --git a/web-app/scripts/perses-boundary-contract.test.mjs b/web-app/scripts/perses-boundary-contract.test.mjs index 8c1f014660..1b880b3659 100644 --- a/web-app/scripts/perses-boundary-contract.test.mjs +++ b/web-app/scripts/perses-boundary-contract.test.mjs @@ -18,6 +18,10 @@ test('Perses has one public platform boundary with the required runtime ownershi assert.ok(existsSync(join(persesRoot, 'index.ts')), 'src/platform/perses/index.ts must own the public API'); assert.ok(existsSync(join(persesRoot, 'runtime')), 'src/platform/perses/runtime must own React integration'); assert.ok(existsSync(join(persesRoot, 'plugins')), 'src/platform/perses/plugins must own plugin registration'); + assert.ok( + existsSync(join(persesRoot, 'datasource')), + 'src/platform/perses/datasource must own HertzBeat API queries' + ); }); test('production code imports Perses packages only inside the platform boundary', () => { diff --git a/web-app/src/platform/perses/datasource/hertzbeat-query-client.test.ts b/web-app/src/platform/perses/datasource/hertzbeat-query-client.test.ts new file mode 100644 index 0000000000..94b1013c84 --- /dev/null +++ b/web-app/src/platform/perses/datasource/hertzbeat-query-client.test.ts @@ -0,0 +1,465 @@ +/* + * 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. + */ + +import { beforeEach, describe, expect, it, vi } from 'vitest'; + +import { ApiMessageError, apiMessageGet } from '@/core/http/api-message'; + +import { HERTZBEAT_QUERY_LIMITS, queryHertzBeatData } from './hertzbeat-query-client'; + +vi.mock('@/core/http/api-message', async importOriginal => { + const actual = await importOriginal<typeof import('@/core/http/api-message')>(); + return { ...actual, apiMessageGet: vi.fn() }; +}); + +const request = vi.mocked(apiMessageGet); +const timeWindow = { from: 1_000, to: 2_000 } as const; +const context = { + entityId: '42', + entityType: 'service', + serviceName: 'checkout', + serviceNamespace: 'commerce', + environment: 'prod', + collectorId: 'collector-a', + instance: 'checkout-01', + endpoint: '/orders' +} as const; + +describe('HertzBeat Perses query client', () => { + beforeEach(() => { + request.mockReset(); + }); + + it('queries a simple metric through the authenticated same-origin console path', async () => { + const signal = new AbortController().signal; + request.mockResolvedValue(metricConsole()); + + const result = await queryHertzBeatData( + { + signal: 'metrics', + queryKind: 'time-series', + timeWindow, + context, + metric: { + name: 'http_server_duration_seconds', + aggregation: 'avg', + temporalAggregation: 'rate', + stepSeconds: 15, + operationName: 'POST /orders' + }, + limit: 20 + }, + { signal } + ); + + const [path, options] = request.mock.calls[0] ?? []; + expect(path).toMatch(/^\/api\/ingestion\/otlp\/metrics\/console\?/u); + const params = new URL(path as string, 'https://hertzbeat.local').searchParams; + expect(Object.fromEntries(params)).toEqual({ + entityId: '42', + entityType: 'service', + start: '1000', + end: '2000', + serviceName: 'checkout', + serviceNamespace: 'commerce', + environment: 'prod', + collectorId: 'collector-a', + instance: 'checkout-01', + endpoint: '/orders', + query: 'http_server_duration_seconds', + aggregation: 'avg', + temporalAggregation: 'rate', + step: '15', + limit: '20', + operationName: 'POST /orders' + }); + expect(options).toEqual({ signal }); + expect(path).not.toMatch(/workspace|sql|promql|secret|greptime/iu); + expect(result).toEqual({ + state: 'ready', + data: { + timeWindow, + source: 'Greptime-promql', + series: [ + { + key: 'http_server_duration_seconds-0', + name: 'http_server_duration_seconds', + unit: 'seconds', + labels: { __name__: 'http_server_duration_seconds', service_name: 'checkout' }, + points: [ + { timestamp: 1_000, value: 12 }, + { timestamp: 2_000, value: 14 } + ] + } + ] + }, + truncated: 'unknown' + }); + }); + + it('keeps a valid empty metric response distinct from fake zero evidence', async () => { + request.mockResolvedValue(metricConsole({ frames: [], totalSeries: 0, nonEmptySeries: 0 })); + + await expect( + queryHertzBeatData({ + signal: 'metrics', + queryKind: 'time-series', + timeWindow, + metric: { name: 'http_requests_total' } + }) + ).resolves.toEqual({ state: 'empty', truncated: false }); + }); + + it('rejects metric frame and statistics contradictions instead of reporting false empty evidence', async () => { + request + .mockResolvedValueOnce(metricConsole({ frames: [], totalSeries: 1, nonEmptySeries: 0 })) + .mockResolvedValueOnce(metricConsole({ frames: [metricFrame()], totalSeries: 0, nonEmptySeries: 0 })); + const query = { + signal: 'metrics', + queryKind: 'time-series', + timeWindow, + metric: { name: 'up' } + } as const; + + await expect(queryHertzBeatData(query)).resolves.toMatchObject({ + state: 'error', + error: { kind: 'contract_error' } + }); + await expect(queryHertzBeatData(query)).resolves.toMatchObject({ + state: 'error', + error: { kind: 'contract_error' } + }); + }); + + it('queries bounded log and trace tables through their typed endpoints', async () => { + request + .mockResolvedValueOnce({ content: [logRow()], totalElements: 3, pageIndex: 0, pageSize: 2 }) + .mockResolvedValueOnce({ content: [traceRow()], totalElements: 1, totalPages: 1, number: 0, size: 2 }); + + const logs = await queryHertzBeatData({ + signal: 'logs', + queryKind: 'table', + timeWindow, + context, + search: 'checkout failed', + severity: 'ERROR', + traceId: 'trace-1', + hideInternal: true, + limit: 2 + }); + const traces = await queryHertzBeatData({ + signal: 'traces', + queryKind: 'table', + timeWindow, + context, + operationName: 'POST /orders', + errorOnly: true, + spanScope: 'entrypoint', + limit: 2 + }); + + expect(request.mock.calls.map(([path]) => path)).toEqual([ + '/api/logs/list?entityId=42&entityType=service&start=1000&end=2000&serviceName=checkout&serviceNamespace=commerce&environment=prod&collectorId=collector-a&instance=checkout-01&endpoint=%2Forders&pageIndex=0&pageSize=2&search=checkout+failed&severityText=ERROR&traceId=trace-1&hideInternal=true', + '/api/traces/list?entityId=42&entityType=service&start=1000&end=2000&serviceName=checkout&serviceNamespace=commerce&environment=prod&collectorId=collector-a&instance=checkout-01&endpoint=%2Forders&pageIndex=0&pageSize=2&operationName=POST+%2Forders&errorOnly=true&spanScope=entrypoint' + ]); + expect(logs).toMatchObject({ state: 'ready', truncated: true, data: { rows: [{ severityText: 'ERROR' }] } }); + expect(traces).toMatchObject({ state: 'ready', truncated: false, data: { rows: [{ traceId: 'trace-1' }] } }); + }); + + it('rejects an empty first page when log or trace totals claim missing rows', async () => { + request + .mockResolvedValueOnce({ content: [], totalElements: 1, pageIndex: 0, pageSize: 2 }) + .mockResolvedValueOnce({ content: [], totalElements: 1, totalPages: 1, number: 0, size: 2 }); + + const logs = await queryHertzBeatData({ signal: 'logs', queryKind: 'table', timeWindow, limit: 2 }); + const traces = await queryHertzBeatData({ signal: 'traces', queryKind: 'table', timeWindow, limit: 2 }); + + expect(logs).toMatchObject({ state: 'error', error: { kind: 'contract_error' } }); + expect(traces).toMatchObject({ state: 'error', error: { kind: 'contract_error' } }); + }); + + it('loads one trace gantt primitive without exposing a free-form endpoint', async () => { + request.mockResolvedValue(traceDetail()); + + const result = await queryHertzBeatData({ + signal: 'traces', + queryKind: 'gantt', + timeWindow, + context, + traceId: 'trace-1', + spanId: 'span-2' + }); + + expect(request.mock.calls[0]?.[0]).toBe( + '/api/traces/trace-1?entityId=42&start=1000&end=2000&serviceName=checkout&serviceNamespace=commerce&environment=prod&collectorId=collector-a&instance=checkout-01&endpoint=%2Forders&spanId=span-2' + ); + expect(result).toMatchObject({ + state: 'ready', + truncated: false, + data: { traceId: 'trace-1', spans: [{ spanId: 'span-1' }, { spanId: 'span-2' }] } + }); + }); + + it('rejects unbounded or transport-shaped input before issuing a request', async () => { + const invalid = [ + { signal: 'metrics', queryKind: 'time-series', timeWindow: { from: 2_000, to: 1_000 }, metric: { name: 'up' } }, + { signal: 'metrics', queryKind: 'time-series', timeWindow, metric: { name: 'sum(rate(up[5m]))' } }, + { + signal: 'logs', + queryKind: 'table', + timeWindow, + limit: HERTZBEAT_QUERY_LIMITS.tableRows + 1, + url: 'https://greptime.invalid', + sql: 'select * from secrets' + } + ]; + + for (const candidate of invalid) { + await expect(queryHertzBeatData(candidate as never)).resolves.toMatchObject({ + state: 'error', + error: { kind: 'invalid_request', retryable: false } + }); + } + expect(request).not.toHaveBeenCalled(); + }); + + it('accepts the full positive Java Long entity id range and rejects overflow', async () => { + request.mockResolvedValue(metricConsole({ entityId: Number('9223372036854775807') })); + const query = { + signal: 'metrics', + queryKind: 'time-series', + timeWindow, + metric: { name: 'up' } + } as const; + + const accepted = await queryHertzBeatData({ + ...query, + context: { entityId: '9223372036854775807' } + }); + const rejected = await queryHertzBeatData({ + ...query, + context: { entityId: '9223372036854775808' } + }); + + expect(accepted.state).toBe('ready'); + expect(request.mock.calls[0]?.[0]).toContain('entityId=9223372036854775807'); + expect(rejected).toMatchObject({ state: 'error', error: { kind: 'invalid_request' } }); + expect(request).toHaveBeenCalledTimes(1); + }); + + it('rejects metric responses that exceed the defensive series or point budgets', async () => { + request + .mockResolvedValueOnce(metricConsole({ frames: Array.from({ length: 33 }, () => metricFrame()) })) + .mockResolvedValueOnce( + metricConsole({ + frames: [ + { + ...metricFrame(), + data: Array.from({ length: HERTZBEAT_QUERY_LIMITS.metricPointsPerSeries + 1 }, (_, index) => [ + index, + index + ]) + } + ] + }) + ); + const query = { + signal: 'metrics', + queryKind: 'time-series', + timeWindow, + metric: { name: 'up' } + } as const; + + await expect(queryHertzBeatData(query)).resolves.toMatchObject({ + state: 'error', + error: { kind: 'contract_error', messageKey: 'perses.query.contract' } + }); + await expect(queryHertzBeatData(query)).resolves.toMatchObject({ + state: 'error', + error: { kind: 'contract_error', messageKey: 'perses.query.contract' } + }); + }); + + it('returns sanitized typed failures and never exposes backend diagnostics', async () => { + request + .mockRejectedValueOnce(new ApiMessageError('permission sql=SELECT secret', { status: 403 })) + .mockRejectedValueOnce(new ApiMessageError('capacity query=private_metric', { status: 429 })) + .mockResolvedValueOnce({ unexpected: 'private payload' }); + + const query = { + signal: 'metrics', + queryKind: 'time-series', + timeWindow, + metric: { name: 'up' } + } as const; + const permission = await queryHertzBeatData(query); + const overloaded = await queryHertzBeatData(query); + const contract = await queryHertzBeatData(query); + + expect(permission).toEqual({ + state: 'error', + error: { + kind: 'permission', + messageKey: 'perses.query.permission', + retryable: false + } + }); + expect(overloaded).toEqual({ + state: 'error', + error: { kind: 'overloaded', messageKey: 'perses.query.overloaded', retryable: true } + }); + expect(contract).toEqual({ + state: 'error', + error: { + kind: 'contract_error', + messageKey: 'perses.query.contract', + retryable: false + } + }); + expect(JSON.stringify([permission, overloaded, contract])).not.toMatch( + /SELECT|secret|private_metric|private payload/u + ); + }); + + it('propagates caller cancellation instead of converting it into an unavailable state', async () => { + const controller = new AbortController(); + request.mockImplementation( + (_path, options) => + new Promise((_resolve, reject) => { + options?.signal?.addEventListener('abort', () => reject(new ApiMessageError('cancelled'))); + }) + ); + const pending = queryHertzBeatData( + { signal: 'metrics', queryKind: 'time-series', timeWindow, metric: { name: 'up' } }, + { signal: controller.signal } + ); + + controller.abort(new DOMException('Cancelled', 'AbortError')); + + await expect(pending).rejects.toMatchObject({ name: 'AbortError' }); + expect(request.mock.calls[0]?.[1]).toEqual({ signal: controller.signal }); + }); +}); + +function metricConsole( + options: { entityId?: number; frames?: unknown[]; totalSeries?: number; nonEmptySeries?: number } = {} +) { + return { + context: { + entityId: options.entityId ?? 42, + entityType: 'service', + entityName: 'Checkout API', + serviceName: 'checkout', + serviceNamespace: 'commerce', + environment: 'prod', + collectorId: 'collector-a', + instance: 'checkout-01', + endpoint: '/orders', + operationName: 'POST /orders', + start: 1_000, + end: 2_000 + }, + query: 'http_server_duration_seconds', + datasource: 'Greptime-promql', + queryMode: 'promql', + results: { + refId: 'otlp-metrics-console', + status: 200, + msg: null, + frames: options.frames ?? [metricFrame()] + }, + stats: { + totalSeries: options.totalSeries ?? 1, + nonEmptySeries: options.nonEmptySeries ?? 1, + latestObservedAt: 2_000 + }, + emptyStateReason: null, + errorMessage: null + }; +} + +function metricFrame() { + return { + schema: { + fields: [ + { name: '__ts__', type: 'time', unit: null }, + { name: '__value__', type: 'number', unit: 'seconds' } + ], + labels: { __name__: 'http_server_duration_seconds', service_name: 'checkout' }, + meta: {} + }, + data: [ + [1_000, 12], + [2_000, 14] + ] + }; +} + +function logRow() { + return { + timeUnixNano: 1_000_000_000, + observedTimeUnixNano: 1_000_000_000, + severityNumber: 17, + severityText: 'ERROR', + body: 'checkout failed', + attributes: {}, + droppedAttributesCount: 0, + traceId: 'trace-1', + spanId: 'span-1', + traceFlags: 1, + resource: { 'service.name': 'checkout' }, + resourceSchemaUrl: null, + instrumentationScope: null, + scopeSchemaUrl: null + }; +} + +function traceRow() { + return { + traceId: 'trace-1', + rootSpanId: 'span-1', + serviceName: 'checkout', + serviceNamespace: 'commerce', + rootSpanName: 'POST /orders', + durationNanos: 10_000_000, + status: 'ERROR', + startTime: 1_000, + errorSpanCount: 1, + resourceAttributes: { 'service.name': 'checkout' } + }; +} + +function traceDetail() { + return { + ...traceRow(), + spans: [traceSpan('span-1', null, false), traceSpan('span-2', 'span-1', true)] + }; +} + +function traceSpan(spanId: string, parentSpanId: string | null, highlighted: boolean) { + return { + traceId: 'trace-1', + spanId, + parentSpanId, + spanName: spanId === 'span-1' ? 'POST /orders' : 'SELECT cart', + serviceName: 'checkout', + status: highlighted ? 'ERROR' : 'OK', + spanKind: highlighted ? 'CLIENT' : 'SERVER', + statusMessage: null, + traceState: null, + scopeName: 'checkout', + scopeVersion: '1.0.0', + durationNanos: 5_000_000, + startTime: 1_000, + highlighted, + resourceAttributes: {}, + spanAttributes: {}, + events: [], + links: [], + codeNavigationHint: null + }; +} diff --git a/web-app/src/platform/perses/datasource/hertzbeat-query-client.ts b/web-app/src/platform/perses/datasource/hertzbeat-query-client.ts new file mode 100644 index 0000000000..a9e29f7db2 --- /dev/null +++ b/web-app/src/platform/perses/datasource/hertzbeat-query-client.ts @@ -0,0 +1,212 @@ +/* + * 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. + */ + +import { ApiMessageError, apiMessageGet } from '@/core/http/api-message'; + +import { + hertzBeatQuerySchema, + type HertzBeatLogTableQuery, + type HertzBeatMetricQuery, + type HertzBeatQuery, + type HertzBeatQueryFailure, + type HertzBeatQueryOutcome, + type HertzBeatTraceGanttQuery, + type HertzBeatTraceTableQuery +} from './hertzbeat-query-contract'; +import { + HertzBeatResponseContractError, + HertzBeatResponseStateError, + parseLogTable, + parseMetricResponse, + parseTraceGantt, + parseTraceTable, + type HertzBeatLogRow, + type HertzBeatMetricData, + type HertzBeatTableData, + type HertzBeatTraceDetail, + type HertzBeatTraceRow +} from './hertzbeat-query-schema'; + +export { HERTZBEAT_QUERY_LIMITS } from './hertzbeat-query-contract'; + +const DEFAULT_TABLE_LIMIT = 100; + +type QueryOptions = { signal?: AbortSignal | undefined }; + +export function queryHertzBeatData( + query: HertzBeatMetricQuery, + options?: QueryOptions +): Promise<HertzBeatQueryOutcome<HertzBeatMetricData>>; +export function queryHertzBeatData( + query: HertzBeatLogTableQuery, + options?: QueryOptions +): Promise<HertzBeatQueryOutcome<HertzBeatTableData<HertzBeatLogRow>>>; +export function queryHertzBeatData( + query: HertzBeatTraceTableQuery, + options?: QueryOptions +): Promise<HertzBeatQueryOutcome<HertzBeatTableData<HertzBeatTraceRow>>>; +export function queryHertzBeatData( + query: HertzBeatTraceGanttQuery, + options?: QueryOptions +): Promise<HertzBeatQueryOutcome<HertzBeatTraceDetail>>; +export async function queryHertzBeatData( + query: HertzBeatQuery, + options: QueryOptions = {} +): Promise<HertzBeatQueryOutcome<unknown>> { + const parsed = hertzBeatQuerySchema.safeParse(query); + if (!parsed.success) return failure('invalid_request'); + try { + return await executeQuery(parsed.data, options.signal); + } catch (error) { + if (options.signal?.aborted) throw options.signal.reason ?? new DOMException('Aborted', 'AbortError'); + return mapFailure(error); + } +} + +export const hertzBeatPersesQueryClient = { query: queryHertzBeatData }; + +async function executeQuery(query: HertzBeatQuery, signal?: AbortSignal): Promise<HertzBeatQueryOutcome<unknown>> { + if (query.signal === 'metrics') { + const data = parseMetricResponse(await request(buildMetricPath(query), signal), query.timeWindow); + return data ? { state: 'ready', data, truncated: 'unknown' } : { state: 'empty', truncated: false }; + } + if (query.signal === 'logs') { + const limit = query.limit ?? DEFAULT_TABLE_LIMIT; + const data = parseLogTable(await request(buildLogTablePath(query, limit), signal), limit); + return data + ? { state: 'ready', data, truncated: data.total > data.rows.length } + : { state: 'empty', truncated: false }; + } + if (query.queryKind === 'table') { + const limit = query.limit ?? DEFAULT_TABLE_LIMIT; + const data = parseTraceTable(await request(buildTraceTablePath(query, limit), signal), limit); + return data + ? { state: 'ready', data, truncated: data.total > data.rows.length } + : { state: 'empty', truncated: false }; + } + const data = parseTraceGantt(await request(buildTraceGanttPath(query), signal), query.traceId); + return data ? { state: 'ready', data, truncated: false } : { state: 'empty', truncated: false }; +} + +function request(path: string, signal?: AbortSignal) { + return apiMessageGet(path, { signal: signal ?? null }); +} + +function buildMetricPath(query: HertzBeatMetricQuery) { + const params = baseParams(query, true); + params.set('query', query.metric.name); + set(params, 'aggregation', query.metric.aggregation); + set(params, 'temporalAggregation', query.metric.temporalAggregation); + setNumber(params, 'step', query.metric.stepSeconds); + setNumber(params, 'limit', query.limit); + set(params, 'operationName', query.metric.operationName); + return `/api/ingestion/otlp/metrics/console?${params.toString()}`; +} + +function buildLogTablePath(query: HertzBeatLogTableQuery, limit: number) { + const params = baseParams(query, true); + params.set('pageIndex', '0'); + params.set('pageSize', String(limit)); + set(params, 'search', query.search); + set(params, 'severityText', query.severity); + set(params, 'traceId', query.traceId); + set(params, 'spanId', query.spanId); + setBoolean(params, 'hideInternal', query.hideInternal); + setBoolean(params, 'hideNoise', query.hideNoise); + return `/api/logs/list?${params.toString()}`; +} + +function buildTraceTablePath(query: HertzBeatTraceTableQuery, limit: number) { + const params = baseParams(query, true); + params.set('pageIndex', '0'); + params.set('pageSize', String(limit)); + set(params, 'traceId', query.traceId); + set(params, 'operationName', query.operationName); + setBoolean(params, 'errorOnly', query.errorOnly); + setNumber(params, 'minDurationMs', query.minDurationMs); + setNumber(params, 'maxDurationMs', query.maxDurationMs); + set(params, 'spanScope', query.spanScope); + setBoolean(params, 'hideInternal', query.hideInternal); + return `/api/traces/list?${params.toString()}`; +} + +function buildTraceGanttPath(query: HertzBeatTraceGanttQuery) { + const params = baseParams(query, false); + set(params, 'spanId', query.spanId); + setNumber(params, 'minDurationMs', query.minDurationMs); + setNumber(params, 'maxDurationMs', query.maxDurationMs); + return `/api/traces/${encodeURIComponent(query.traceId)}?${params.toString()}`; +} + +function baseParams(query: HertzBeatQuery, includeEntityType: boolean) { + const params = new URLSearchParams(); + set(params, 'entityId', query.context?.entityId); + if (includeEntityType) set(params, 'entityType', query.context?.entityType); + params.set('start', String(query.timeWindow.from)); + params.set('end', String(query.timeWindow.to)); + set(params, 'serviceName', query.context?.serviceName); + set(params, 'serviceNamespace', query.context?.serviceNamespace); + set(params, 'environment', query.context?.environment); + set(params, 'collectorId', query.context?.collectorId); + set(params, 'instance', query.context?.instance); + set(params, 'endpoint', query.context?.endpoint); + return params; +} + +function set(params: URLSearchParams, key: string, value: string | undefined) { + if (value) params.set(key, value); +} + +function setNumber(params: URLSearchParams, key: string, value: number | undefined) { + if (value != null) params.set(key, String(value)); +} + +function setBoolean(params: URLSearchParams, key: string, value: boolean | undefined) { + if (value) params.set(key, 'true'); +} + +function mapFailure(error: unknown): HertzBeatQueryOutcome<never> { + if (error instanceof HertzBeatResponseContractError) return failure('contract_error'); + if (error instanceof HertzBeatResponseStateError) return failure(error.kind); + if (error instanceof ApiMessageError) { + if (error.status === 401 || error.status === 403) return failure('permission'); + if (error.status === 400 || error.status === 422 || error.code === 3) return failure('invalid_request'); + if (error.status === 429) return failure('overloaded'); + } + return failure('unavailable'); +} + +function failure(kind: HertzBeatQueryFailure['kind']): HertzBeatQueryOutcome<never> { + const errors: Record<HertzBeatQueryFailure['kind'], HertzBeatQueryFailure> = { + invalid_request: { + kind: 'invalid_request', + messageKey: 'perses.query.invalid', + retryable: false + }, + permission: { + kind: 'permission', + messageKey: 'perses.query.permission', + retryable: false + }, + overloaded: { + kind: 'overloaded', + messageKey: 'perses.query.overloaded', + retryable: true + }, + unavailable: { + kind: 'unavailable', + messageKey: 'perses.query.unavailable', + retryable: true + }, + contract_error: { + kind: 'contract_error', + messageKey: 'perses.query.contract', + retryable: false + } + }; + return { state: 'error', error: errors[kind] }; +} diff --git a/web-app/src/platform/perses/datasource/hertzbeat-query-contract.ts b/web-app/src/platform/perses/datasource/hertzbeat-query-contract.ts new file mode 100644 index 0000000000..c54de95ddc --- /dev/null +++ b/web-app/src/platform/perses/datasource/hertzbeat-query-contract.ts @@ -0,0 +1,155 @@ +/* + * 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. + */ + +import { z } from 'zod'; + +const boundedText = z.string().trim().min(1).max(512); +const identifier = z.string().trim().min(1).max(256); +const JAVA_LONG_MAX = '9223372036854775807'; +const entityId = z + .string() + .regex(/^[1-9]\d{0,18}$/u) + .refine(value => value.length < JAVA_LONG_MAX.length || value <= JAVA_LONG_MAX); + +const contextSchema = z + .object({ + entityId: entityId.optional(), + entityType: z + .string() + .regex(/^[A-Za-z0-9_.:-]{1,128}$/u) + .optional(), + serviceName: identifier.optional(), + serviceNamespace: identifier.optional(), + environment: identifier.optional(), + collectorId: z + .string() + .regex(/^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/u) + .optional(), + instance: identifier.optional(), + endpoint: boundedText.optional() + }) + .strict(); + +const timeWindowSchema = z + .object({ + from: z.number().int().safe().positive(), + to: z.number().int().safe().positive() + }) + .strict() + .refine(window => window.from < window.to, 'Query time window must be ordered') + .refine(window => window.to - window.from <= 24 * 60 * 60 * 1_000, 'Query time window is too large'); + +export const HERTZBEAT_QUERY_LIMITS = { + metricSeries: 32, + metricPointsPerSeries: 1_200, + tableRows: 1_000, + maximumWindowMs: 24 * 60 * 60 * 1_000 +} as const; + +const baseQueryShape = { + timeWindow: timeWindowSchema, + context: contextSchema.optional() +}; + +const metricQuerySchema = z + .object({ + signal: z.literal('metrics'), + queryKind: z.literal('time-series'), + ...baseQueryShape, + metric: z + .object({ + name: z.string().regex(/^[A-Za-z_:][A-Za-z0-9_:]{0,254}$/u), + aggregation: z.enum(['avg', 'sum', 'min', 'max', 'count']).optional(), + temporalAggregation: z.enum(['raw', 'rate', 'increase', 'delta']).optional(), + stepSeconds: z.number().int().positive().max(86_400).optional(), + operationName: boundedText.optional() + }) + .strict(), + limit: z.number().int().positive().max(HERTZBEAT_QUERY_LIMITS.metricSeries).optional() + }) + .strict(); + +const logTableQuerySchema = z + .object({ + signal: z.literal('logs'), + queryKind: z.literal('table'), + ...baseQueryShape, + search: boundedText.optional(), + severity: z.enum(['TRACE', 'DEBUG', 'INFO', 'WARN', 'ERROR', 'FATAL']).optional(), + traceId: identifier.optional(), + spanId: identifier.optional(), + hideInternal: z.boolean().optional(), + hideNoise: z.boolean().optional(), + limit: z.number().int().positive().max(HERTZBEAT_QUERY_LIMITS.tableRows).optional() + }) + .strict(); + +const traceTableQuerySchema = z + .object({ + signal: z.literal('traces'), + queryKind: z.literal('table'), + ...baseQueryShape, + traceId: identifier.optional(), + operationName: boundedText.optional(), + errorOnly: z.boolean().optional(), + minDurationMs: z.number().int().nonnegative().safe().optional(), + maxDurationMs: z.number().int().nonnegative().safe().optional(), + spanScope: z.enum(['root', 'entrypoint']).optional(), + hideInternal: z.boolean().optional(), + limit: z.number().int().positive().max(HERTZBEAT_QUERY_LIMITS.tableRows).optional() + }) + .strict() + .refine( + query => query.minDurationMs == null || query.maxDurationMs == null || query.minDurationMs <= query.maxDurationMs + ); + +const traceGanttQuerySchema = z + .object({ + signal: z.literal('traces'), + queryKind: z.literal('gantt'), + ...baseQueryShape, + traceId: identifier, + spanId: identifier.optional(), + minDurationMs: z.number().int().nonnegative().safe().optional(), + maxDurationMs: z.number().int().nonnegative().safe().optional() + }) + .strict() + .refine( + query => query.minDurationMs == null || query.maxDurationMs == null || query.minDurationMs <= query.maxDurationMs + ); + +export const hertzBeatQuerySchema = z.union([ + metricQuerySchema, + logTableQuerySchema, + traceTableQuerySchema, + traceGanttQuerySchema +]); + +export type HertzBeatQuery = z.infer<typeof hertzBeatQuerySchema>; +export type HertzBeatMetricQuery = z.infer<typeof metricQuerySchema>; +export type HertzBeatLogTableQuery = z.infer<typeof logTableQuerySchema>; +export type HertzBeatTraceTableQuery = z.infer<typeof traceTableQuerySchema>; +export type HertzBeatTraceGanttQuery = z.infer<typeof traceGanttQuerySchema>; + +export type HertzBeatQueryFailureKind = + 'invalid_request' | 'permission' | 'overloaded' | 'unavailable' | 'contract_error'; + +export type HertzBeatQueryFailure = { + kind: HertzBeatQueryFailureKind; + messageKey: + | 'perses.query.invalid' + | 'perses.query.permission' + | 'perses.query.overloaded' + | 'perses.query.unavailable' + | 'perses.query.contract'; + retryable: boolean; +}; + +export type HertzBeatQueryOutcome<T> = + | { state: 'ready'; data: T; truncated: boolean | 'unknown' } + | { state: 'empty'; truncated: false } + | { state: 'error'; error: HertzBeatQueryFailure }; diff --git a/web-app/src/platform/perses/datasource/hertzbeat-query-schema.ts b/web-app/src/platform/perses/datasource/hertzbeat-query-schema.ts new file mode 100644 index 0000000000..f224f842c3 --- /dev/null +++ b/web-app/src/platform/perses/datasource/hertzbeat-query-schema.ts @@ -0,0 +1,316 @@ +/* + * 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. + */ + +import { z } from 'zod'; + +import type { ExactTimeWindow } from '@/shared/query-context'; + +import { HERTZBEAT_QUERY_LIMITS } from './hertzbeat-query-contract'; + +const nullableText = z.string().nullable(); +const safeInteger = z.number().int().safe(); +const nonNegativeInteger = safeInteger.nonnegative(); +const nullableNonNegativeInteger = nonNegativeInteger.nullable(); +const javaLong = z + .number() + .finite() + .refine(Number.isInteger) + .refine(value => value >= 0); +const nullableJavaLong = javaLong.nullable(); +const stringMap = z.record(z.string(), z.string()); +const nullableStringMap = stringMap.nullable(); + +type JsonValue = null | boolean | number | string | JsonValue[] | { [key: string]: JsonValue }; +const jsonValue: z.ZodType<JsonValue> = z.lazy(() => + z.union([z.null(), z.boolean(), z.number().finite(), z.string(), z.array(jsonValue), z.record(z.string(), jsonValue)]) +); +const nullableJsonMap = z.record(z.string(), jsonValue).nullable(); + +const metricField = z.object({ + name: nullableText, + type: z.enum(['number', 'string', 'time', 'bool']).nullable(), + unit: nullableText +}); +const metricFrame = z.object({ + schema: z + .object({ fields: z.array(metricField).nullable(), labels: nullableStringMap, meta: nullableStringMap }) + .nullable(), + data: z.array(z.array(jsonValue)).max(HERTZBEAT_QUERY_LIMITS.metricPointsPerSeries).nullable() +}); +const metricConsole = z.object({ + context: z + .object({ + entityId: nullableJavaLong, + entityType: nullableText, + entityName: nullableText, + serviceName: nullableText, + serviceNamespace: nullableText, + environment: nullableText, + collectorId: nullableText, + instance: nullableText, + endpoint: nullableText, + operationName: nullableText, + start: nullableNonNegativeInteger, + end: nullableNonNegativeInteger + }) + .nullable(), + query: nullableText, + datasource: nullableText, + queryMode: nullableText, + results: z + .object({ + refId: nullableText, + status: safeInteger.nullable(), + msg: nullableText, + frames: z.array(metricFrame).max(HERTZBEAT_QUERY_LIMITS.metricSeries).nullable() + }) + .nullable(), + stats: z + .object({ + totalSeries: nonNegativeInteger, + nonEmptySeries: nonNegativeInteger, + latestObservedAt: nullableNonNegativeInteger + }) + .refine(stats => stats.nonEmptySeries <= stats.totalSeries) + .nullable(), + emptyStateReason: nullableText, + errorMessage: nullableText +}); + +const instrumentationScope = z.object({ + name: nullableText, + version: nullableText, + attributes: nullableJsonMap, + droppedAttributesCount: nullableNonNegativeInteger +}); +export const logRowSchema = z.object({ + timeUnixNano: nullableJavaLong, + observedTimeUnixNano: nullableJavaLong, + severityNumber: nullableNonNegativeInteger, + severityText: nullableText, + body: jsonValue, + attributes: nullableJsonMap, + droppedAttributesCount: nullableNonNegativeInteger, + traceId: nullableText, + spanId: nullableText, + traceFlags: nullableNonNegativeInteger, + resource: nullableJsonMap, + resourceSchemaUrl: nullableText, + instrumentationScope: instrumentationScope.nullable(), + scopeSchemaUrl: nullableText +}); + +const traceRowShape = { + traceId: z.string().min(1), + rootSpanId: nullableText, + serviceName: nullableText, + serviceNamespace: nullableText, + rootSpanName: nullableText, + durationNanos: nullableJavaLong, + status: nullableText, + startTime: nullableNonNegativeInteger, + errorSpanCount: nonNegativeInteger, + resourceAttributes: nullableStringMap +}; +export const traceRowSchema = z.object(traceRowShape); +const traceEvent = z.object({ + timeUnixNano: nullableJavaLong, + name: nullableText, + attributes: nullableJsonMap, + droppedAttributesCount: nullableNonNegativeInteger +}); +const traceLink = z.object({ + traceId: nullableText, + spanId: nullableText, + traceState: nullableText, + attributes: nullableJsonMap, + droppedAttributesCount: nullableNonNegativeInteger +}); +const codeNavigationHint = z.object({ + repositoryUrl: nullableText, + provider: nullableText, + defaultPath: nullableText, + searchQuery: nullableText, + label: nullableText +}); +const traceSpan = z.object({ + traceId: nullableText, + spanId: nullableText, + parentSpanId: nullableText, + spanName: nullableText, + serviceName: nullableText, + status: nullableText, + spanKind: nullableText, + statusMessage: nullableText, + traceState: nullableText, + scopeName: nullableText, + scopeVersion: nullableText, + durationNanos: nullableJavaLong, + startTime: nullableNonNegativeInteger, + highlighted: z.boolean(), + resourceAttributes: nullableStringMap, + spanAttributes: nullableStringMap, + events: z.array(traceEvent).nullable(), + links: z.array(traceLink).nullable(), + codeNavigationHint: codeNavigationHint.nullable() +}); +const traceDetail = z.object({ ...traceRowShape, spans: z.array(traceSpan).nullable() }); + +export type HertzBeatMetricSeries = { + key: string; + name: string; + unit?: string | undefined; + labels: Record<string, string>; + points: Array<{ timestamp: number; value: number }>; +}; +export type HertzBeatMetricData = { + timeWindow: ExactTimeWindow; + source: string | null; + series: HertzBeatMetricSeries[]; +}; +export type HertzBeatLogRow = z.infer<typeof logRowSchema>; +export type HertzBeatTraceRow = z.infer<typeof traceRowSchema>; +export type HertzBeatTraceDetail = z.infer<typeof traceDetail>; +export type HertzBeatTableData<T> = { rows: T[]; total: number }; + +export class HertzBeatResponseContractError extends Error { + constructor() { + super('Unexpected observability response'); + this.name = 'HertzBeatResponseContractError'; + } +} + +export class HertzBeatResponseStateError extends Error { + constructor(readonly kind: 'invalid_request' | 'unavailable') { + super('Observability response is not ready'); + this.name = 'HertzBeatResponseStateError'; + } +} + +export function parseMetricResponse(value: unknown, window: ExactTimeWindow): HertzBeatMetricData | undefined { + const parsed = metricConsole.safeParse(value); + if (!parsed.success) throw new HertzBeatResponseContractError(); + const console = parsed.data; + requireMetricReadyState(console); + if (console.context?.start !== window.from || console.context.end !== window.to) { + throw new HertzBeatResponseContractError(); + } + const frames = requireMetricFrames(console); + if (frames.length === 0) return undefined; + const series = frames.map((frame, index) => metricSeries(frame, index)); + return series.some(item => item.points.length > 0) + ? { timeWindow: window, source: console.datasource, series } + : undefined; +} + +function requireMetricReadyState(console: z.infer<typeof metricConsole>) { + if (console.emptyStateReason === 'no_context' || console.emptyStateReason === 'unsupported_query') { + throw new HertzBeatResponseStateError('invalid_request'); + } + if (console.errorMessage != null || console.emptyStateReason === 'load_failed') { + throw new HertzBeatResponseStateError('unavailable'); + } +} + +function requireMetricFrames(console: z.infer<typeof metricConsole>) { + if (!console.results || console.results.status !== 200 || !console.results.frames) { + throw new HertzBeatResponseStateError('unavailable'); + } + const frames = console.results.frames; + const nonEmptyFrames = frames.filter(frame => (frame.data?.length ?? 0) > 0).length; + if ( + !console.stats || + console.stats.totalSeries !== frames.length || + console.stats.nonEmptySeries !== nonEmptyFrames + ) { + throw new HertzBeatResponseContractError(); + } + return frames; +} + +export function parseLogTable(value: unknown, limit: number): HertzBeatTableData<HertzBeatLogRow> | undefined { + const result = z + .object({ + content: z.array(logRowSchema), + totalElements: nonNegativeInteger, + pageIndex: z.literal(0), + pageSize: z.literal(limit) + }) + .safeParse(value); + if ( + !result.success || + (result.data.content.length === 0 && result.data.totalElements > 0) || + result.data.content.length > Math.min(limit, result.data.totalElements) + ) { + throw new HertzBeatResponseContractError(); + } + return result.data.content.length > 0 ? { rows: result.data.content, total: result.data.totalElements } : undefined; +} + +export function parseTraceTable(value: unknown, limit: number): HertzBeatTableData<HertzBeatTraceRow> | undefined { + const result = z + .object({ + content: z.array(traceRowSchema), + totalElements: nonNegativeInteger, + totalPages: nonNegativeInteger, + number: z.literal(0), + size: z.literal(limit) + }) + .safeParse(value); + if (!result.success || result.data.totalPages !== Math.ceil(result.data.totalElements / limit)) { + throw new HertzBeatResponseContractError(); + } + if ( + (result.data.content.length === 0 && result.data.totalElements > 0) || + result.data.content.length > Math.min(limit, result.data.totalElements) + ) { + throw new HertzBeatResponseContractError(); + } + return result.data.content.length > 0 ? { rows: result.data.content, total: result.data.totalElements } : undefined; +} + +export function parseTraceGantt(value: unknown, traceId: string): HertzBeatTraceDetail | undefined { + if (value == null) return undefined; + const result = traceDetail.safeParse(value); + if (!result.success || result.data.traceId !== traceId) throw new HertzBeatResponseContractError(); + const spanIds = result.data.spans?.map(span => { + if (!span.spanId || (span.traceId !== null && span.traceId !== traceId)) throw new HertzBeatResponseContractError(); + return span.spanId; + }); + if (spanIds && new Set(spanIds).size !== spanIds.length) throw new HertzBeatResponseContractError(); + return result.data; +} + +function metricSeries(frame: z.infer<typeof metricFrame>, index: number): HertzBeatMetricSeries { + if (!frame.schema?.fields || !frame.schema.labels || !frame.data) throw new HertzBeatResponseContractError(); + const timeIndex = frame.schema.fields.findIndex(field => field.type === 'time'); + const valueIndex = frame.schema.fields.findIndex(field => field.type === 'number'); + if (timeIndex < 0 || valueIndex < 0) throw new HertzBeatResponseContractError(); + const valueField = frame.schema.fields[valueIndex]; + if (!valueField) throw new HertzBeatResponseContractError(); + const name = frame.schema.labels.__name__ ?? valueField.name ?? `series-${index + 1}`; + const points = frame.data.map(row => { + const timestamp = finiteNumber(row[timeIndex]); + const value = finiteNumber(row[valueIndex]); + if (timestamp == null || value == null) throw new HertzBeatResponseContractError(); + return { timestamp, value }; + }); + return { + key: `${name}-${index}`, + name, + ...(valueField.unit ? { unit: valueField.unit } : {}), + labels: frame.schema.labels, + points + }; +} + +function finiteNumber(value: JsonValue | undefined) { + if (typeof value === 'number') return Number.isFinite(value) ? value : undefined; + if (typeof value !== 'string' || !value.trim()) return undefined; + const parsed = Number(value); + return Number.isFinite(parsed) ? parsed : undefined; +} diff --git a/web-app/src/platform/perses/index.ts b/web-app/src/platform/perses/index.ts index ba6eddf200..a98cff4290 100644 --- a/web-app/src/platform/perses/index.ts +++ b/web-app/src/platform/perses/index.ts @@ -1,3 +1,26 @@ /* Licensed to the Apache Software Foundation (ASF) under the Apache License, Version 2.0. */ export { PersesTimeSeries } from './runtime/perses-time-series'; +export { + HERTZBEAT_QUERY_LIMITS, + hertzBeatPersesQueryClient, + queryHertzBeatData +} from './datasource/hertzbeat-query-client'; +export type { + HertzBeatLogTableQuery, + HertzBeatMetricQuery, + HertzBeatQuery, + HertzBeatQueryFailure, + HertzBeatQueryFailureKind, + HertzBeatQueryOutcome, + HertzBeatTraceGanttQuery, + HertzBeatTraceTableQuery +} from './datasource/hertzbeat-query-contract'; +export type { + HertzBeatLogRow, + HertzBeatMetricData, + HertzBeatMetricSeries, + HertzBeatTableData, + HertzBeatTraceDetail, + HertzBeatTraceRow +} from './datasource/hertzbeat-query-schema'; --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
