Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,42 @@ 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);

List<Map<String, Object>> docs = searchDocs(metricName);
assertThat("Only finite samples should be indexed", docs, hasSize(1));
assertThat(new ObjectPath(docs.getFirst()).evaluate("metrics." + metricName), equalTo(42.5));
Comment thread
felixbarny marked this conversation as resolved.
}

public void testRemoteWriteMissingNameLabelReturns400() throws Exception {
long timestamp = System.currentTimeMillis();

Expand Down Expand Up @@ -288,4 +324,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;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -111,14 +111,22 @@ protected void doExecute(Task task, RemoteWriteRequest request, ActionListener<R
}

for (Sample sample : timeSeries.getSamplesList()) {
if (Double.isFinite(sample.getValue()) == false) {
continue;
}
IndexRequest indexRequest = buildIndexRequest(timeSeries, sample, metricName, request.dataset, request.namespace);
bulkRequestBuilder.add(indexRequest);
}
}

if (bulkRequestBuilder.numberOfActions() == 0) {
String message = buildFailureSummary(totalSamples, droppedMissingName, droppedMissingName, Map.of());
listener.onFailure(new ElasticsearchStatusException(message, RestStatus.BAD_REQUEST));
if (droppedMissingName > 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;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -251,6 +251,61 @@ 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);
RemoteWriteResponse response = executeRequest(createWriteRequest("stale_metric", stalenessMarker, System.currentTimeMillis()));

assertThat(response.getStatus(), equalTo(RestStatus.NO_CONTENT));
assertNull(response.getMessage());
verify(client, never()).execute(any(), any(), any());
}

public void testNaNSamplesAreDropped() {
RemoteWriteResponse response = executeRequest(createWriteRequest("nan_metric", Double.NaN, System.currentTimeMillis()));

assertThat(response.getStatus(), equalTo(RestStatus.NO_CONTENT));
assertNull(response.getMessage());
verify(client, never()).execute(any(), any(), any());
}

public void testPositiveInfinitySamplesAreDropped() {
RemoteWriteResponse response = executeRequest(
createWriteRequest("inf_metric", Double.POSITIVE_INFINITY, System.currentTimeMillis())
);

assertThat(response.getStatus(), equalTo(RestStatus.NO_CONTENT));
assertNull(response.getMessage());
verify(client, never()).execute(any(), any(), any());
}

public void testNegativeInfinitySamplesAreDropped() {
RemoteWriteResponse response = executeRequest(
createWriteRequest("neg_inf_metric", Double.NEGATIVE_INFINITY, System.currentTimeMillis())
);

assertThat(response.getStatus(), equalTo(RestStatus.NO_CONTENT));
assertNull(response.getMessage());
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();

RemoteWriteResponse response = executeRequest(createWriteRequest(writeRequest, "generic", "default"));

assertThat(response.getStatus(), equalTo(RestStatus.NO_CONTENT));
}

public void testCustomDatasetAndNamespace() {
executeRequest(createWriteRequest("test_metric", 42.0, System.currentTimeMillis(), "myapp", "production"));
}
Expand Down
Loading