diff --git a/x-pack/plugin/prometheus/src/javaRestTest/java/org/elasticsearch/xpack/prometheus/PrometheusRemoteWriteRestIT.java b/x-pack/plugin/prometheus/src/javaRestTest/java/org/elasticsearch/xpack/prometheus/PrometheusRemoteWriteRestIT.java index 203aa2a04bd8f..e7b160e831b0e 100644 --- a/x-pack/plugin/prometheus/src/javaRestTest/java/org/elasticsearch/xpack/prometheus/PrometheusRemoteWriteRestIT.java +++ b/x-pack/plugin/prometheus/src/javaRestTest/java/org/elasticsearch/xpack/prometheus/PrometheusRemoteWriteRestIT.java @@ -142,6 +142,43 @@ public void testRemoteWriteWithMultipleTimeseriesAndSamples() throws Exception { assertThat(new ObjectPath(hits2.getFirst()).evaluate("metrics." + metric2), equalTo(100.0)); } + public void testRemoteWriteDropsNaNSamples() throws Exception { + long timestamp = System.currentTimeMillis(); + String metricName = "nan_metric"; + + RemoteWrite.WriteRequest writeRequest = RemoteWrite.WriteRequest.newBuilder() + .addTimeseries(timeSeries(metricName, Map.of("job", "test_job"), sample(Double.NaN, timestamp))) + .build(); + + sendAndAssertSuccess(writeRequest); + assertFalse("NaN-only request should not create a data stream", dataStreamExists(DEFAULT_DATA_STREAM)); + } + + public void testRemoteWriteDropsNaNButKeepsValidSamples() throws Exception { + long timestamp = System.currentTimeMillis(); + String metricName = "nan_mixed_metric"; + + RemoteWrite.WriteRequest writeRequest = RemoteWrite.WriteRequest.newBuilder() + .addTimeseries( + timeSeries( + metricName, + Map.of("job", "test_job"), + sample(Double.NaN, timestamp - 1000), + sample(42.5, timestamp), + sample(Double.POSITIVE_INFINITY, timestamp + 1000), + sample(Double.NEGATIVE_INFINITY, timestamp + 2000) + ) + ) + .build(); + + sendAndAssertSuccess(writeRequest); + assertTrue("Data stream should exist for valid samples", dataStreamExists(DEFAULT_DATA_STREAM)); + + List> docs = searchDocs(metricName); + assertThat("Only finite samples should be indexed", docs, hasSize(1)); + assertThat(new ObjectPath(docs.getFirst()).evaluate("metrics." + metricName), equalTo(42.5)); + } + public void testRemoteWriteMissingNameLabelReturns400() throws Exception { long timestamp = System.currentTimeMillis(); @@ -288,4 +325,16 @@ private ObjectPath searchSingleDoc(String dataStream, String metricName) throws return new ObjectPath(docs.getFirst()); } + private boolean dataStreamExists(String dataStream) throws IOException { + Request request = new Request("GET", "/_data_stream/" + dataStream); + try { + client().performRequest(request); + return true; + } catch (ResponseException e) { + if (e.getResponse().getStatusLine().getStatusCode() == 404) { + return false; + } + throw e; + } + } } diff --git a/x-pack/plugin/prometheus/src/main/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportAction.java b/x-pack/plugin/prometheus/src/main/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportAction.java index 2714f27a1328b..76754b2d47fed 100644 --- a/x-pack/plugin/prometheus/src/main/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportAction.java +++ b/x-pack/plugin/prometheus/src/main/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportAction.java @@ -111,14 +111,22 @@ protected void doExecute(Task task, RemoteWriteRequest request, ActionListener 0) { + String message = buildFailureSummary(totalSamples, droppedMissingName, droppedMissingName, Map.of()); + listener.onFailure(new ElasticsearchStatusException(message, RestStatus.BAD_REQUEST)); + } else { + // All samples were non-finite (NaN/Infinity) and silently dropped — not a client error + listener.onResponse(new RemoteWriteResponse()); + } return; } diff --git a/x-pack/plugin/prometheus/src/test/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportActionTests.java b/x-pack/plugin/prometheus/src/test/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportActionTests.java index 4ff40bed5871f..0a4c7731fe282 100644 --- a/x-pack/plugin/prometheus/src/test/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportActionTests.java +++ b/x-pack/plugin/prometheus/src/test/java/org/elasticsearch/xpack/prometheus/rest/PrometheusRemoteWriteTransportActionTests.java @@ -251,6 +251,43 @@ protected boolean apply(String actionName, ActionRequest actionRequest, ActionLi assertRegisteredIndexingPressureReleased("indexing pressure should be released when execution short-circuits"); } + public void testStalenessMarkerIsDropped() { + double stalenessMarker = Double.longBitsToDouble(0x7ff0000000000002L); + executeRequest(createWriteRequest("stale_metric", stalenessMarker, System.currentTimeMillis())); + verify(client, never()).execute(any(), any(), any()); + } + + public void testNaNSamplesAreDropped() { + executeRequest(createWriteRequest("nan_metric", Double.NaN, System.currentTimeMillis())); + verify(client, never()).execute(any(), any(), any()); + } + + public void testPositiveInfinitySamplesAreDropped() { + executeRequest(createWriteRequest("inf_metric", Double.POSITIVE_INFINITY, System.currentTimeMillis())); + verify(client, never()).execute(any(), any(), any()); + } + + public void testNegativeInfinitySamplesAreDropped() { + executeRequest(createWriteRequest("neg_inf_metric", Double.NEGATIVE_INFINITY, System.currentTimeMillis())); + verify(client, never()).execute(any(), any(), any()); + } + + public void testMixedFiniteAndNonFiniteSamples() { + long now = System.currentTimeMillis(); + RemoteWrite.WriteRequest writeRequest = RemoteWrite.WriteRequest.newBuilder() + .addTimeseries( + RemoteWrite.TimeSeries.newBuilder() + .addLabels(RemoteWrite.Label.newBuilder().setName("__name__").setValue("mixed_metric").build()) + .addSamples(RemoteWrite.Sample.newBuilder().setValue(Double.NaN).setTimestamp(now - 2000).build()) + .addSamples(RemoteWrite.Sample.newBuilder().setValue(42.0).setTimestamp(now - 1000).build()) + .addSamples(RemoteWrite.Sample.newBuilder().setValue(Double.POSITIVE_INFINITY).setTimestamp(now).build()) + .build() + ) + .build(); + + executeRequest(createWriteRequest(writeRequest, "generic", "default")); + } + public void testCustomDatasetAndNamespace() { executeRequest(createWriteRequest("test_metric", 42.0, System.currentTimeMillis(), "myapp", "production")); }