@cryptotaxi247 / netdata-1 / commits / 7d0ca0b83

Adds some introspection into the MQTT_WSS (#14039)

Timotej S committed Jan 16, 2023 at 15:32 UTC 7d0ca0b83e96fdebed5e66b73362ece4fa0ad62e
3 files changed +73 -2
CMakeLists.txt
+2 -1
@@ -1282,7 +1282,8 @@ ADD_LIBRARY(mqttwebsockets STATIC
1282
1283 target_compile_options(mqttwebsockets PUBLIC
1284 -DMQTT_WSS_CUSTOM_ALLOC
1285 - -DRBUF_CUSTOM_MALLOC)
1285 + -DRBUF_CUSTOM_MALLOC
1286 + -DMQTT_WSS_CPUSTATS)
1287
1288 target_include_directories(mqttwebsockets PUBLIC
1289 ${CMAKE_SOURCE_DIR}/aclk/helpers)
Makefile.am
+1 -1
@@ -742,7 +742,7 @@ libmqttwebsockets_a_SOURCES = \
742 mqtt_websockets/c_rhash/include/c_rhash.h \
743 mqtt_websockets/c_rhash/src/c_rhash_internal.h
744
745 -libmqttwebsockets_a_CFLAGS = $(CFLAGS) -DMQTT_WSS_CUSTOM_ALLOC -DRBUF_CUSTOM_MALLOC -I$(srcdir)/aclk/helpers -I$(srcdir)/mqtt_websockets/c_rhash/include
745 +libmqttwebsockets_a_CFLAGS = $(CFLAGS) -DMQTT_WSS_CUSTOM_ALLOC -DRBUF_CUSTOM_MALLOC -DMQTT_WSS_CPUSTATS -I$(srcdir)/aclk/helpers -I$(srcdir)/mqtt_websockets/c_rhash/include
746
747 if MQTT_WSS_DEBUG
748 libmqttwebsockets_a_CFLAGS += -DMQTT_WSS_DEBUG
aclk/aclk_stats.c
+70
@@ -1,5 +1,7 @@
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 +#define MQTT_WSS_CPUSTATS
4 +
5 #include "aclk_stats.h"
6
7 #include "aclk_query.h"
@@ -259,6 +261,23 @@ static void aclk_stats_mqtt_wss(struct mqtt_wss_stats *stats)
261 static uint64_t sent = 0;
262 static uint64_t recvd = 0;
263
264 + static RRDSET *st_txbuf_perc = NULL;
265 + static RRDDIM *rd_txbuf_perc = NULL;
266 +
267 + static RRDSET *st_txbuf = NULL;
268 + static RRDDIM *rd_tx_buffer_usable = NULL;
269 + static RRDDIM *rd_tx_buffer_reclaimable = NULL;
270 + static RRDDIM *rd_tx_buffer_used = NULL;
271 + static RRDDIM *rd_tx_buffer_free = NULL;
272 + static RRDDIM *rd_tx_buffer_size = NULL;
273 +
274 + static RRDSET *st_timing = NULL;
275 + static RRDDIM *rd_keepalive = NULL;
276 + static RRDDIM *rd_read_socket = NULL;
277 + static RRDDIM *rd_write_socket = NULL;
278 + static RRDDIM *rd_process_websocket = NULL;
279 + static RRDDIM *rd_process_mqtt = NULL;
280 +
281 sent += stats->bytes_tx;
282 recvd += stats->bytes_rx;
283
@@ -271,10 +290,61 @@ static void aclk_stats_mqtt_wss(struct mqtt_wss_stats *stats)
290 rd_recvd = rrddim_add(st, "received", NULL, 1, 1, RRD_ALGORITHM_INCREMENTAL);
291 }
292
293 + if (unlikely(!st_txbuf_perc)) {
294 + st_txbuf_perc = rrdset_create_localhost(
295 + "netdata", "aclk_mqtt_tx_perc", NULL, "aclk", NULL, "Actively used percentage of MQTT Tx Buffer,", "%",
296 + "netdata", "stats", 200012, localhost->rrd_update_every, RRDSET_TYPE_LINE);
297 +
298 + rd_txbuf_perc = rrddim_add(st_txbuf_perc, "used", NULL, 1, 100, RRD_ALGORITHM_ABSOLUTE);
299 + }
300 +
301 + if (unlikely(!st_txbuf)) {
302 + st_txbuf = rrdset_create_localhost(
303 + "netdata", "aclk_mqtt_tx_queue", NULL, "aclk", NULL, "State of transmit MQTT queue.", "B",
304 + "netdata", "stats", 200013, localhost->rrd_update_every, RRDSET_TYPE_LINE);
305 +
306 + rd_tx_buffer_usable = rrddim_add(st_txbuf, "usable", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
307 + rd_tx_buffer_reclaimable = rrddim_add(st_txbuf, "reclaimable", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
308 + rd_tx_buffer_used = rrddim_add(st_txbuf, "used", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
309 + rd_tx_buffer_free = rrddim_add(st_txbuf, "free", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
310 + rd_tx_buffer_size = rrddim_add(st_txbuf, "size", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
311 + }
312 +
313 + if (unlikely(!st_timing)) {
314 + st_timing = rrdset_create_localhost(
315 + "netdata", "aclk_mqtt_wss_time", NULL, "aclk", NULL, "Time spent handling MQTT, WSS, SSL and network communication.", "us",
316 + "netdata", "stats", 200014, localhost->rrd_update_every, RRDSET_TYPE_STACKED);
317 +
318 + rd_keepalive = rrddim_add(st_timing, "keep-alive", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
319 + rd_read_socket = rrddim_add(st_timing, "socket_read_ssl", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
320 + rd_write_socket = rrddim_add(st_timing, "socket_write_ssl", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
321 + rd_process_websocket = rrddim_add(st_timing, "process_websocket", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
322 + rd_process_mqtt = rrddim_add(st_timing, "process_mqtt", NULL, 1, 1, RRD_ALGORITHM_ABSOLUTE);
323 + }
324 +
325 rrddim_set_by_pointer(st, rd_sent, sent);
326 rrddim_set_by_pointer(st, rd_recvd, recvd);
327
328 + float usage = ((float)stats->mqtt.tx_buffer_free + stats->mqtt.tx_buffer_reclaimable) / stats->mqtt.tx_buffer_size;
329 + usage = (1 - usage) * 10000;
330 + rrddim_set_by_pointer(st_txbuf_perc, rd_txbuf_perc, usage);
331 +
332 + rrddim_set_by_pointer(st_txbuf, rd_tx_buffer_usable, stats->mqtt.tx_buffer_reclaimable + stats->mqtt.tx_buffer_free);
333 + rrddim_set_by_pointer(st_txbuf, rd_tx_buffer_reclaimable, stats->mqtt.tx_buffer_reclaimable);
334 + rrddim_set_by_pointer(st_txbuf, rd_tx_buffer_used, stats->mqtt.tx_buffer_used);
335 + rrddim_set_by_pointer(st_txbuf, rd_tx_buffer_free, stats->mqtt.tx_buffer_free);
336 + rrddim_set_by_pointer(st_txbuf, rd_tx_buffer_size, stats->mqtt.tx_buffer_size);
337 +
338 + rrddim_set_by_pointer(st_timing, rd_keepalive, stats->time_keepalive);
339 + rrddim_set_by_pointer(st_timing, rd_read_socket, stats->time_read_socket);
340 + rrddim_set_by_pointer(st_timing, rd_write_socket, stats->time_write_socket);
341 + rrddim_set_by_pointer(st_timing, rd_process_websocket, stats->time_process_websocket);
342 + rrddim_set_by_pointer(st_timing, rd_process_mqtt, stats->time_process_mqtt);
343 +
344 rrdset_done(st);
345 + rrdset_done(st_txbuf_perc);
346 + rrdset_done(st_txbuf);
347 + rrdset_done(st_timing);
348 }
349
350 void aclk_stats_thread_prepare(int query_thread_count, unsigned int proto_hdl_cnt)