diff --git a/.dlc.json b/.dlc.json index 4a0bb02cc4f2..38fb4d029672 100644 --- a/.dlc.json +++ b/.dlc.json @@ -23,6 +23,9 @@ }, { "pattern": "^https://x.com/AsfSkyWalking" + }, + { + "pattern": "^https://docs.konghq.com/hub/kong-inc/prometheus/$" } ], "timeout": "10s", diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index 1841774216b7..6e0d39842b72 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -29,6 +29,7 @@ * Refuse unsupported TraceQL with `400` instead of dropping the predicate and answering with unfiltered traces: syntax errors (`||` inside a spanset, regex, an attribute without a scope or leading dot), `!=` and other unsupported operators, negation, attribute-existence checks, multiple spansets, unknown intrinsics and values, and conditions a datasource cannot filter (`kind` on Zipkin and SkyWalking, `resource.instance` on Zipkin, `resource.remote.service` and `status = unset` on SkyWalking, `resource.instance` or `name` without a service on SkyWalking). The deprecated `tags` search parameter is parsed as logfmt on every datasource and goes through the same mapping and refusals. `kind = server` and `status = error` accept the bare keyword as Tempo does. Every search-result span carries a `status` attribute (`error`, `ok`, `unset`) next to `service.name` and `span.kind`, so a trace list can show failures. A trace the storage matched but none of whose spans satisfies every condition is left out instead of listed with every span, and the `/otlp` datasource compares names the way the receiver indexed them and keeps the `resource.` and `span.` scopes apart in its tag index, so `span.env` no longer matches a resource attribute. `duration >` and `<` are strict at microsecond precision. (#14093) * Support the OpenTelemetry Collector `hostmetrics` receiver as an alternative source for Linux and Windows host monitoring. `vm.yaml` and `windows.yaml` map it to the same `meter_vm_*` / `meter_win_*` metrics as node-exporter and windows_exporter, whose metrics keep their meaning, and add metrics only hostmetrics provides: CPU core count and normalized CPU usage, plus CPU load, file-system usage, system handle count and pagefile usage on Windows. node-exporter metrics without a matching hostmetrics source (`tcp_alloc`, `sockets_used`, `udp_inuse`, `filefd_allocated`) stay node-exporter only. The new `process-hostmetrics-linux` and `process-hostmetrics-windows` rules, enabled by default, report each process name as an instance of its host (`mp_process_linux_*` / `mp_process_windows_*`: process count, threads, CPU, resident memory, open handles and oldest process uptime). Reference Collector configurations for both systems are under `docs/en/setup/backend/`; they set the `job_name` and normalized `process_name` labels the rules route and group on, and sum same-named processes before export. * AI agent conversations adopt AI Sessionizer `130601c`. Calls to MCP servers: the `execution` kind the Sessionizer's Claude Code plugin writes, `streams//execution--.sd`, is stored like any other file; the document lists every `execution/1` record under `tool_executions`, joined to its step by tool-use id and kept once by its own id, a tool step names its records under `executions`, and a call to an MCP server carries `mcp_server` and `mcp_tool` in its `attrs`. The new `otel-rules/ai-agent/mcp_endpoint.yaml` gives one endpoint per MCP server and tool, `/`, under the agent's service, with `meter_ai_agent_mcp_calls`, `meter_ai_agent_mcp_calls_by_outcome` and `meter_ai_agent_mcp_duration` from the Sessionizer's `agent.mcp.calls` and `agent.mcp.duration`. The token rules read `agent.token.usage`, the Sessionizer's name for the metric. The document follows the Sessionizer's: an `llm.call` takes its provider bodies from the round's `provider_bodies` attribute and the session its count from `provider_bodies_landed`, and a round whose bodies do not read or point past its range is refused; talks are in the order they began across streams; only a `child` stream's talk is a child's, not an `auxiliary` one's; a stream's `opened_by`, the relations and a step's edges are in the order they happened, by record position inside one stream or workflow run and by time across them, never by id; change and execution records of one instant are in the order they were read; a record is read by the fields its format lists, and one of another shape is skipped; a round with a frame field of another type, or a negative sequence, row or round number, does not read; a record time is RFC 3339 as Session Data defines it, to the second, and anything else is no time; record times compare as instants. +* Fix the Elasticsearch `BulkProcessor` treating a bulk write as fully successful whenever the overall HTTP status was 200, even when Elasticsearch rejected individual items in the same response, for example a 429 `cluster_block_exception` when an index is switched to `read_only_allow_delete` by the flood-stage disk watermark. The response body's `errors` flag and per-item `status`/`error` are now decoded and checked on every 200 response; a rejected request's future now completes exceptionally instead of being treated as a success, while the rest of the batch still completes normally. One ERROR line per `_bulk` request summarizes the rejections grouped by status code and error type, never a document id or the raw ES error reason. A response whose `items` array is shorter than the request (a malformed or truncated response) also fails those requests instead of silently treating them as succeeded. This also fixes a related batching bug where, once a flush was split into multiple `_bulk` HTTP requests by `batchOfBytes`, the completion of one chunk's request completed or failed every request in the whole flush instead of just its own chunk, and a bug where requests whose bulk failed to even build (e.g. an encoding error) were left with a future that never completed, which could block a persistence round forever. #### UI * Add a Virtual GenAI evaluation-record page and evaluation-score chart in Horizon UI, so operators can inspect evaluation result, level, reason, judge model, timestamp, trace linkage, and the `gen_ai_model_evaluation_score_ppm` trend for evaluated records. diff --git a/docs/en/setup/backend/backend-kong-monitoring.md b/docs/en/setup/backend/backend-kong-monitoring.md index 594621680a73..1e3dc2f4ad00 100644 --- a/docs/en/setup/backend/backend-kong-monitoring.md +++ b/docs/en/setup/backend/backend-kong-monitoring.md @@ -7,13 +7,13 @@ SkyWalking leverages OpenTelemetry Collector to transfer the metrics to[OpenTele and into the [Meter System](./../../concepts-and-designs/mal.md). ### Data flow -1. [KONG Prometheus plugin](https://docs.konghq.com/hub/kong-inc/prometheus/) collects metrics data from KONG. -2. OpenTelemetry Collector fetches metrics from [KONG Prometheus plugin](https://docs.konghq.com/hub/kong-inc/prometheus/) via +1. [KONG Prometheus plugin](https://developer.konghq.com/plugins/prometheus/) collects metrics data from KONG. +2. OpenTelemetry Collector fetches metrics from [KONG Prometheus plugin](https://developer.konghq.com/plugins/prometheus/) via Prometheus Receiver and pushes metrics to SkyWalking OAP Server via OpenTelemetry gRPC exporter. 3. The SkyWalking OAP Server parses the expression with [MAL](../../concepts-and-designs/mal.md) to filter/calculate/aggregate and store the results. ### Set up -1. Enable KONG [KONG Prometheus plugin](https://docs.konghq.com/hub/kong-inc/prometheus/). Note that if need to monitor per_consumer, +1. Enable KONG [KONG Prometheus plugin](https://developer.konghq.com/plugins/prometheus/). Note that if need to monitor per_consumer, status_code_metrics, ai_metrics, latency_metrics, bandwidth_metrics or upstream_health_metrics, **need to enable them manually as needed**, which can be enabled in the [konga](https://pantsel.github.io/konga/) dashboard or through the Admin API, such as the following command ~~~bash @@ -32,7 +32,7 @@ and into the [Meter System](./../../concepts-and-designs/mal.md). ### KONG Monitoring -[KONG prometheus plugin](https://docs.konghq.com/hub/kong-inc/prometheus/) provide multiple dimensions metrics for KONG server, upstream, route etc. +[KONG prometheus plugin](https://developer.konghq.com/plugins/prometheus/) provide multiple dimensions metrics for KONG server, upstream, route etc. Accordingly, SkyWalking observes the status, requests, and latency of the KONG server, which is cataloged as a `LAYER: KONG` `Service` in the OAP. Each Kong server is cataloged as a `LAYER: KONG` `instance`, meanwhile, the route rules would be recognized as a `LAYER: KONG` `endpoint`. diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java index dee6f3590b65..b1442e3f96b3 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessor.java @@ -17,20 +17,25 @@ package org.apache.skywalking.library.elasticsearch.bulk; +import com.linecorp.armeria.common.HttpData; import com.linecorp.armeria.common.HttpStatus; import com.linecorp.armeria.common.util.Exceptions; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; +import java.io.InputStream; import java.time.Duration; import java.util.ArrayList; import java.util.Collections; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.concurrent.Semaphore; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; +import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; @@ -39,6 +44,8 @@ import org.apache.skywalking.library.elasticsearch.requests.UpdateRequest; import org.apache.skywalking.library.elasticsearch.requests.factory.Codec; import org.apache.skywalking.library.elasticsearch.requests.factory.RequestFactory; +import org.apache.skywalking.library.elasticsearch.response.bulk.BulkItemResult; +import org.apache.skywalking.library.elasticsearch.response.bulk.BulkResponse; import org.apache.skywalking.oap.server.library.util.CollectionUtils; import org.apache.skywalking.oap.server.library.util.RunnableWithExceptionProtection; import org.apache.skywalking.oap.server.telemetry.api.HistogramMetrics; @@ -161,23 +168,31 @@ private List> doFlush(final List batch) { final List bs = new ArrayList<>(); List> futures = new ArrayList<>(); List byteBufList = new ArrayList<>(); + List> holderChunks = new ArrayList<>(); + List currentChunkHolders = new ArrayList<>(); for (final Holder holder : batch) { byte[] bytes = codec.encode(holder.request); bs.add(bytes); bs.add("\n".getBytes()); bufferOfBytes += bytes.length + 1; + currentChunkHolders.add(holder); if (bufferOfBytes >= batchOfBytes) { final ByteBuf content = Unpooled.wrappedBuffer(bs.toArray(new byte[0][])); byteBufList.add(content); + holderChunks.add(currentChunkHolders); bs.clear(); bufferOfBytes = 0; + currentChunkHolders = new ArrayList<>(); } } if (CollectionUtils.isNotEmpty(bs)) { final ByteBuf content = Unpooled.wrappedBuffer(bs.toArray(new byte[0][])); byteBufList.add(content); + holderChunks.add(currentChunkHolders); } - for (final ByteBuf content : byteBufList) { + for (int i = 0; i < byteBufList.size(); i++) { + final ByteBuf content = byteBufList.get(i); + final List chunkHolders = holderChunks.get(i); CompletableFuture future = es.get().version().thenCompose(v -> { try { final RequestFactory rf = v.requestFactory(); @@ -186,6 +201,12 @@ private List> doFlush(final List batch) { if (status != HttpStatus.OK) { throw new RuntimeException(response.contentUtf8()); } + try (final HttpData responseContent = response.content(); + final InputStream is = responseContent.toInputStream()) { + completeHolders(chunkHolders, v.codec().decode(is, BulkResponse.class)); + } catch (Exception e) { + Exceptions.throwUnsafely(e); + } }); } catch (Exception e) { return Exceptions.throwUnsafely(e); @@ -193,12 +214,9 @@ private List> doFlush(final List batch) { }); future.whenComplete((ignored, exception) -> { if (exception != null) { - batch.stream().map(it -> it.future) - .forEach(it -> it.completeExceptionally((Throwable) exception)); + chunkHolders.stream().map(it -> it.future) + .forEach(it -> it.completeExceptionally((Throwable) exception)); log.error("Failed to execute requests in bulk", exception); - } else { - log.debug("Succeeded to execute {} requests in bulk", batch.size()); - batch.stream().map(it -> it.future).forEach(it -> it.complete(null)); } }); futures.add(future); @@ -207,10 +225,56 @@ private List> doFlush(final List batch) { } catch (Exception e) { log.error("Failed to execute requests in bulk", e); + batch.forEach(it -> it.future.completeExceptionally(e)); return Collections.emptyList(); } } + static void completeHolders(final List holders, final BulkResponse bulkResponse) { + if (!bulkResponse.isErrors()) { + holders.forEach(it -> it.future.complete(null)); + return; + } + + final Map rejections = new LinkedHashMap<>(); + final List rejected = new ArrayList<>(); + final List> items = bulkResponse.getItems(); + for (int i = 0; i < holders.size(); i++) { + final BulkItemResult result = itemResultOf(items, i); + if (result != null && result.getError() == null) { + holders.get(i).future.complete(null); + continue; + } + final String group = result == null + ? "no response item" + : "code=" + result.getStatus() + " type=" + result.getError().getType(); + rejections.merge(group, 1, Integer::sum); + rejected.add(holders.get(i)); + } + if (rejected.isEmpty()) { + return; + } + + final String message = "Bulk request to ES had " + rejected.size() + " of " + holders.size() + + " items rejected: " + rejections.entrySet().stream() + .map(e -> e.getKey() + " count=" + e.getValue()) + .collect(Collectors.joining("; ")); + log.error(message); + final RuntimeException failure = new RuntimeException(message); + rejected.forEach(it -> it.future.completeExceptionally(failure)); + } + + static BulkItemResult itemResultOf(final List> items, final int index) { + if (items == null || index >= items.size()) { + return null; + } + final Map item = items.get(index); + if (item == null || item.isEmpty()) { + return null; + } + return item.values().iterator().next(); + } + @RequiredArgsConstructor static class Holder { private final CompletableFuture future; diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemError.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemError.java new file mode 100644 index 000000000000..2f249f1c4046 --- /dev/null +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemError.java @@ -0,0 +1,27 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.skywalking.library.elasticsearch.response.bulk; + +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +public final class BulkItemError { + private String type; +} diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemResult.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemResult.java new file mode 100644 index 000000000000..c120543802e9 --- /dev/null +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemResult.java @@ -0,0 +1,28 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.skywalking.library.elasticsearch.response.bulk; + +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +public final class BulkItemResult { + private int status; + private BulkItemError error; +} diff --git a/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkResponse.java b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkResponse.java new file mode 100644 index 000000000000..93535f25709b --- /dev/null +++ b/oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkResponse.java @@ -0,0 +1,30 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.skywalking.library.elasticsearch.response.bulk; + +import java.util.List; +import java.util.Map; +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +public final class BulkResponse { + private boolean errors; + private List> items; +} diff --git a/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java b/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java index 659954f60b12..a23005e34aeb 100644 --- a/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java +++ b/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/ElasticSearchIT.java @@ -17,6 +17,7 @@ package org.apache.skywalking.library.elasticsearch; +import org.apache.skywalking.library.elasticsearch.bulk.BulkProcessor; import org.apache.skywalking.library.elasticsearch.client.TemplateClient; import org.apache.skywalking.library.elasticsearch.requests.IndexRequest; import org.apache.skywalking.library.elasticsearch.requests.UpdateRequest; @@ -28,6 +29,7 @@ import org.apache.skywalking.library.elasticsearch.response.IndexTemplate; import org.apache.skywalking.library.elasticsearch.response.Mappings; import org.apache.skywalking.library.elasticsearch.response.search.SearchResponse; +import org.apache.skywalking.oap.server.telemetry.api.HistogramMetrics; import org.awaitility.Duration; import org.junit.jupiter.api.Tag; import org.junit.jupiter.params.ParameterizedTest; @@ -42,6 +44,10 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; import static org.assertj.core.api.Assertions.assertThat; import static org.awaitility.Awaitility.await; @@ -49,6 +55,7 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @Tag("slow") @@ -450,4 +457,87 @@ public void testClientBuilder(final String ignored, server.close(); } + + /** + * https://github.com/apache/skywalking/issues/14113: Elasticsearch answers a bulk write with overall HTTP 200 + * while rejecting individual items, e.g. when an index is blocked for writes. Run against every ES/OpenSearch + * version in the matrix, and with a chunked ({@code batchOfBytes(1)}) and unchunked bulk request, to also cover + * the fix that scoped a chunk's completion to its own requests instead of the whole flush. + */ + @ParameterizedTest(name = "version: {0}") + @MethodSource("es") + public void testBulkProcessorItemLevelFailure(final String ignored, + final ElasticsearchContainer server) throws Exception { + server.start(); + + final ElasticSearch client = + ElasticSearch.builder() + .endpoints(server.getHttpHostAddress()) + .build(); + client.connect(); + + final String normalIndex = "bulk-item-failure-normal"; + final String blockedIndex = "bulk-item-failure-blocked"; + assertTrue(client.index().create(normalIndex, null, null)); + final Map blockedSettings = new HashMap<>(); + blockedSettings.put("index.blocks.write", true); + assertTrue(client.index().create(blockedIndex, null, blockedSettings)); + + // One `_bulk` request carrying both a succeeding and a rejected item. + runBulkAndAssertItemLevelFailure(client, normalIndex, blockedIndex, "one-chunk", 5 * 1024 * 1024); + // `batchOfBytes(1)` forces each item into its own chunk/`_bulk` request. + runBulkAndAssertItemLevelFailure(client, normalIndex, blockedIndex, "many-chunks", 1); + + server.close(); + } + + private static void runBulkAndAssertItemLevelFailure(final ElasticSearch client, + final String normalIndex, + final String blockedIndex, + final String idSuffix, + final int batchOfBytes) throws Exception { + final BulkProcessor bulkProcessor = BulkProcessor.builder() + .bulkActions(2) + .batchOfBytes(batchOfBytes) + .flushInterval(java.time.Duration.ofSeconds(30)) + .concurrentRequests(2) + .bulkMetrics(new HistogramMetrics() { + @Override + public void observe(final double value) { + } + }) + .build(new AtomicReference<>(client)); + + final String type = "type"; + final String normalId = "normal-" + idSuffix; + final String blockedId = "blocked-" + idSuffix; + + final CompletableFuture normalFuture = bulkProcessor.add( + IndexRequest.builder() + .index(normalIndex) + .type(type) + .id(normalId) + .doc(ImmutableMap.of("key", "val")) + .build()); + final CompletableFuture blockedFuture = bulkProcessor.add( + IndexRequest.builder() + .index(blockedIndex) + .type(type) + .id(blockedId) + .doc(ImmutableMap.of("key", "val")) + .build()); + + bulkProcessor.flush(); + + // The scheduler thread's own periodical flush can race this method's explicit flush() and drain either + // or both requests first, so wait on the futures themselves instead of asserting isDone() right after + // flush() returns. + normalFuture.get(30, TimeUnit.SECONDS); + assertTrue(client.documents().get(normalIndex, type, normalId).isPresent()); + + final String message = assertThrows( + ExecutionException.class, () -> blockedFuture.get(30, TimeUnit.SECONDS)).getCause().getMessage(); + assertTrue(message.contains("type=cluster_block_exception"), message); + assertFalse(message.contains(blockedId), "the rejected document id must not be logged: " + message); + } } diff --git a/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessorTest.java b/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessorTest.java new file mode 100644 index 000000000000..ae3482428fd5 --- /dev/null +++ b/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessorTest.java @@ -0,0 +1,163 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.skywalking.library.elasticsearch.bulk; + +import java.io.ByteArrayInputStream; +import java.io.InputStream; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutionException; +import org.apache.skywalking.library.elasticsearch.requests.factory.v7plus.codec.V78Codec; +import org.apache.skywalking.library.elasticsearch.response.bulk.BulkItemResult; +import org.apache.skywalking.library.elasticsearch.response.bulk.BulkResponse; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class BulkProcessorTest { + /** + * A response reported in https://github.com/apache/skywalking/issues/14113: Elasticsearch returns overall + * HTTP 200 while an item is rejected because the index is blocked by the flood-stage disk watermark. + */ + private static final String BLOCKED_ITEM_RESPONSE = "{" + + "\"took\":0," + + "\"errors\":true," + + "\"items\":[{" + + " \"update\":{" + + " \"_index\":\"blocked-test-index\"," + + " \"_id\":\"oap-block-check-1\"," + + " \"status\":429," + + " \"error\":{" + + " \"type\":\"cluster_block_exception\"," + + " \"reason\":\"index [blocked-test-index] blocked by: [TOO_MANY_REQUESTS/12/disk usage exceeded flood-stage watermark, index has read-only-allow-delete block];\"" + + " }" + + " }" + + "}]}"; + + private static final String ALL_SUCCEEDED_RESPONSE = "{" + + "\"took\":1," + + "\"errors\":false," + + "\"items\":[{\"index\":{\"_index\":\"test-index\",\"_id\":\"1\",\"status\":201}}]}"; + + @Test + void decodesItemLevelFailureFromHttp200Response() throws Exception { + final BulkResponse response = decode(BLOCKED_ITEM_RESPONSE); + + assertTrue(response.isErrors()); + assertEquals(1, response.getItems().size()); + final BulkItemResult item = response.getItems().get(0).values().iterator().next(); + assertEquals(429, item.getStatus()); + assertEquals("cluster_block_exception", item.getError().getType()); + } + + @Test + void decodesAllSucceededResponse() throws Exception { + final BulkResponse response = decode(ALL_SUCCEEDED_RESPONSE); + + assertFalse(response.isErrors()); + assertNull(response.getItems().get(0).values().iterator().next().getError()); + } + + @Test + void completesAllHoldersWhenNoErrors() throws Exception { + final BulkResponse response = decode(ALL_SUCCEEDED_RESPONSE); + final CompletableFuture future = new CompletableFuture<>(); + final List holders = holdersOf(future); + + BulkProcessor.completeHolders(holders, response); + + assertTrue(future.isDone()); + assertFalse(future.isCompletedExceptionally()); + } + + @Test + void failsOnlyTheItemsThatAreRejectedByElasticsearchWithASummaryMessage() throws Exception { + final String twoItemResponse = "{" + + "\"errors\":true," + + "\"items\":[" + + " {\"index\":{\"_index\":\"idx\",\"_id\":\"1\",\"status\":201}}," + + " {\"update\":{\"_index\":\"idx\",\"_id\":\"rejected-doc-2\",\"status\":429," + + " \"error\":{\"type\":\"cluster_block_exception\",\"reason\":\"disk usage exceeded flood-stage watermark\"}}}" + + "]}"; + final BulkResponse response = decode(twoItemResponse); + + final CompletableFuture succeeded = new CompletableFuture<>(); + final CompletableFuture failed = new CompletableFuture<>(); + final List holders = Arrays.asList( + new BulkProcessor.Holder(succeeded, "req-1"), + new BulkProcessor.Holder(failed, "req-2")); + + BulkProcessor.completeHolders(holders, response); + + assertTrue(succeeded.isDone()); + assertFalse(succeeded.isCompletedExceptionally()); + + assertTrue(failed.isDone()); + assertTrue(failed.isCompletedExceptionally()); + final ExecutionException ex = assertThrows(ExecutionException.class, failed::get); + final String message = ex.getCause().getMessage(); + assertTrue(message.contains("1 of 2 items rejected"), message); + assertTrue(message.contains("code=429 type=cluster_block_exception count=1"), message); + assertFalse(message.contains("rejected-doc-2"), "the rejected document id must not be logged: " + message); + assertFalse(message.contains("flood-stage"), "the raw ES error reason must not be logged: " + message); + } + + @Test + void failsRequestsThatHaveNoCorrespondingResponseItem() throws Exception { + // "errors: true" but only one item reported for two requests in the chunk. + final BulkResponse response = decode(BLOCKED_ITEM_RESPONSE); + + final CompletableFuture withItem = new CompletableFuture<>(); + final CompletableFuture withoutItem = new CompletableFuture<>(); + final List holders = Arrays.asList( + new BulkProcessor.Holder(withItem, "req-1"), + new BulkProcessor.Holder(withoutItem, "req-2")); + + BulkProcessor.completeHolders(holders, response); + + assertTrue(withItem.isCompletedExceptionally()); + assertTrue(withoutItem.isCompletedExceptionally()); + final String message = assertThrows(ExecutionException.class, withoutItem::get).getCause().getMessage(); + assertTrue(message.contains("no response item count=1"), message); + } + + @Test + void itemResultOfReturnsNullWhenIndexOutOfBounds() { + assertNull(BulkProcessor.itemResultOf(null, 0)); + assertNull(BulkProcessor.itemResultOf(new ArrayList<>(), 0)); + } + + private static BulkResponse decode(final String json) throws Exception { + try (final InputStream is = new ByteArrayInputStream(json.getBytes(StandardCharsets.UTF_8))) { + return V78Codec.INSTANCE.decode(is, BulkResponse.class); + } + } + + private static List holdersOf(final CompletableFuture future) { + final List holders = new ArrayList<>(); + holders.add(new BulkProcessor.Holder(future, "req")); + return holders; + } +} diff --git a/test/e2e-v2/cases/event/banyandb/e2e.yaml b/test/e2e-v2/cases/event/banyandb/e2e.yaml index 7642e4ccf7a8..c1cd26eb0857 100644 --- a/test/e2e-v2/cases/event/banyandb/e2e.yaml +++ b/test/e2e-v2/cases/event/banyandb/e2e.yaml @@ -27,6 +27,17 @@ setup: command: bash test/e2e-v2/script/prepare/setup-e2e-shell/install.sh yq - name: install swctl command: bash test/e2e-v2/script/prepare/setup-e2e-shell/install.sh swctl + # No service here waits on OAP's healthcheck, so wait for OAP to answer before verifying; + # it opens its ports only after installing the BanyanDB schemas, which outlasts the verify retries. + - name: wait for OAP + command: | + export PATH=/tmp/skywalking-infra-e2e/bin:$PATH + for i in $(seq 1 100); do + swctl --base-url=http://${oap_host}:${oap_12800}/graphql service ls > /dev/null 2>&1 && exit 0 + sleep 3 + done + echo "OAP did not answer within 300s" + exit 1 verify: retry: