improvement(go.d/nats): add basic jetstream metrics (#19285)
Ilya Mashchenko committed
Dec 25, 2024 at 18:04 UTC
38d641a509794400fb5f5ea20aa6a88b2d092187
6 files changed
+294
-5
src/go/plugin/go.d/collector/nats/charts.go
+131
-4
@@ -25,6 +25,16 @@ const (
25
prioServerMemoryUsage
26
prioServerUptime
27
28
+ prioJetStreamStreams
29
+ prioJetStreamConsumers
30
+ prioJetStreamBytes
31
+ prioJetStreamMessages
32
+ prioJetStreamApiRequests
33
+ prioJetStreamApiErrors
34
+ prioJetStreamApiInflight
35
+ prioJetStreamMemoryUsed
36
+ prioJetStreamStorageUsed
37
+
38
prioAccountTraffic
39
prioAccountMessages
40
prioAccountConnections
@@ -48,7 +58,7 @@ const (
58
prioLeafRTT
59
)
60
51
-var serverCharts = func() module.Charts {
61
+func serverCharts() *module.Charts {
62
charts := module.Charts{
63
chartServerConnectionsCurrent.Copy(),
64
chartServerConnectionsRate.Copy(),
@@ -60,8 +70,9 @@ var serverCharts = func() module.Charts {
70
chartServerUptime.Copy(),
71
}
72
charts = append(charts, httpEndpointsCharts()...)
63
- return charts
64
-}()
73
+ charts = append(charts, *jetStreamCharts.Copy()...)
74
+ return charts.Copy()
75
+}
76
77
var (
78
chartServerTraffic = module.Chart{
@@ -191,6 +202,122 @@ var httpEndpointRequestsChartTmpl = module.Chart{
202
},
203
}
204
205
+var jetStreamCharts = module.Charts{
206
+ jetStreamStreams.Copy(),
207
+ jetStreamStreamsStorageBytes.Copy(),
208
+ jetStreamStreamsStorageMessages.Copy(),
209
+ jetStreamConsumers.Copy(),
210
+ jetStreamApiRequests.Copy(),
211
+ jetStreamApiInflightRequests.Copy(),
212
+ jetStreamApiErrors.Copy(),
213
+ jetStreamMemoryUsed.Copy(),
214
+ jetStreamStorageUsed.Copy(),
215
+}
216
+
217
+var (
218
+ jetStreamStreams = module.Chart{
219
+ ID: "jetstream_streams",
220
+ Title: "JetStream Streams",
221
+ Units: "streams",
222
+ Fam: "jstream streams",
223
+ Ctx: "nats.jetstream_streams",
224
+ Priority: prioJetStreamStreams,
225
+ Dims: module.Dims{
226
+ {ID: "jsz_streams", Name: "active"},
227
+ },
228
+ }
229
+ jetStreamStreamsStorageBytes = module.Chart{
230
+ ID: "jetstream_streams_storage_bytes",
231
+ Title: "JetStream Bytes",
232
+ Units: "bytes",
233
+ Fam: "jstream streams",
234
+ Ctx: "nats.jetstream_streams_storage_bytes",
235
+ Priority: prioJetStreamBytes,
236
+ Type: module.Area,
237
+ Dims: module.Dims{
238
+ {ID: "jsz_bytes", Name: "used"},
239
+ },
240
+ }
241
+ jetStreamStreamsStorageMessages = module.Chart{
242
+ ID: "jetstream_streams_storage_messages",
243
+ Title: "JetStream Messages",
244
+ Units: "messages",
245
+ Fam: "jstream streams",
246
+ Ctx: "nats.jetstream_streams_storage_messages",
247
+ Priority: prioJetStreamMessages,
248
+ Dims: module.Dims{
249
+ {ID: "jsz_messages", Name: "stored"},
250
+ },
251
+ }
252
+ jetStreamConsumers = module.Chart{
253
+ ID: "jetstream_consumers",
254
+ Title: "JetStream Consumers",
255
+ Units: "consumers",
256
+ Fam: "jstream consumers",
257
+ Ctx: "nats.jetstream_consumers",
258
+ Priority: prioJetStreamConsumers,
259
+ Dims: module.Dims{
260
+ {ID: "jsz_consumers", Name: "active"},
261
+ },
262
+ }
263
+ jetStreamApiRequests = module.Chart{
264
+ ID: "jetstream_api_requests",
265
+ Title: "JetStream API Requests",
266
+ Units: "requests/s",
267
+ Fam: "jstream api",
268
+ Ctx: "nats.jetstream_api_requests",
269
+ Priority: prioJetStreamApiRequests,
270
+ Dims: module.Dims{
271
+ {ID: "jsz_api_total", Name: "requests", Algo: module.Incremental},
272
+ },
273
+ }
274
+ jetStreamApiErrors = module.Chart{
275
+ ID: "jetstream_api_errors",
276
+ Title: "JetStream API Errors",
277
+ Units: "errors/s",
278
+ Fam: "jstream api",
279
+ Ctx: "nats.jetstream_api_errors",
280
+ Priority: prioJetStreamApiErrors,
281
+ Dims: module.Dims{
282
+ {ID: "jsz_api_errors", Name: "errors", Algo: module.Incremental},
283
+ },
284
+ }
285
+ jetStreamApiInflightRequests = module.Chart{
286
+ ID: "jetstream_api_inflight",
287
+ Title: "JetStream API Inflight",
288
+ Units: "requests",
289
+ Fam: "jstream api",
290
+ Ctx: "nats.jetstream_api_inflight",
291
+ Priority: prioJetStreamApiInflight,
292
+ Dims: module.Dims{
293
+ {ID: "jsz_api_inflight", Name: "inflight"},
294
+ },
295
+ }
296
+ jetStreamMemoryUsed = module.Chart{
297
+ ID: "jetstream_memory_used",
298
+ Title: "JetStream Used Memory",
299
+ Units: "bytes",
300
+ Fam: "jstream rusage",
301
+ Ctx: "nats.jetstream_memory_used",
302
+ Priority: prioJetStreamMemoryUsed,
303
+ Type: module.Area,
304
+ Dims: module.Dims{
305
+ {ID: "jsz_memory_used", Name: "used"},
306
+ },
307
+ }
308
+ jetStreamStorageUsed = module.Chart{
309
+ ID: "jetstream_storage_used",
310
+ Title: "JetStream Used Storage",
311
+ Units: "bytes",
312
+ Fam: "jstream rusage",
313
+ Ctx: "nats.jetstream_storage_used",
314
+ Priority: prioJetStreamStorageUsed,
315
+ Dims: module.Dims{
316
+ {ID: "jsz_store_used", Name: "used"},
317
+ },
318
+ }
319
+)
320
+
321
var accountChartsTmpl = module.Charts{
322
accountTrafficTmpl.Copy(),
323
accountMessagesTmpl.Copy(),
@@ -522,7 +649,7 @@ func (c *Collector) updateCharts() {
649
}
650
651
func (c *Collector) addServerCharts() {
525
- charts := serverCharts.Copy()
652
+ charts := serverCharts()
653
654
for _, chart := range *charts {
655
chart.Labels = []module.Label{
src/go/plugin/go.d/collector/nats/collect.go
+27
@@ -45,6 +45,9 @@ func (c *Collector) collect() (map[string]int64, error) {
45
if err := c.collectLeafz(mx); err != nil {
46
return mx, err
47
}
48
+ if err := c.collectJsz(mx); err != nil {
49
+ return mx, err
50
+ }
51
52
c.updateCharts()
53
@@ -260,6 +263,30 @@ func (c *Collector) collectLeafz(mx map[string]int64) error {
263
return nil
264
}
265
266
+func (c *Collector) collectJsz(mx map[string]int64) error {
267
+ req, err := web.NewHTTPRequestWithPath(c.RequestConfig, urlPathJsz)
268
+ if err != nil {
269
+ return err
270
+ }
271
+
272
+ var resp jszResponse
273
+ if err := web.DoHTTP(c.httpClient).RequestJSON(req, &resp); err != nil {
274
+ return err
275
+ }
276
+
277
+ mx["jsz_streams"] = int64(resp.Streams)
278
+ mx["jsz_consumers"] = int64(resp.Consumers)
279
+ mx["jsz_bytes"] = int64(resp.Bytes)
280
+ mx["jsz_messages"] = int64(resp.Messages)
281
+ mx["jsz_memory_used"] = int64(resp.Memory)
282
+ mx["jsz_store_used"] = int64(resp.Store)
283
+ mx["jsz_api_total"] = int64(resp.Api.Total)
284
+ mx["jsz_api_errors"] = int64(resp.Api.Errors)
285
+ mx["jsz_api_inflight"] = int64(resp.Api.Inflight)
286
+
287
+ return nil
288
+}
289
+
290
func parseUptime(uptime string) (time.Duration, error) {
291
// https://github.com/nats-io/nats-server/blob/v2.10.24/server/monitor.go#L1354
292
src/go/plugin/go.d/collector/nats/collector_test.go
+14
-1
@@ -26,6 +26,7 @@ var (
26
dataVer210Routez, _ = os.ReadFile("testdata/v2.10.24/routez.json")
27
dataVer210Gatewayz, _ = os.ReadFile("testdata/v2.10.24/gatewayz.json")
28
dataVer210Leafz, _ = os.ReadFile("testdata/v2.10.24/leafz.json")
29
+ dataVer210Jsz, _ = os.ReadFile("testdata/v2.10.24/jsz.json")
30
)
31
32
func Test_testDataIsValid(t *testing.T) {
@@ -38,6 +39,7 @@ func Test_testDataIsValid(t *testing.T) {
39
"dataVer210Routez": dataVer210Routez,
40
"dataVer210Gatewayz": dataVer210Gatewayz,
41
"dataVer210Leafz": dataVer210Leafz,
42
+ "dataVer210Jsz": dataVer210Jsz,
43
} {
44
require.NotNil(t, data, name)
45
}
@@ -133,7 +135,7 @@ func TestCollector_Collect(t *testing.T) {
135
}{
136
"success on valid response": {
137
prepare: caseOk,
136
- wantNumOfCharts: len(serverCharts) +
138
+ wantNumOfCharts: len(*serverCharts()) +
139
len(accountChartsTmpl)*3 +
140
len(routeChartsTmpl)*1 +
141
len(gatewayConnChartsTmpl)*5 +
@@ -196,6 +198,15 @@ func TestCollector_Collect(t *testing.T) {
198
"gatewayz_outbound_gw_region3_cid_5_out_bytes": 0,
199
"gatewayz_outbound_gw_region3_cid_5_out_msgs": 0,
200
"gatewayz_outbound_gw_region3_cid_5_uptime": 6,
201
+ "jsz_api_errors": 588,
202
+ "jsz_api_inflight": 0,
203
+ "jsz_api_total": 936916,
204
+ "jsz_bytes": 114419224,
205
+ "jsz_consumers": 9,
206
+ "jsz_memory_used": 128,
207
+ "jsz_messages": 5670,
208
+ "jsz_store_used": 114419224,
209
+ "jsz_streams": 198,
210
"leafz_leaf__$G_127.0.0.1_6223_in_bytes": 0,
211
"leafz_leaf__$G_127.0.0.1_6223_in_msgs": 0,
212
"leafz_leaf__$G_127.0.0.1_6223_num_subs": 1,
@@ -297,6 +308,8 @@ func caseOk(t *testing.T) (*Collector, func()) {
308
_, _ = w.Write(dataVer210Gatewayz)
309
case urlPathLeafz:
310
_, _ = w.Write(dataVer210Leafz)
311
+ case urlPathJsz:
312
+ _, _ = w.Write(dataVer210Jsz)
313
default:
314
w.WriteHeader(http.StatusNotFound)
315
}
src/go/plugin/go.d/collector/nats/metadata.yaml
+54
@@ -240,6 +240,60 @@ modules:
240
chart_type: line
241
dimensions:
242
- name: uptime
243
+ - name: nats.jetstream_streams
244
+ description: JetStream Streams
245
+ unit: streams
246
+ chart_type: line
247
+ dimensions:
248
+ - name: active
249
+ - name: nats.jetstream_streams_storage_bytes
250
+ description: JetStream Bytes
251
+ unit: bytes
252
+ chart_type: area
253
+ dimensions:
254
+ - name: used
255
+ - name: nats.jetstream_streams_storage_messages
256
+ description: JetStream Messages
257
+ unit: messaged
258
+ chart_type: line
259
+ dimensions:
260
+ - name: stored
261
+ - name: nats.jetstream_consumers
262
+ description: JetStream Consumers
263
+ unit: consumers
264
+ chart_type: line
265
+ dimensions:
266
+ - name: active
267
+ - name: nats.jetstream_api_requests
268
+ description: JetStream API Requests
269
+ unit: requests/s
270
+ chart_type: line
271
+ dimensions:
272
+ - name: requests
273
+ - name: nats.jetstream_api_errors
274
+ description: JetStream API Errors
275
+ unit: errors/s
276
+ chart_type: line
277
+ dimensions:
278
+ - name: errors
279
+ - name: nats.jetstream_api_inflight
280
+ description: JetStream API Inflight
281
+ unit: requests
282
+ chart_type: line
283
+ dimensions:
284
+ - name: inflight
285
+ - name: nats.jetstream_memory_used
286
+ description: JetStream Used Memory
287
+ unit: bytes
288
+ chart_type: area
289
+ dimensions:
290
+ - name: used
291
+ - name: nats.jetstream_storage_used
292
+ description: JetStream Used Storage
293
+ unit: bytes
294
+ chart_type: line
295
+ dimensions:
296
+ - name: used
297
- name: http endpoint
298
description: These metrics refer to HTTP endpoints.
299
labels:
src/go/plugin/go.d/collector/nats/restapi.go
+43
@@ -3,6 +3,8 @@
3
package nats
4
5
import (
6
+ "time"
7
+
8
"github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/web"
9
)
10
@@ -21,6 +23,8 @@ const (
23
urlPathGatewayz = "/gatewayz"
24
// https://docs.nats.io/running-a-nats-service/nats_admin/monitoring#leaf-node-information
25
urlPathLeafz = "/leafz"
26
+ // https://docs.nats.io/running-a-nats-service/nats_admin/monitoring#jetstream-information
27
+ urlPathJsz = "/jsz"
28
)
29
30
var (
@@ -159,3 +163,42 @@ type (
163
NumSubs uint32 `json:"subscriptions"`
164
}
165
)
166
+
167
+// https://github.com/nats-io/nats-server/blob/v2.10.24/server/monitor.go#L2801
168
+type (
169
+ jszResponse struct {
170
+ Disabled bool `json:"disabled"`
171
+ Streams int `json:"streams"`
172
+ Consumers int `json:"consumers"`
173
+ Messages uint64 `json:"messages"`
174
+ Bytes uint64 `json:"bytes"`
175
+ Memory uint64 `json:"memory"`
176
+ Store uint64 `json:"storage"`
177
+ ReservedMemory uint64 `json:"reserved_memory"`
178
+ ReservedStore uint64 `json:"reserved_storage"`
179
+ Accounts int `json:"accounts"`
180
+ HAAssets int `json:"ha_assets"`
181
+ Api struct {
182
+ Total uint64 `json:"total"`
183
+ Errors uint64 `json:"errors"`
184
+ Inflight uint64 `json:"inflight"`
185
+ } `json:"api"`
186
+ Meta *jszMetaClusterInfo `json:"meta_cluster"`
187
+ }
188
+ jszMetaClusterInfo struct {
189
+ Name string `json:"name"`
190
+ Leader string `json:"leader"`
191
+ Peer string `json:"peer"`
192
+ Replicas []*jszPeerInfo `json:"replicas"`
193
+ Size int `json:"cluster_size"`
194
+ Pending int `json:"pending"`
195
+ }
196
+ jszPeerInfo struct {
197
+ Name string `json:"name"`
198
+ Current bool `json:"current"`
199
+ Offline bool `json:"offline"`
200
+ Active time.Duration `json:"active"`
201
+ Lag uint64 `json:"lag"`
202
+ Peer string `json:"peer"`
203
+ }
204
+)
src/go/plugin/go.d/collector/nats/testdata/v2.10.24/jsz.json
new
+25
@@ -0,0 +1,25 @@
1
+{
2
+ "server_id": "NDR5FR76SWSTAP5LSKUNYG7ADXPFLXXWNC7ALU3WGX6WLFMMCIDQAD4J",
3
+ "now": "2024-12-25T14:20:19.733407331Z",
4
+ "config": {
5
+ "max_memory": 10737418240,
6
+ "max_storage": 440234147840,
7
+ "store_dir": "/var/jetstream/jetstream",
8
+ "sync_interval": 120000000000,
9
+ "compress_ok": true
10
+ },
11
+ "memory": 128,
12
+ "storage": 114419224,
13
+ "reserved_memory": 0,
14
+ "reserved_storage": 1620615736,
15
+ "accounts": 1,
16
+ "ha_assets": 0,
17
+ "api": {
18
+ "total": 936916,
19
+ "errors": 588
20
+ },
21
+ "streams": 198,
22
+ "consumers": 9,
23
+ "messages": 5670,
24
+ "bytes": 114419224
25
+}