From fd700487c366e3272f2f02144e39ee9e942d7df4 Mon Sep 17 00:00:00 2001 From: kezhenxu94 Date: Mon, 28 Sep 2026 14:48:41 +0800 Subject: [PATCH 1/5] Fix Elasticsearch BulkProcessor ignoring item-level failures in HTTP 200 bulk responses BulkProcessor.doFlush only checked the overall HTTP status of a _bulk response, so an item rejected by Elasticsearch (e.g. 429 cluster_block_exception when an index is switched to read_only_allow_delete by the flood-stage disk watermark) inside an HTTP 200 response was silently treated as a success. Decode the response body's errors flag and per-item status/error on every 200 response, log rejected items with their index/id/status/ error detail, and complete only the failed item's future exceptionally while the rest of the batch still completes normally. Also fixes a related bug where, once a flush was split into multiple _bulk HTTP requests by batchOfBytes, one chunk's completion was applied to every request in the whole flush instead of just its own chunk. Fixes #14113 --- docs/en/changes/changes.md | 1 + .../elasticsearch/bulk/BulkProcessor.java | 68 ++++++++- .../response/bulk/BulkItemError.java | 28 ++++ .../response/bulk/BulkItemResult.java | 33 +++++ .../response/bulk/BulkResponse.java | 30 ++++ .../elasticsearch/bulk/BulkProcessorTest.java | 137 ++++++++++++++++++ 6 files changed, 291 insertions(+), 6 deletions(-) create mode 100644 oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemError.java create mode 100644 oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkItemResult.java create mode 100644 oap-server/server-library/library-elasticsearch-client/src/main/java/org/apache/skywalking/library/elasticsearch/response/bulk/BulkResponse.java create mode 100644 oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessorTest.java diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index 1841774216b7..6ee0d90b910d 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 item is logged with its index, id, status and error detail and completes only that item's future exceptionally, while the rest of the batch still completes normally. 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. #### 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/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..c943a13df6a0 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,14 +17,17 @@ 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.List; +import java.util.Map; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ScheduledThreadPoolExecutor; @@ -39,6 +42,9 @@ 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.BulkItemError; +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 +167,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 +200,14 @@ private List> doFlush(final List batch) { if (status != HttpStatus.OK) { throw new RuntimeException(response.contentUtf8()); } + final BulkResponse bulkResponse; + try (final HttpData responseContent = response.content(); + final InputStream is = responseContent.toInputStream()) { + bulkResponse = v.codec().decode(is, BulkResponse.class); + } catch (Exception e) { + throw new RuntimeException(e); + } + completeHolders(chunkHolders, bulkResponse); }); } catch (Exception e) { return Exceptions.throwUnsafely(e); @@ -193,12 +215,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); @@ -211,6 +230,43 @@ private List> doFlush(final List batch) { } } + static void completeHolders(final List holders, final BulkResponse bulkResponse) { + if (!bulkResponse.isErrors()) { + log.debug("Succeeded to execute {} requests in bulk", holders.size()); + holders.forEach(it -> it.future.complete(null)); + return; + } + + final List> items = bulkResponse.getItems(); + for (int i = 0; i < holders.size(); i++) { + final Holder holder = holders.get(i); + final BulkItemResult result = itemResultOf(items, i); + final BulkItemError error = result == null ? null : result.getError(); + if (error == null) { + holder.future.complete(null); + continue; + } + log.error( + "Failed to execute bulk item, index: {}, id: {}, status: {}, error type: {}, reason: {}", + result.getIndex(), result.getId(), result.getStatus(), error.getType(), error.getReason()); + holder.future.completeExceptionally(new RuntimeException( + "Bulk item failed, index: " + result.getIndex() + ", id: " + result.getId() + + ", status: " + result.getStatus() + ", error type: " + error.getType() + + ", reason: " + error.getReason())); + } + } + + 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..300dc34e25db --- /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,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 BulkItemError { + private String type; + private String reason; +} 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..c9589841fce1 --- /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,33 @@ +/* + * 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 com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Getter; +import lombok.Setter; + +@Getter +@Setter +public final class BulkItemResult { + private int status; + @JsonProperty("_index") + private String index; + @JsonProperty("_id") + private String id; + 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/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..7ed0c5764e61 --- /dev/null +++ b/oap-server/server-library/library-elasticsearch-client/src/test/java/org/apache/skywalking/library/elasticsearch/bulk/BulkProcessorTest.java @@ -0,0 +1,137 @@ +/* + * 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 com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.ObjectMapper; +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.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 { + private static final ObjectMapper MAPPER = new ObjectMapper() + .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + + /** + * 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 = MAPPER.readValue(BLOCKED_ITEM_RESPONSE, BulkResponse.class); + + assertTrue(response.isErrors()); + assertEquals(1, response.getItems().size()); + final BulkItemResult item = response.getItems().get(0).values().iterator().next(); + assertEquals(429, item.getStatus()); + assertEquals("blocked-test-index", item.getIndex()); + assertEquals("oap-block-check-1", item.getId()); + assertEquals("cluster_block_exception", item.getError().getType()); + } + + @Test + void decodesAllSucceededResponse() throws Exception { + final BulkResponse response = MAPPER.readValue(ALL_SUCCEEDED_RESPONSE, BulkResponse.class); + + assertFalse(response.isErrors()); + assertNull(response.getItems().get(0).values().iterator().next().getError()); + } + + @Test + void completesAllHoldersWhenNoErrors() throws Exception { + final BulkResponse response = MAPPER.readValue(ALL_SUCCEEDED_RESPONSE, BulkResponse.class); + final CompletableFuture future = new CompletableFuture<>(); + final List holders = holdersOf(future); + + BulkProcessor.completeHolders(holders, response); + + assertTrue(future.isDone()); + assertFalse(future.isCompletedExceptionally()); + } + + @Test + void failsOnlyTheItemThatIsRejectedByElasticsearch() throws Exception { + final String twoItemResponse = "{" + + "\"errors\":true," + + "\"items\":[" + + " {\"index\":{\"_index\":\"idx\",\"_id\":\"1\",\"status\":201}}," + + " {\"update\":{\"_index\":\"idx\",\"_id\":\"2\",\"status\":429," + + " \"error\":{\"type\":\"cluster_block_exception\",\"reason\":\"disk usage exceeded flood-stage watermark\"}}}" + + "]}"; + final BulkResponse response = MAPPER.readValue(twoItemResponse, BulkResponse.class); + + 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); + assertTrue(ex.getCause().getMessage().contains("cluster_block_exception")); + } + + @Test + void itemResultOfReturnsNullWhenIndexOutOfBounds() { + assertNull(BulkProcessor.itemResultOf(null, 0)); + assertNull(BulkProcessor.itemResultOf(new ArrayList<>(), 0)); + } + + private static List holdersOf(final CompletableFuture future) { + final List holders = new ArrayList<>(); + holders.add(new BulkProcessor.Holder(future, "req")); + return holders; + } +} From b6ba3c7f3c39f0db7c35c37fd78ae0d583ee1679 Mon Sep 17 00:00:00 2001 From: kezhenxu94 Date: Mon, 28 Sep 2026 17:28:33 +0800 Subject: [PATCH 2/5] Address review: summarize bulk item failures, fail missing-item and pre-send failures, add IT Per wu-sheng's review on #14113: - Log one ERROR summary line per _bulk request grouped by status code and error type, instead of one ERROR per rejected item. Never log a document id or the raw Elasticsearch error reason, since ES embeds the document id (and sometimes a preview of the stored value) into it; with the default bulkActions, a sustained block would otherwise flood the log with one line per item. All rejected items in a request share one exception built from the summary message. - Fail requests whose bulk response has no corresponding item (a malformed/truncated response body), instead of silently completing them as successful. - Decode the response body using Exceptions.throwUnsafely inside the try, matching the idiom used elsewhere in this client (SearchClient, ElasticSearch.connect) instead of wrapping in a new RuntimeException. - Complete every future in the batch exceptionally when doFlush fails before a request is even sent (e.g. encoding or version lookup failure), so PersistenceTimer's join on the round can't hang forever on an uncompleted future. - Add an ElasticSearchIT case that runs BulkProcessor against a normal index and one with index.blocks.write, both as a single _bulk request and chunked via batchOfBytes(1), covering the chunk-scoping fix across the real ES/OpenSearch version matrix. - Update BulkProcessorTest to decode through the production V78Codec instead of a local ObjectMapper, assert the new summary message, and cover a response missing an item. --- docs/en/changes/changes.md | 2 +- .../elasticsearch/bulk/BulkProcessor.java | 42 +++++---- .../elasticsearch/ElasticSearchIT.java | 88 +++++++++++++++++++ .../elasticsearch/bulk/BulkProcessorTest.java | 50 ++++++++--- 4 files changed, 153 insertions(+), 29 deletions(-) diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index 6ee0d90b910d..375c2ad82266 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -29,7 +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 item is logged with its index, id, status and error detail and completes only that item's future exceptionally, while the rest of the batch still completes normally. 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. +* 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, which can embed the id and, for mapping errors, a preview of the stored field value) — with the default `bulkActions` of thousands of requests per flush, logging per item would flood the log when every item of a bulk is rejected, as happens under a sustained disk watermark block. 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/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 c943a13df6a0..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 @@ -26,6 +26,7 @@ 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; @@ -34,6 +35,7 @@ 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; @@ -42,7 +44,6 @@ 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.BulkItemError; 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; @@ -200,14 +201,12 @@ private List> doFlush(final List batch) { if (status != HttpStatus.OK) { throw new RuntimeException(response.contentUtf8()); } - final BulkResponse bulkResponse; try (final HttpData responseContent = response.content(); final InputStream is = responseContent.toInputStream()) { - bulkResponse = v.codec().decode(is, BulkResponse.class); + completeHolders(chunkHolders, v.codec().decode(is, BulkResponse.class)); } catch (Exception e) { - throw new RuntimeException(e); + Exceptions.throwUnsafely(e); } - completeHolders(chunkHolders, bulkResponse); }); } catch (Exception e) { return Exceptions.throwUnsafely(e); @@ -226,34 +225,43 @@ 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()) { - log.debug("Succeeded to execute {} requests in bulk", holders.size()); 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 Holder holder = holders.get(i); final BulkItemResult result = itemResultOf(items, i); - final BulkItemError error = result == null ? null : result.getError(); - if (error == null) { - holder.future.complete(null); + if (result != null && result.getError() == null) { + holders.get(i).future.complete(null); continue; } - log.error( - "Failed to execute bulk item, index: {}, id: {}, status: {}, error type: {}, reason: {}", - result.getIndex(), result.getId(), result.getStatus(), error.getType(), error.getReason()); - holder.future.completeExceptionally(new RuntimeException( - "Bulk item failed, index: " + result.getIndex() + ", id: " + result.getId() - + ", status: " + result.getStatus() + ", error type: " + error.getType() - + ", reason: " + error.getReason())); + 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) { 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..69bee847aa83 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,9 @@ 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.atomic.AtomicReference; import static org.assertj.core.api.Assertions.assertThat; import static org.awaitility.Awaitility.await; @@ -49,6 +54,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 +456,86 @@ 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(); + + assertTrue(normalFuture.isDone()); + assertFalse(normalFuture.isCompletedExceptionally()); + assertTrue(client.documents().get(normalIndex, type, normalId).isPresent()); + + assertTrue(blockedFuture.isDone()); + assertTrue(blockedFuture.isCompletedExceptionally()); + final String message = assertThrows(ExecutionException.class, blockedFuture::get).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 index 7ed0c5764e61..4177bfbd5eef 100644 --- 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 @@ -17,13 +17,15 @@ package org.apache.skywalking.library.elasticsearch.bulk; -import com.fasterxml.jackson.databind.DeserializationFeature; -import com.fasterxml.jackson.databind.ObjectMapper; +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; @@ -35,9 +37,6 @@ import static org.junit.jupiter.api.Assertions.assertTrue; class BulkProcessorTest { - private static final ObjectMapper MAPPER = new ObjectMapper() - .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); - /** * 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. @@ -64,7 +63,7 @@ class BulkProcessorTest { @Test void decodesItemLevelFailureFromHttp200Response() throws Exception { - final BulkResponse response = MAPPER.readValue(BLOCKED_ITEM_RESPONSE, BulkResponse.class); + final BulkResponse response = decode(BLOCKED_ITEM_RESPONSE); assertTrue(response.isErrors()); assertEquals(1, response.getItems().size()); @@ -77,7 +76,7 @@ void decodesItemLevelFailureFromHttp200Response() throws Exception { @Test void decodesAllSucceededResponse() throws Exception { - final BulkResponse response = MAPPER.readValue(ALL_SUCCEEDED_RESPONSE, BulkResponse.class); + final BulkResponse response = decode(ALL_SUCCEEDED_RESPONSE); assertFalse(response.isErrors()); assertNull(response.getItems().get(0).values().iterator().next().getError()); @@ -85,7 +84,7 @@ void decodesAllSucceededResponse() throws Exception { @Test void completesAllHoldersWhenNoErrors() throws Exception { - final BulkResponse response = MAPPER.readValue(ALL_SUCCEEDED_RESPONSE, BulkResponse.class); + final BulkResponse response = decode(ALL_SUCCEEDED_RESPONSE); final CompletableFuture future = new CompletableFuture<>(); final List holders = holdersOf(future); @@ -96,7 +95,7 @@ void completesAllHoldersWhenNoErrors() throws Exception { } @Test - void failsOnlyTheItemThatIsRejectedByElasticsearch() throws Exception { + void failsOnlyTheItemsThatAreRejectedByElasticsearchWithASummaryMessage() throws Exception { final String twoItemResponse = "{" + "\"errors\":true," + "\"items\":[" @@ -104,7 +103,7 @@ void failsOnlyTheItemThatIsRejectedByElasticsearch() throws Exception { + " {\"update\":{\"_index\":\"idx\",\"_id\":\"2\",\"status\":429," + " \"error\":{\"type\":\"cluster_block_exception\",\"reason\":\"disk usage exceeded flood-stage watermark\"}}}" + "]}"; - final BulkResponse response = MAPPER.readValue(twoItemResponse, BulkResponse.class); + final BulkResponse response = decode(twoItemResponse); final CompletableFuture succeeded = new CompletableFuture<>(); final CompletableFuture failed = new CompletableFuture<>(); @@ -120,7 +119,30 @@ void failsOnlyTheItemThatIsRejectedByElasticsearch() throws Exception { assertTrue(failed.isDone()); assertTrue(failed.isCompletedExceptionally()); final ExecutionException ex = assertThrows(ExecutionException.class, failed::get); - assertTrue(ex.getCause().getMessage().contains("cluster_block_exception")); + 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("\"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 @@ -129,6 +151,12 @@ void itemResultOfReturnsNullWhenIndexOutOfBounds() { 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")); From 92e416348dfc872845691a0c8dde469af5cb2c16 Mon Sep 17 00:00:00 2001 From: kezhenxu94 Date: Tue, 29 Sep 2026 14:26:59 +0800 Subject: [PATCH 3/5] Address second review round: fix IT race, tighten id assertion, drop unused fields - ElasticSearchIT#testBulkProcessorItemLevelFailure asserted isDone() right after flush(), but BulkProcessor's own periodical-flush thread can race the test's explicit flush() and drain either request first. Wait on the futures with a timeout instead. - BulkProcessorTest's id-leak assertion checked for a quoted id that the message format never produces, so it could never fail. Give the rejected item a distinctive id and assert the message excludes it. - Remove BulkItemResult.index/id and BulkItemError.reason: nothing reads them now that failures are logged as a per-status/error-type summary, and keeping them around risks a future change quietly reintroducing an id or reason into a log line. - Trim the changes.md entry to the behavior change, dropping the "why" (volume/PR-description material). --- docs/en/changes/changes.md | 2 +- .../elasticsearch/response/bulk/BulkItemError.java | 1 - .../elasticsearch/response/bulk/BulkItemResult.java | 5 ----- .../library/elasticsearch/ElasticSearchIT.java | 12 +++++++----- .../elasticsearch/bulk/BulkProcessorTest.java | 6 ++---- 5 files changed, 10 insertions(+), 16 deletions(-) diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index 375c2ad82266..6e0d39842b72 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -29,7 +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, which can embed the id and, for mapping errors, a preview of the stored field value) — with the default `bulkActions` of thousands of requests per flush, logging per item would flood the log when every item of a bulk is rejected, as happens under a sustained disk watermark block. 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. +* 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/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 index 300dc34e25db..2f249f1c4046 100644 --- 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 @@ -24,5 +24,4 @@ @Setter public final class BulkItemError { private String type; - private String reason; } 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 index c9589841fce1..c120543802e9 100644 --- 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 @@ -17,7 +17,6 @@ package org.apache.skywalking.library.elasticsearch.response.bulk; -import com.fasterxml.jackson.annotation.JsonProperty; import lombok.Getter; import lombok.Setter; @@ -25,9 +24,5 @@ @Setter public final class BulkItemResult { private int status; - @JsonProperty("_index") - private String index; - @JsonProperty("_id") - private String id; private BulkItemError error; } 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 69bee847aa83..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 @@ -46,6 +46,7 @@ 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; @@ -528,13 +529,14 @@ public void observe(final double value) { bulkProcessor.flush(); - assertTrue(normalFuture.isDone()); - assertFalse(normalFuture.isCompletedExceptionally()); + // 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()); - assertTrue(blockedFuture.isDone()); - assertTrue(blockedFuture.isCompletedExceptionally()); - final String message = assertThrows(ExecutionException.class, blockedFuture::get).getCause().getMessage(); + 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 index 4177bfbd5eef..ae3482428fd5 100644 --- 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 @@ -69,8 +69,6 @@ void decodesItemLevelFailureFromHttp200Response() throws Exception { assertEquals(1, response.getItems().size()); final BulkItemResult item = response.getItems().get(0).values().iterator().next(); assertEquals(429, item.getStatus()); - assertEquals("blocked-test-index", item.getIndex()); - assertEquals("oap-block-check-1", item.getId()); assertEquals("cluster_block_exception", item.getError().getType()); } @@ -100,7 +98,7 @@ void failsOnlyTheItemsThatAreRejectedByElasticsearchWithASummaryMessage() throws + "\"errors\":true," + "\"items\":[" + " {\"index\":{\"_index\":\"idx\",\"_id\":\"1\",\"status\":201}}," - + " {\"update\":{\"_index\":\"idx\",\"_id\":\"2\",\"status\":429," + + " {\"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); @@ -122,7 +120,7 @@ void failsOnlyTheItemsThatAreRejectedByElasticsearchWithASummaryMessage() throws 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("\"2\""), "the rejected document id must not be logged: " + 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); } From 8a1cbe0c8d585d45d84a20bc81e94bc2f2b718c5 Mon Sep 17 00:00:00 2001 From: Wu Sheng Date: Tue, 29 Sep 2026 15:48:23 +0800 Subject: [PATCH 4/5] Point the Kong monitoring doc at the moved Prometheus plugin page docs.konghq.com/hub/kong-inc/prometheus/ now redirects to developer.konghq.com/plugins/prometheus/, and the dead-link checker fails on it. Link the operator doc to the new page directly. SWIP-8 still carries the old URL and is frozen since it shipped in 10.2.0, so the old URL is ignored in .dlc.json instead of editing it. --- .dlc.json | 3 +++ docs/en/setup/backend/backend-kong-monitoring.md | 8 ++++---- 2 files changed, 7 insertions(+), 4 deletions(-) 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/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`. From f7e0b997e3f7c304bbb630a15790f6e669a32f7a Mon Sep 17 00:00:00 2001 From: Wu Sheng Date: Tue, 29 Sep 2026 18:57:03 +0800 Subject: [PATCH 5/5] Wait for OAP before verifying the Event BanyanDB e2e case Nothing in this compose file depends on OAP's healthcheck, so verification starts as soon as the e2e runner reports OAP's ports ready. In CI that check passes within a second of OAP starting, while OAP opens 11800 and 12800 only after installing the BanyanDB schemas, which now takes longer than the 20 x 3s verify retries. Every retry got "connection reset by peer", and the case fails on master as well. Add a setup step that waits until OAP answers a GraphQL query, bounded at 300s. --- test/e2e-v2/cases/event/banyandb/e2e.yaml | 11 +++++++++++ 1 file changed, 11 insertions(+) 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: