From 0c8a48ce1cb12160e7e6a61af98a2fd113bc32a0 Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Mon, 31 Aug 2026 15:51:30 +0200 Subject: [PATCH 1/9] Add JetStream data stream to NATS integration Add a new `jetstream` data stream that wraps the Metricbeat `nats/jetstream` metricset (available since Elastic Stack 9.1). The data stream collects stats, account, stream, and consumer metrics from the NATS `/jsz` monitoring endpoint, with optional name filters for accounts, streams, and consumers. - New data stream `packages/nats/data_stream/jetstream/` modeled after the existing `connection` data stream, with TSDB field mappings (metric_type, dimensions). - Package version bumped to 1.13.0; Kibana constraint raised to ^9.1.0 because the jetstream metricset is only in Agent 9.1+. - Docker test setup updated to NATS 2.10.27 with `-js` enabled. - Docs, changelog, and sample event updated. Closes elastic/integrations#10748 Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- packages/nats/_dev/build/docs/README.md | 30 +- packages/nats/_dev/deploy/docker/Dockerfile | 12 +- packages/nats/_dev/deploy/docker/run.sh | 6 +- packages/nats/_dev/deploy/variants.yml | 6 +- packages/nats/changelog.yml | 5 + .../data_stream/connection/fields/fields.yml | 6 + .../_dev/test/system/test-default-config.yml | 5 + .../jetstream/agent/stream/stream.yml.hbs | 33 ++ .../jetstream/fields/base-fields.yml | 20 + .../nats/data_stream/jetstream/fields/ecs.yml | 27 + .../data_stream/jetstream/fields/fields.yml | 473 ++++++++++++++++++ .../jetstream/fields/package-fields.yml | 12 + .../nats/data_stream/jetstream/manifest.yml | 70 +++ .../data_stream/jetstream/sample_event.json | 86 ++++ .../nats/data_stream/stats/fields/fields.yml | 3 + packages/nats/docs/README.md | 220 +++++++- packages/nats/manifest.yml | 6 +- 17 files changed, 998 insertions(+), 22 deletions(-) create mode 100644 packages/nats/data_stream/jetstream/_dev/test/system/test-default-config.yml create mode 100644 packages/nats/data_stream/jetstream/agent/stream/stream.yml.hbs create mode 100644 packages/nats/data_stream/jetstream/fields/base-fields.yml create mode 100644 packages/nats/data_stream/jetstream/fields/ecs.yml create mode 100644 packages/nats/data_stream/jetstream/fields/fields.yml create mode 100644 packages/nats/data_stream/jetstream/fields/package-fields.yml create mode 100644 packages/nats/data_stream/jetstream/manifest.yml create mode 100644 packages/nats/data_stream/jetstream/sample_event.json diff --git a/packages/nats/_dev/build/docs/README.md b/packages/nats/_dev/build/docs/README.md index 7d8a4b27ded..bb99329da4d 100644 --- a/packages/nats/_dev/build/docs/README.md +++ b/packages/nats/_dev/build/docs/README.md @@ -6,7 +6,7 @@ The integration collects metrics from [NATS monitoring server APIs](https://docs ## Compatibility -The Nats package is tested with Nats 1.3.0, 2.0.4 and 2.1.4 +The Nats package is tested with NATS 2.10.27. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+) and Elastic Agent 9.1+. Consumer metrics require JetStream 2.9+. ## Logs @@ -24,8 +24,8 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur ## Metrics -The default datasets are `stats`, `connections`, `routes` and `subscriptions` while `connection` and `route` -datasets can be enabled to collect detailed metrics per connection/route. +The default datasets are `stats`, `connections`, `routes` and `subscriptions` while `connection`, `route` +and `jetstream` datasets can be enabled to collect detailed metrics per connection/route and JetStream monitoring. ### stats @@ -93,6 +93,30 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur {{fields "connection"}} +### jetstream + +This is the `jetstream` dataset of the Nats package, in charge of retrieving +JetStream metrics from a NATS server. It collects data from the [/jsz](https://docs.nats.io/running-a-nats-service/nats_admin/monitoring#jetstream-information-jsz) monitoring endpoint. + +The `jetstream` dataset supports four categories of metrics that can be enabled independently: + +* `stats` — General JetStream server stats (streams, consumers, messages, memory, storage). +* `account` — Per-account JetStream metrics (memory, storage, API stats). +* `stream` — Per-stream metrics (state, config, cluster info). +* `consumer` — Per-consumer metrics (delivered, ack floor, pending, config). Requires JetStream 2.9+. + +Account, stream, and consumer metrics can be filtered by name. Filters are cumulative and apply even if a category is not enabled but name filters are configured. When no names are configured, all entities are reported. + +This dataset requires Elastic Agent 9.1+ and a NATS server with JetStream enabled. + +{{event "jetstream"}} + +**ECS Field Reference** + +Please refer to the following [document](https://www.elastic.co/guide/en/ecs/current/ecs-field-reference.html) for detailed information on ECS fields. + +{{fields "jetstream"}} + ### route This is the `route` dataset of the Nats package, in charge of retrieving detailed diff --git a/packages/nats/_dev/deploy/docker/Dockerfile b/packages/nats/_dev/deploy/docker/Dockerfile index cbdc0f698e1..2033d6d431d 100644 --- a/packages/nats/_dev/deploy/docker/Dockerfile +++ b/packages/nats/_dev/deploy/docker/Dockerfile @@ -1,19 +1,19 @@ -ARG NATS_VERSION=2.0.4 +ARG NATS_VERSION=2.10.27 FROM nats:$NATS_VERSION # build stage -FROM golang:1.20.2-alpine3.17 AS build-env -RUN apk --no-cache add build-base git mercurial gcc -RUN go install github.com/nats-io/nats.go/examples/nats-bench@v1.10.0 +FROM golang:1.25-alpine AS build-env +RUN apk --no-cache add build-base git gcc +RUN go install github.com/nats-io/nats.go/examples/nats-bench@v1.36.0 # create an enhanced container with nc command available since nats is based # on scratch image making healthcheck impossible -FROM alpine:3.17 +FROM alpine:3.20 COPY --from=0 / /opt/nats COPY --from=build-env /go/bin/nats-bench /nats-bench COPY run.sh /run.sh # Expose client, management, and cluster ports EXPOSE 4222 8222 6222 -HEALTHCHECK --interval=1s --retries=10 CMD nc -w 1 0.0.0.0 8222 Date: Tue, 1 Sep 2026 06:28:58 +0200 Subject: [PATCH 2/9] fix(nats): align JetStream field definitions, docs, and stack version --- packages/nats/data_stream/jetstream/fields/fields.yml | 9 +++------ packages/nats/docs/README.md | 6 +++--- packages/nats/manifest.yml | 2 +- 3 files changed, 7 insertions(+), 10 deletions(-) diff --git a/packages/nats/data_stream/jetstream/fields/fields.yml b/packages/nats/data_stream/jetstream/fields/fields.yml index 7d6bce8da83..bf9fd5bda4b 100644 --- a/packages/nats/data_stream/jetstream/fields/fields.yml +++ b/packages/nats/data_stream/jetstream/fields/fields.yml @@ -43,7 +43,7 @@ format: bytes metric_type: gauge description: | - The of memory (bytes) reserved by the JetStream server. + The amount of memory (bytes) reserved by the JetStream server. - name: storage type: long format: bytes @@ -80,7 +80,6 @@ The maximum amount of storage (bytes) the JetStream server can use. - name: store_dir type: keyword - dimension: true description: | The path on disk where the JetStream storage lives. - name: sync_interval @@ -104,7 +103,7 @@ description: | The name of the JetStream account. - name: accounts - type: integer + type: long metric_type: gauge description: | The number of accounts using JetStream on the server. @@ -173,7 +172,6 @@ fields: - name: leader type: keyword - dimension: true description: | The ID of the leader in the cluster. - name: state @@ -255,7 +253,7 @@ description: | The retention policy for the stream. - name: num_replicas - type: integer + type: long metric_type: gauge description: | How many replicas to keep for each message in a clustered JetStream. @@ -331,7 +329,6 @@ fields: - name: leader type: keyword - dimension: true description: | The ID of the leader in the cluster. - name: ack_floor diff --git a/packages/nats/docs/README.md b/packages/nats/docs/README.md index e91ef05d9f8..d89cd737a7d 100644 --- a/packages/nats/docs/README.md +++ b/packages/nats/docs/README.md @@ -892,7 +892,7 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur | event.dataset | Event dataset | constant_keyword | | | event.module | Event module | constant_keyword | | | host.name | Name of the host. It can contain what hostname returns on Unix systems, the fully qualified domain name (FQDN), or a name specified by the user. The recommended value is the lowercase FQDN of the host. | keyword | | -| nats.jetstream.account.accounts | The number of accounts using JetStream on the server. | integer | gauge | +| nats.jetstream.account.accounts | The number of accounts using JetStream on the server. | long | gauge | | nats.jetstream.account.api.errors | The total number of JetStream API errors encountered by this account. | long | counter | | nats.jetstream.account.api.total | The total number of JetStream API calls made by this account. | long | counter | | nats.jetstream.account.high_availability_assets | Indicates the number of JetStream high-availability (HA) assets allocated for an account. | integer | gauge | @@ -940,7 +940,7 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur | nats.jetstream.stats.consumers | The total number of consumers on the JetStream server. | long | gauge | | nats.jetstream.stats.memory | The total amount of memory (bytes) used by the JetStream server. | long | gauge | | nats.jetstream.stats.messages | The total number of messages on the JetStream server. | long | gauge | -| nats.jetstream.stats.reserved_memory | The of memory (bytes) reserved by the JetStream server. | long | gauge | +| nats.jetstream.stats.reserved_memory | The amount of memory (bytes) reserved by the JetStream server. | long | gauge | | nats.jetstream.stats.reserved_storage | The total amount of storage (bytes) reserved by the JetStream server. | long | gauge | | nats.jetstream.stats.storage | The total amount of storage (bytes) used by the JetStream server. | long | gauge | | nats.jetstream.stats.streams | The total number of streams on the JetStream server. | long | gauge | @@ -954,7 +954,7 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur | nats.jetstream.stream.config.max_msg_size | The largest message (bytes) that will be accepted by the stream. The size of a message is a sum of payload and headers. | long | gauge | | nats.jetstream.stream.config.max_msgs | Maximum number of messages stored in the stream. Adheres to Discard Policy, removing oldest or refusing new messages if the Stream exceeds this number of messages. | long | gauge | | nats.jetstream.stream.config.max_msgs_per_subject | Limits maximum number of messages in the stream to retain per subject. | long | gauge | -| nats.jetstream.stream.config.num_replicas | How many replicas to keep for each message in a clustered JetStream. | integer | gauge | +| nats.jetstream.stream.config.num_replicas | How many replicas to keep for each message in a clustered JetStream. | long | gauge | | nats.jetstream.stream.config.retention | The retention policy for the stream. | keyword | | | nats.jetstream.stream.config.storage | The storage type for stream data. | keyword | | | nats.jetstream.stream.config.subjects | The list of subjects bound to the stream. | keyword | | diff --git a/packages/nats/manifest.yml b/packages/nats/manifest.yml index 808ab9d92ee..924cf9767ac 100644 --- a/packages/nats/manifest.yml +++ b/packages/nats/manifest.yml @@ -16,7 +16,7 @@ categories: - stream_processing conditions: kibana: - version: "^9.1.0" + version: "^8.13.0 || ^9.0.0" elastic: subscription: basic screenshots: From ff4bf95395c0fe146dd2f2938c1a4e8746b4ac31 Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Tue, 1 Sep 2026 06:38:23 +0200 Subject: [PATCH 3/9] docs(nats): update JetStream monitoring documentation URL --- packages/nats/_dev/build/docs/README.md | 2 +- packages/nats/docs/README.md | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/nats/_dev/build/docs/README.md b/packages/nats/_dev/build/docs/README.md index bb99329da4d..1d024801b96 100644 --- a/packages/nats/_dev/build/docs/README.md +++ b/packages/nats/_dev/build/docs/README.md @@ -96,7 +96,7 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur ### jetstream This is the `jetstream` dataset of the Nats package, in charge of retrieving -JetStream metrics from a NATS server. It collects data from the [/jsz](https://docs.nats.io/running-a-nats-service/nats_admin/monitoring#jetstream-information-jsz) monitoring endpoint. +JetStream metrics from a NATS server. It collects data from the [/jsz](https://docs.nats.io/learn/monitoring/monitoring-endpoints#jsz-reports-jetstream-state) monitoring endpoint. The `jetstream` dataset supports four categories of metrics that can be enabled independently: diff --git a/packages/nats/docs/README.md b/packages/nats/docs/README.md index d89cd737a7d..cec9722c662 100644 --- a/packages/nats/docs/README.md +++ b/packages/nats/docs/README.md @@ -766,7 +766,7 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur ### jetstream This is the `jetstream` dataset of the Nats package, in charge of retrieving -JetStream metrics from a NATS server. It collects data from the [/jsz](https://docs.nats.io/running-a-nats-service/nats_admin/monitoring#jetstream-information-jsz) monitoring endpoint. +JetStream metrics from a NATS server. It collects data from the [/jsz](https://docs.nats.io/learn/monitoring/monitoring-endpoints#jsz-reports-jetstream-state) monitoring endpoint. The `jetstream` dataset supports four categories of metrics that can be enabled independently: From d2534c993867dbb2d94bb216e864e69a80abb95c Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Tue, 1 Sep 2026 10:37:26 +0200 Subject: [PATCH 4/9] fix(nats): use POSIX test in run.sh and reduce TSDB dimensions --- packages/nats/_dev/deploy/docker/Dockerfile | 4 + .../_dev/deploy/docker/jetstream-traffic.go | 118 ++++++++++++++++++ packages/nats/_dev/deploy/docker/run.sh | 7 +- .../_dev/test/system/test-default-config.yml | 6 +- .../data_stream/jetstream/fields/fields.yml | 14 --- 5 files changed, 132 insertions(+), 17 deletions(-) create mode 100644 packages/nats/_dev/deploy/docker/jetstream-traffic.go diff --git a/packages/nats/_dev/deploy/docker/Dockerfile b/packages/nats/_dev/deploy/docker/Dockerfile index 2033d6d431d..57ba49ff370 100644 --- a/packages/nats/_dev/deploy/docker/Dockerfile +++ b/packages/nats/_dev/deploy/docker/Dockerfile @@ -5,12 +5,16 @@ FROM nats:$NATS_VERSION FROM golang:1.25-alpine AS build-env RUN apk --no-cache add build-base git gcc RUN go install github.com/nats-io/nats.go/examples/nats-bench@v1.36.0 +WORKDIR /src +COPY jetstream-traffic.go . +RUN go mod init jetstream-traffic && go get github.com/nats-io/nats.go@v1.36.0 && go mod tidy && CGO_ENABLED=0 go build -o /jetstream-traffic jetstream-traffic.go # create an enhanced container with nc command available since nats is based # on scratch image making healthcheck impossible FROM alpine:3.20 COPY --from=0 / /opt/nats COPY --from=build-env /go/bin/nats-bench /nats-bench +COPY --from=build-env /jetstream-traffic /jetstream-traffic COPY run.sh /run.sh # Expose client, management, and cluster ports EXPOSE 4222 8222 6222 diff --git a/packages/nats/_dev/deploy/docker/jetstream-traffic.go b/packages/nats/_dev/deploy/docker/jetstream-traffic.go new file mode 100644 index 00000000000..8b9cf9aeeef --- /dev/null +++ b/packages/nats/_dev/deploy/docker/jetstream-traffic.go @@ -0,0 +1,118 @@ +package main + +import ( + "context" + "fmt" + "log" + "time" + + "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" +) + +func main() { + var nc *nats.Conn + var err error + + // Wait for NATS to become ready + for i := 0; i < 30; i++ { + nc, err = nats.Connect("nats://127.0.0.1:4222") + if err == nil { + break + } + time.Sleep(1 * time.Second) + } + if err != nil { + log.Fatalf("Failed to connect to NATS: %v", err) + } + defer nc.Close() + + var js jetstream.JetStream + js, err = jetstream.New(nc) + if err != nil { + log.Fatalf("Failed to create JetStream context: %v", err) + } + + cfg := jetstream.StreamConfig{ + Name: "ORDERS", + Description: "Orders stream for system testing", + Subjects: []string{"orders.*"}, + Storage: jetstream.FileStorage, + Retention: jetstream.LimitsPolicy, + MaxMsgs: 10000, + MaxBytes: 10485760, + MaxAge: 24 * time.Hour, + MaxMsgSize: 1048576, + Replicas: 1, + } + + var stream jetstream.Stream + for i := 0; i < 30; i++ { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + stream, err = js.CreateOrUpdateStream(ctx, cfg) + cancel() + if err == nil { + fmt.Printf("Stream created: %s\n", stream.CachedInfo().Config.Name) + break + } + log.Printf("CreateOrUpdateStream attempt %d failed: %v, retrying...", i+1, err) + time.Sleep(1 * time.Second) + } + + consumerCfg := jetstream.ConsumerConfig{ + Name: "PROCESSOR", + Durable: "PROCESSOR", + Description: "Order processor consumer", + DeliverPolicy: jetstream.DeliverAllPolicy, + AckPolicy: jetstream.AckExplicitPolicy, + AckWait: 30 * time.Second, + MaxDeliver: 5, + FilterSubject: "orders.*", + ReplayPolicy: jetstream.ReplayInstantPolicy, + MaxAckPending: 1000, + } + + var consumer jetstream.Consumer + for i := 0; i < 30; i++ { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + consumer, err = js.CreateOrUpdateConsumer(ctx, "ORDERS", consumerCfg) + cancel() + if err == nil { + fmt.Printf("Consumer created: %s\n", consumer.CachedInfo().Config.Name) + break + } + log.Printf("CreateOrUpdateConsumer attempt %d failed: %v, retrying...", i+1, err) + time.Sleep(1 * time.Second) + } + + // Initial batch of messages + for i := 0; i < 20; i++ { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + msg := fmt.Sprintf("order data %d", i) + _, _ = js.Publish(ctx, "orders.created", []byte(msg)) + cancel() + } + + // Continuous loop + ticker := time.NewTicker(200 * time.Millisecond) + count := 20 + for range ticker.C { + count++ + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + msg := fmt.Sprintf("order data %d", count) + _, err := js.Publish(ctx, "orders.created", []byte(msg)) + cancel() + if err != nil { + continue + } + + if count%2 == 0 && consumer != nil { + msgs, err := consumer.Fetch(1, jetstream.FetchMaxWait(50*time.Millisecond)) + if err == nil { + for m := range msgs.Messages() { + _ = m.Ack() + } + } + } + } +} diff --git a/packages/nats/_dev/deploy/docker/run.sh b/packages/nats/_dev/deploy/docker/run.sh index 2d0c26ff762..5b74e22bc02 100755 --- a/packages/nats/_dev/deploy/docker/run.sh +++ b/packages/nats/_dev/deploy/docker/run.sh @@ -9,17 +9,20 @@ mkdir -p /var/log/nats # NATS 2.X if [ -x /opt/nats/nats-server ]; then - if [[ -z "${ROUTES}" ]]; then + if [ -z "${ROUTES}" ]; then (/opt/nats/nats-server -DV -js --server_name nats --cluster_name nats-cluster -l /var/log/nats/nats.log --cluster nats://0.0.0.0:6222 --http_port 8222 --port 4222 --routes nats://nats-routes:6222) & else (/opt/nats/nats-server -DV -js --server_name nats-routes --cluster_name nats-cluster -l /var/log/nats/nats.log --cluster nats://0.0.0.0:6222 --http_port 8222 --port 4222 --routes nats://nats:6222) & fi + if [ -x /jetstream-traffic ]; then + /jetstream-traffic & + fi while true; do /nats-bench -np 1 -n 100000000 -ms 16 foo; done fi # NATS 1.X if [ -x /opt/nats/gnatsd ]; then - if [[ -z "${ROUTES}" ]]; then + if [ -z "${ROUTES}" ]; then (/opt/nats/gnatsd -DV -l /var/log/nats/nats.log --cluster nats://0.0.0.0:6222 --http_port 8222 --port 4222) & else (/opt/nats/gnatsd -DV -l /var/log/nats/nats.log --cluster nats://0.0.0.0:6222 --http_port 8222 --port 4222 --routes nats://nats:6222) & diff --git a/packages/nats/data_stream/jetstream/_dev/test/system/test-default-config.yml b/packages/nats/data_stream/jetstream/_dev/test/system/test-default-config.yml index 83fc84f62ea..31ac4622ca6 100644 --- a/packages/nats/data_stream/jetstream/_dev/test/system/test-default-config.yml +++ b/packages/nats/data_stream/jetstream/_dev/test/system/test-default-config.yml @@ -2,4 +2,8 @@ vars: hosts: - http://{{Hostname}}:{{Port}} data_stream: - vars: ~ + vars: + jetstream_stats_enabled: true + jetstream_account_enabled: true + jetstream_stream_enabled: true + jetstream_consumer_enabled: true diff --git a/packages/nats/data_stream/jetstream/fields/fields.yml b/packages/nats/data_stream/jetstream/fields/fields.yml index bf9fd5bda4b..7f8e6c386a1 100644 --- a/packages/nats/data_stream/jetstream/fields/fields.yml +++ b/packages/nats/data_stream/jetstream/fields/fields.yml @@ -94,7 +94,6 @@ fields: - name: id type: keyword - dimension: true description: | The ID of the JetStream account. - name: name @@ -230,12 +229,10 @@ fields: - name: id type: keyword - dimension: true description: | The ID of the account. - name: name type: keyword - dimension: true description: | The name of the account. - name: config @@ -249,7 +246,6 @@ The description of the stream. - name: retention type: keyword - dimension: true description: | The retention policy for the stream. - name: num_replicas @@ -259,7 +255,6 @@ How many replicas to keep for each message in a clustered JetStream. - name: storage type: keyword - dimension: true description: | The storage type for stream data. - name: max_consumers @@ -319,7 +314,6 @@ fields: - name: name type: keyword - dimension: true description: | The name of the stream. - name: cluster @@ -400,12 +394,10 @@ fields: - name: id type: keyword - dimension: true description: | The ID of the account. - name: name type: keyword - dimension: true description: | The name of the account. - name: config @@ -415,32 +407,26 @@ fields: - name: name type: keyword - dimension: true description: | The name of the consumer. - name: durable_name type: keyword - dimension: true description: | The durable name of the consumer. If set, clients can have subscriptions bind to the consumer and resume until the consumer is explicitly deleted. - name: deliver_policy type: keyword - dimension: true description: | The point in the stream from which to receive messages. - name: filter_subject type: keyword - dimension: true description: | A subject that overlaps with the subjects bound to the stream to filter delivery to subscribers. - name: replay_policy type: keyword - dimension: true description: | The configured replay policy for the consumer. - name: ack_policy type: keyword - dimension: true description: | The configured ack policy for the consumer. - name: ack_wait From aa863749f7a92e5f8d997c5bbb114665cb6c111f Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Tue, 1 Sep 2026 10:56:28 +0200 Subject: [PATCH 5/9] fix(nats): add entity scoping dimensions, bump TSDB limit, and fail fast in traffic generator --- packages/nats/_dev/deploy/docker/jetstream-traffic.go | 6 ++++++ packages/nats/data_stream/jetstream/fields/fields.yml | 3 +++ packages/nats/data_stream/jetstream/manifest.yml | 8 +++++++- 3 files changed, 16 insertions(+), 1 deletion(-) diff --git a/packages/nats/_dev/deploy/docker/jetstream-traffic.go b/packages/nats/_dev/deploy/docker/jetstream-traffic.go index 8b9cf9aeeef..befb971790f 100644 --- a/packages/nats/_dev/deploy/docker/jetstream-traffic.go +++ b/packages/nats/_dev/deploy/docker/jetstream-traffic.go @@ -58,6 +58,9 @@ func main() { log.Printf("CreateOrUpdateStream attempt %d failed: %v, retrying...", i+1, err) time.Sleep(1 * time.Second) } + if err != nil { + log.Fatalf("Failed to create stream after retries: %v", err) + } consumerCfg := jetstream.ConsumerConfig{ Name: "PROCESSOR", @@ -84,6 +87,9 @@ func main() { log.Printf("CreateOrUpdateConsumer attempt %d failed: %v, retrying...", i+1, err) time.Sleep(1 * time.Second) } + if err != nil { + log.Fatalf("Failed to create consumer after retries: %v", err) + } // Initial batch of messages for i := 0; i < 20; i++ { diff --git a/packages/nats/data_stream/jetstream/fields/fields.yml b/packages/nats/data_stream/jetstream/fields/fields.yml index 7f8e6c386a1..373d5b2aac9 100644 --- a/packages/nats/data_stream/jetstream/fields/fields.yml +++ b/packages/nats/data_stream/jetstream/fields/fields.yml @@ -233,6 +233,7 @@ The ID of the account. - name: name type: keyword + dimension: true description: | The name of the account. - name: config @@ -314,6 +315,7 @@ fields: - name: name type: keyword + dimension: true description: | The name of the stream. - name: cluster @@ -398,6 +400,7 @@ The ID of the account. - name: name type: keyword + dimension: true description: | The name of the account. - name: config diff --git a/packages/nats/data_stream/jetstream/manifest.yml b/packages/nats/data_stream/jetstream/manifest.yml index 23dd2018742..b658bec7cdc 100644 --- a/packages/nats/data_stream/jetstream/manifest.yml +++ b/packages/nats/data_stream/jetstream/manifest.yml @@ -65,6 +65,12 @@ streams: description: Filter consumer metrics by name. When empty, all consumers are reported. title: NATS JetStream metrics enabled: false - description: Collect JetStream metrics (stats, account, stream, consumer) from NATS servers + description: Collect JetStream metrics (stats, account, stream, consumer) from NATS servers. Requires Elastic Agent 9.1+. elasticsearch: index_mode: "time_series" + index_template: + settings: + index: + mapping: + dimension_fields: + limit: 24 From 2de06f17468326573e48a364bd1d6971cf7c8725 Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Tue, 1 Sep 2026 11:04:48 +0200 Subject: [PATCH 6/9] fix(nats): add duration formats, enforce agent condition, and update NATS version wording --- packages/nats/_dev/build/docs/README.md | 4 ++-- packages/nats/data_stream/jetstream/fields/fields.yml | 3 +++ packages/nats/data_stream/jetstream/manifest.yml | 2 +- packages/nats/docs/README.md | 4 ++-- packages/nats/manifest.yml | 2 ++ 5 files changed, 10 insertions(+), 5 deletions(-) diff --git a/packages/nats/_dev/build/docs/README.md b/packages/nats/_dev/build/docs/README.md index 1d024801b96..359ee521b5f 100644 --- a/packages/nats/_dev/build/docs/README.md +++ b/packages/nats/_dev/build/docs/README.md @@ -6,7 +6,7 @@ The integration collects metrics from [NATS monitoring server APIs](https://docs ## Compatibility -The Nats package is tested with NATS 2.10.27. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+) and Elastic Agent 9.1+. Consumer metrics require JetStream 2.9+. +The Nats package is tested with NATS 2.10.27. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+) and Elastic Agent 9.1+. Consumer metrics require NATS 2.9+. ## Logs @@ -103,7 +103,7 @@ The `jetstream` dataset supports four categories of metrics that can be enabled * `stats` — General JetStream server stats (streams, consumers, messages, memory, storage). * `account` — Per-account JetStream metrics (memory, storage, API stats). * `stream` — Per-stream metrics (state, config, cluster info). -* `consumer` — Per-consumer metrics (delivered, ack floor, pending, config). Requires JetStream 2.9+. +* `consumer` — Per-consumer metrics (delivered, ack floor, pending, config). Requires NATS 2.9+. Account, stream, and consumer metrics can be filtered by name. Filters are cumulative and apply even if a category is not enabled but name filters are configured. When no names are configured, all entities are reported. diff --git a/packages/nats/data_stream/jetstream/fields/fields.yml b/packages/nats/data_stream/jetstream/fields/fields.yml index 373d5b2aac9..bf41966767d 100644 --- a/packages/nats/data_stream/jetstream/fields/fields.yml +++ b/packages/nats/data_stream/jetstream/fields/fields.yml @@ -84,6 +84,7 @@ The path on disk where the JetStream storage lives. - name: sync_interval type: long + format: duration metric_type: gauge description: | The fsync/sync interval for page cache in the filestore. @@ -276,6 +277,7 @@ Maximum number of bytes stored in the stream. Adheres to Discard Policy, removing oldest or refusing new messages if the Stream exceeds this size. - name: max_age type: long + format: duration metric_type: gauge description: | Maximum age of any message in the stream, expressed in nanoseconds. @@ -434,6 +436,7 @@ The configured ack policy for the consumer. - name: ack_wait type: long + format: duration metric_type: gauge description: | The duration (in nanoseconds) that the server will wait for an acknowledgment for any individual message once it has been delivered to a consumer. If an acknowledgment is not received in time, the message will be redelivered. diff --git a/packages/nats/data_stream/jetstream/manifest.yml b/packages/nats/data_stream/jetstream/manifest.yml index b658bec7cdc..089e5b57376 100644 --- a/packages/nats/data_stream/jetstream/manifest.yml +++ b/packages/nats/data_stream/jetstream/manifest.yml @@ -55,7 +55,7 @@ streams: required: false show_user: true default: false - description: Enable collection of JetStream consumer metrics. Requires JetStream 2.9+. + description: Enable collection of JetStream consumer metrics. Requires NATS 2.9+. - name: jetstream_consumer_names type: text title: JetStream consumer names filter diff --git a/packages/nats/docs/README.md b/packages/nats/docs/README.md index cec9722c662..f54d61e20ea 100644 --- a/packages/nats/docs/README.md +++ b/packages/nats/docs/README.md @@ -6,7 +6,7 @@ The integration collects metrics from [NATS monitoring server APIs](https://docs ## Compatibility -The Nats package is tested with NATS 2.10.27. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+) and Elastic Agent 9.1+. Consumer metrics require JetStream 2.9+. +The Nats package is tested with NATS 2.10.27. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+) and Elastic Agent 9.1+. Consumer metrics require NATS 2.9+. ## Logs @@ -773,7 +773,7 @@ The `jetstream` dataset supports four categories of metrics that can be enabled * `stats` — General JetStream server stats (streams, consumers, messages, memory, storage). * `account` — Per-account JetStream metrics (memory, storage, API stats). * `stream` — Per-stream metrics (state, config, cluster info). -* `consumer` — Per-consumer metrics (delivered, ack floor, pending, config). Requires JetStream 2.9+. +* `consumer` — Per-consumer metrics (delivered, ack floor, pending, config). Requires NATS 2.9+. Account, stream, and consumer metrics can be filtered by name. Filters are cumulative and apply even if a category is not enabled but name filters are configured. When no names are configured, all entities are reported. diff --git a/packages/nats/manifest.yml b/packages/nats/manifest.yml index 924cf9767ac..f2a9e82ece7 100644 --- a/packages/nats/manifest.yml +++ b/packages/nats/manifest.yml @@ -17,6 +17,8 @@ categories: conditions: kibana: version: "^8.13.0 || ^9.0.0" + agent: + version: "^9.1.0" elastic: subscription: basic screenshots: From bec8a2b182941dde5a508f0035b6e34cb8015660 Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Tue, 1 Sep 2026 11:13:57 +0200 Subject: [PATCH 7/9] fix(nats): align Kibana version floor, fix ack_floor descriptions, and add publish error logging --- packages/nats/_dev/deploy/docker/jetstream-traffic.go | 5 ++++- packages/nats/data_stream/jetstream/fields/fields.yml | 4 ++-- packages/nats/docs/README.md | 4 ++-- packages/nats/manifest.yml | 2 +- 4 files changed, 9 insertions(+), 6 deletions(-) diff --git a/packages/nats/_dev/deploy/docker/jetstream-traffic.go b/packages/nats/_dev/deploy/docker/jetstream-traffic.go index befb971790f..afba8bcb65b 100644 --- a/packages/nats/_dev/deploy/docker/jetstream-traffic.go +++ b/packages/nats/_dev/deploy/docker/jetstream-traffic.go @@ -95,7 +95,9 @@ func main() { for i := 0; i < 20; i++ { ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) msg := fmt.Sprintf("order data %d", i) - _, _ = js.Publish(ctx, "orders.created", []byte(msg)) + if _, err := js.Publish(ctx, "orders.created", []byte(msg)); err != nil { + log.Printf("Initial publish error for message %d: %v", i, err) + } cancel() } @@ -109,6 +111,7 @@ func main() { _, err := js.Publish(ctx, "orders.created", []byte(msg)) cancel() if err != nil { + log.Printf("Publish error: %v", err) continue } diff --git a/packages/nats/data_stream/jetstream/fields/fields.yml b/packages/nats/data_stream/jetstream/fields/fields.yml index bf41966767d..c5fa063db25 100644 --- a/packages/nats/data_stream/jetstream/fields/fields.yml +++ b/packages/nats/data_stream/jetstream/fields/fields.yml @@ -338,12 +338,12 @@ type: long metric_type: gauge description: | - The lowest contiguous consumer sequence number that has been acknowledged. + The highest contiguous consumer sequence number that has been acknowledged. - name: stream_seq type: long metric_type: gauge description: | - The lowest contiguous stream sequence number that has been acknowledged by the consumer. + The highest contiguous stream sequence number that has been acknowledged by the consumer. - name: last_active type: date description: | diff --git a/packages/nats/docs/README.md b/packages/nats/docs/README.md index f54d61e20ea..24db81efa32 100644 --- a/packages/nats/docs/README.md +++ b/packages/nats/docs/README.md @@ -905,9 +905,9 @@ Please refer to the following [document](https://www.elastic.co/guide/en/ecs/cur | nats.jetstream.category | The category of metrics represented in this event (stats, account, stream, or consumer). | keyword | | | nats.jetstream.consumer.account.id | The ID of the account. | keyword | | | nats.jetstream.consumer.account.name | The name of the account. | keyword | | -| nats.jetstream.consumer.ack_floor.consumer_seq | The lowest contiguous consumer sequence number that has been acknowledged. | long | gauge | +| nats.jetstream.consumer.ack_floor.consumer_seq | The highest contiguous consumer sequence number that has been acknowledged. | long | gauge | | nats.jetstream.consumer.ack_floor.last_active | The timestamp of the last acknowledged message. | date | | -| nats.jetstream.consumer.ack_floor.stream_seq | The lowest contiguous stream sequence number that has been acknowledged by the consumer. | long | gauge | +| nats.jetstream.consumer.ack_floor.stream_seq | The highest contiguous stream sequence number that has been acknowledged by the consumer. | long | gauge | | nats.jetstream.consumer.cluster.leader | The ID of the leader in the cluster. | keyword | | | nats.jetstream.consumer.config.ack_policy | The configured ack policy for the consumer. | keyword | | | nats.jetstream.consumer.config.ack_wait | The duration (in nanoseconds) that the server will wait for an acknowledgment for any individual message once it has been delivered to a consumer. If an acknowledgment is not received in time, the message will be redelivered. | long | gauge | diff --git a/packages/nats/manifest.yml b/packages/nats/manifest.yml index f2a9e82ece7..f4ed3dec717 100644 --- a/packages/nats/manifest.yml +++ b/packages/nats/manifest.yml @@ -16,7 +16,7 @@ categories: - stream_processing conditions: kibana: - version: "^8.13.0 || ^9.0.0" + version: "^9.1.0" agent: version: "^9.1.0" elastic: From 3291d1bfa35f7793cc2038ae4efba0a9501abf03 Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Tue, 1 Sep 2026 11:20:17 +0200 Subject: [PATCH 8/9] docs(nats): clarify package compatibility and update changelog link to PR --- packages/nats/_dev/build/docs/README.md | 2 +- packages/nats/changelog.yml | 2 +- packages/nats/docs/README.md | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/nats/_dev/build/docs/README.md b/packages/nats/_dev/build/docs/README.md index 359ee521b5f..7d8d27953c5 100644 --- a/packages/nats/_dev/build/docs/README.md +++ b/packages/nats/_dev/build/docs/README.md @@ -6,7 +6,7 @@ The integration collects metrics from [NATS monitoring server APIs](https://docs ## Compatibility -The Nats package is tested with NATS 2.10.27. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+) and Elastic Agent 9.1+. Consumer metrics require NATS 2.9+. +The NATS package is tested with NATS 2.10.27 and requires Elastic Agent 9.1+. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+), and consumer metrics require NATS 2.9+. ## Logs diff --git a/packages/nats/changelog.yml b/packages/nats/changelog.yml index 48715895508..50618109139 100644 --- a/packages/nats/changelog.yml +++ b/packages/nats/changelog.yml @@ -3,7 +3,7 @@ changes: - description: Add JetStream data stream for monitoring JetStream stats, accounts, streams, and consumers. Requires Elastic Agent 9.1+. type: enhancement - link: https://github.com/elastic/integrations/issues/10748 + link: https://github.com/elastic/integrations/pull/20980 - version: "1.12.0" changes: - description: Add missing fields to `connection` and `stats` data datastreams. diff --git a/packages/nats/docs/README.md b/packages/nats/docs/README.md index 24db81efa32..36b6d8a7390 100644 --- a/packages/nats/docs/README.md +++ b/packages/nats/docs/README.md @@ -6,7 +6,7 @@ The integration collects metrics from [NATS monitoring server APIs](https://docs ## Compatibility -The Nats package is tested with NATS 2.10.27. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+) and Elastic Agent 9.1+. Consumer metrics require NATS 2.9+. +The NATS package is tested with NATS 2.10.27 and requires Elastic Agent 9.1+. The `jetstream` dataset requires NATS with JetStream enabled (NATS 2.2+), and consumer metrics require NATS 2.9+. ## Logs From a6c857315a08d9bcc6b33ee9e69f1dbf39308c99 Mon Sep 17 00:00:00 2001 From: Nikolai Gut Date: Fri, 4 Sep 2026 11:53:22 +0200 Subject: [PATCH 9/9] fix(nats): bump version to 2.0.0, document breaking change, and clean up run.sh routes --- packages/nats/_dev/deploy/docker/run.sh | 2 +- packages/nats/changelog.yml | 7 +++++-- packages/nats/manifest.yml | 2 +- 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/packages/nats/_dev/deploy/docker/run.sh b/packages/nats/_dev/deploy/docker/run.sh index 5b74e22bc02..248262c1502 100755 --- a/packages/nats/_dev/deploy/docker/run.sh +++ b/packages/nats/_dev/deploy/docker/run.sh @@ -10,7 +10,7 @@ mkdir -p /var/log/nats # NATS 2.X if [ -x /opt/nats/nats-server ]; then if [ -z "${ROUTES}" ]; then - (/opt/nats/nats-server -DV -js --server_name nats --cluster_name nats-cluster -l /var/log/nats/nats.log --cluster nats://0.0.0.0:6222 --http_port 8222 --port 4222 --routes nats://nats-routes:6222) & + (/opt/nats/nats-server -DV -js --server_name nats --cluster_name nats-cluster -l /var/log/nats/nats.log --cluster nats://0.0.0.0:6222 --http_port 8222 --port 4222) & else (/opt/nats/nats-server -DV -js --server_name nats-routes --cluster_name nats-cluster -l /var/log/nats/nats.log --cluster nats://0.0.0.0:6222 --http_port 8222 --port 4222 --routes nats://nats:6222) & fi diff --git a/packages/nats/changelog.yml b/packages/nats/changelog.yml index 50618109139..ef833ad2d81 100644 --- a/packages/nats/changelog.yml +++ b/packages/nats/changelog.yml @@ -1,7 +1,10 @@ # newer versions go on top -- version: "1.13.0" +- version: "2.0.0" changes: - - description: Add JetStream data stream for monitoring JetStream stats, accounts, streams, and consumers. Requires Elastic Agent 9.1+. + - description: Raise the minimum required Kibana and Elastic Agent versions to 9.1.0, needed for JetStream support. Deployments on older stacks cannot upgrade to 2.x; the 1.x line is reserved for backports serving older stacks. + type: breaking-change + link: https://github.com/elastic/integrations/pull/20980 + - description: Add JetStream data stream for monitoring JetStream stats, accounts, streams, and consumers. type: enhancement link: https://github.com/elastic/integrations/pull/20980 - version: "1.12.0" diff --git a/packages/nats/manifest.yml b/packages/nats/manifest.yml index f4ed3dec717..83404060ec6 100644 --- a/packages/nats/manifest.yml +++ b/packages/nats/manifest.yml @@ -1,6 +1,6 @@ name: nats title: NATS -version: 1.13.0 +version: 2.0.0 description: Collect logs and metrics from NATS servers with Elastic Agent. type: integration icons: