Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .dlc.json
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,9 @@
},
{
"pattern": "^https://x.com/AsfSkyWalking"
},
{
"pattern": "^https://docs.konghq.com/hub/kong-inc/prometheus/$"
}
],
"timeout": "10s",
Expand Down
1 change: 1 addition & 0 deletions docs/en/changes/changes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/<stream>/execution-<stamp>-<seq>.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, `<server>/<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.
Expand Down
8 changes: 4 additions & 4 deletions docs/en/setup/backend/backend-kong-monitoring.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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`.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -161,23 +168,31 @@ private List<CompletableFuture<Void>> doFlush(final List<Holder> batch) {
final List<byte[]> bs = new ArrayList<>();
List<CompletableFuture<Void>> futures = new ArrayList<>();
List<ByteBuf> byteBufList = new ArrayList<>();
List<List<Holder>> holderChunks = new ArrayList<>();
List<Holder> 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<Holder> chunkHolders = holderChunks.get(i);
CompletableFuture<Void> future = es.get().version().thenCompose(v -> {
try {
final RequestFactory rf = v.requestFactory();
Expand All @@ -186,19 +201,22 @@ private List<CompletableFuture<Void>> doFlush(final List<Holder> 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);
}
});
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);
Expand All @@ -207,10 +225,56 @@ private List<CompletableFuture<Void>> doFlush(final List<Holder> 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<Holder> holders, final BulkResponse bulkResponse) {
if (!bulkResponse.isErrors()) {
holders.forEach(it -> it.future.complete(null));
return;
}

final Map<String, Integer> rejections = new LinkedHashMap<>();
final List<Holder> rejected = new ArrayList<>();
final List<Map<String, BulkItemResult>> 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<Map<String, BulkItemResult>> items, final int index) {
if (items == null || index >= items.size()) {
return null;
}
final Map<String, BulkItemResult> item = items.get(index);
if (item == null || item.isEmpty()) {
return null;
}
return item.values().iterator().next();
}

@RequiredArgsConstructor
static class Holder {
private final CompletableFuture<Void> future;
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
}
Original file line number Diff line number Diff line change
@@ -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;
}
Original file line number Diff line number Diff line change
@@ -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<Map<String, BulkItemResult>> items;
}
Loading
Loading