@cryptotaxi247 / netdata / commits / c51b2b016

Monitor memory reclamation and buffer compact (#20810)

Stelios Fragkakis committed Aug 12, 2025 at 18:31 UTC c51b2b0169949ab1bbd8dcfbf350583c235a565d
3 files changed +19 -5
src/aclk/aclk.c
+2
@@ -840,6 +840,8 @@ void aclk_main(void *ptr)
840 worker_register_job_name(WORKER_ACLK_SEND_FRAGMENT, "send fragment");
841 worker_register_job_name(WORKER_ACLK_MSG_CALLBACK, "msg callback");
842 worker_register_job_name(WORKER_ACLK_WAITING_TO_CONNECT, "conn wait");
843 + worker_register_job_name(WORKER_ACLK_RECLAIM_MEMORY, "reclaim");
844 + worker_register_job_name(WORKER_ACLK_BUFFER_COMPACT, "compact");
845
846 ACLK_PROXY_TYPE proxy_type;
847 aclk_get_proxy(&proxy_type, false);
src/aclk/mqtt_websockets/aclk_mqtt_workers.h
+2
@@ -38,5 +38,7 @@
38 #define WORKER_ACLK_SEND_FRAGMENT 32
39 #define WORKER_ACLK_MSG_CALLBACK 33
40 #define WORKER_ACLK_WAITING_TO_CONNECT 34
41 +#define WORKER_ACLK_RECLAIM_MEMORY 35
42 +#define WORKER_ACLK_BUFFER_COMPACT 36
43
44 #endif //NETDATA_ACLK_MQTT_WORKERS_H
src/aclk/mqtt_websockets/mqtt_ng.c
+15 -5
@@ -486,7 +486,7 @@ static void buffer_rebuild(struct header_buffer *buf)
486 } while(frag);
487 }
488
489 -static void buffer_garbage_collect(struct header_buffer *buf)
489 +static void buffer_garbage_collect(struct header_buffer *buf, bool main_thread)
490 {
491 struct buffer_fragment *frag = BUFFER_FIRST_FRAG(buf);
492 while (frag) {
@@ -515,17 +515,27 @@ static void buffer_garbage_collect(struct header_buffer *buf)
515 #endif
516
517 memmove(buf->data, frag, buf->tail - (unsigned char *) frag);
518 + if (main_thread)
519 + worker_is_busy(WORKER_ACLK_BUFFER_COMPACT);
520 +
521 buffer_rebuild(buf);
522 +
523 + if (main_thread)
524 + worker_is_idle();
525 }
526
521 -static void transaction_buffer_garbage_collect(struct transaction_buffer *buf)
527 +static void transaction_buffer_garbage_collect(struct transaction_buffer *buf, bool main_thread)
528 {
529 + if (main_thread)
530 + worker_is_busy(WORKER_ACLK_RECLAIM_MEMORY);
531 // Invalidate the cached sending fragment
532 // as we will move data around
533 if (buf->sending_frag != &ping_frag)
534 buf->sending_frag = NULL;
535
528 - buffer_garbage_collect(&buf->hdr_buffer);
536 + buffer_garbage_collect(&buf->hdr_buffer, main_thread);
537 + if (main_thread)
538 + worker_is_idle();
539 }
540
541 static int transaction_buffer_grow(struct transaction_buffer *buf, float rate, size_t max)
@@ -838,7 +848,7 @@ static void add_packet_to_timeout_monitor_list(struct mqtt_ng_client *client, ui
848 int _rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
849 if (_rc == MQTT_NG_MSGGEN_BUFFER_OOM) { \
850 LOCK_HDR_BUFFER(&client->main_buffer); \
841 - transaction_buffer_garbage_collect((&client->main_buffer)); \
851 + transaction_buffer_garbage_collect(&client->main_buffer, false); \
852 UNLOCK_HDR_BUFFER(&client->main_buffer); \
853 _rc = generator_function(&client->main_buffer, ##__VA_ARGS__); \
854 if (_rc == MQTT_NG_MSGGEN_BUFFER_OOM && client->max_mem_bytes) { \
@@ -1196,7 +1206,7 @@ static int mark_packet_acked(struct mqtt_ng_client *client, uint16_t packet_id)
1206
1207 size_t used = BUFFER_BYTES_USED(&client->main_buffer.hdr_buffer);
1208 if (reclaimable >= (used / 4))
1199 - transaction_buffer_garbage_collect(&client->main_buffer);
1209 + transaction_buffer_garbage_collect(&client->main_buffer, true);
1210
1211 UNLOCK_HDR_BUFFER(&client->main_buffer);
1212 remove_packet_from_timeout_monitor_list_unsafe(client, packet_id);