Improve ACLK query processing (#19518)
* aclk_query_free handles everything that needs freed * Fix memory leak during shutdown (resolve coverity issues)
Stelios Fragkakis committed
Jan 29, 2025 at 10:37 UTC
68b24ae8a9b907005b726e21615012e662ea657f
4 files changed
+64
-33
src/aclk/aclk_query_queue.c
+24
@@ -11,12 +11,36 @@ aclk_query_t aclk_query_new(aclk_query_type_t type)
11
12
void aclk_query_free(aclk_query_t query)
13
{
14
+ struct ctxs_checkpoint *cmd;
15
switch (query->type) {
16
case HTTP_API_V2:
17
freez(query->data.http_api_v2.payload);
18
if (query->data.http_api_v2.query != query->dedup_id)
19
freez(query->data.http_api_v2.query);
20
break;
21
+ case ALERT_START_STREAMING:
22
+ freez(query->data.node_id);
23
+ break;
24
+ case ALERT_CHECKPOINT:
25
+ freez(query->data.node_id);
26
+ freez(query->claim_id);
27
+ break;
28
+ case CREATE_NODE_INSTANCE:
29
+ freez(query->data.node_id);
30
+ freez(query->machine_guid);
31
+ break;
32
+ case CTX_STOP_STREAMING:
33
+ cmd = query->data.payload;
34
+ freez(cmd->claim_id);
35
+ freez(cmd->node_id);
36
+ freez(cmd);
37
+ break;
38
+ case CTX_CHECKPOINT:
39
+ cmd = query->data.payload;
40
+ freez(cmd->claim_id);
41
+ freez(cmd->node_id);
42
+ freez(cmd);
43
+ break;
44
45
default:
46
break;
src/aclk/aclk_rx_msgs.c
+7
-13
@@ -256,12 +256,9 @@ int create_node_instance_result(const char *msg, size_t msg_len)
256
netdata_log_debug(D_ACLK, "CreateNodeInstanceResult: guid:%s nodeid:%s", res.machine_guid, res.node_id);
257
258
aclk_query_t query = aclk_query_new(CREATE_NODE_INSTANCE);
259
- query->data.node_id = strdupz(res.node_id);
260
- query->machine_guid = strdupz(res.machine_guid);
259
+ query->data.node_id = res.node_id; // Will be freed on query free
260
+ query->machine_guid = res.machine_guid; // Will be freed on query free
261
aclk_add_job(query);
262
-
263
- freez(res.node_id);
264
- freez(res.machine_guid);
262
return 0;
263
}
264
@@ -306,10 +303,9 @@ int start_alarm_streaming(const char *msg, size_t msg_len)
303
return 1;
304
}
305
aclk_query_t query = aclk_query_new(ALERT_START_STREAMING);
309
- query->data.node_id = strdupz(res.node_id);
306
+ query->data.node_id = res.node_id; // Will be freed on query free
307
query->version = res.version;
308
aclk_add_job(query);
312
- freez(res.node_id);
309
return 0;
310
}
311
@@ -323,12 +319,10 @@ int send_alarm_checkpoint(const char *msg, size_t msg_len)
319
return 1;
320
}
321
aclk_query_t query = aclk_query_new(ALERT_CHECKPOINT);
326
- query->data.node_id = strdupz(sac.node_id);
327
- query->claim_id = strdupz(sac.claim_id);
322
+ query->data.node_id = sac.node_id; // Will be freed on query free
323
+ query->claim_id = sac.claim_id;
324
query->version = sac.version;
325
aclk_add_job(query);
330
- freez(sac.node_id);
331
- freez(sac.claim_id);
326
return 0;
327
}
328
@@ -354,8 +348,8 @@ int send_alarm_snapshot(const char *msg, size_t msg_len)
348
return 1;
349
}
350
aclk_query_t query = aclk_query_new(ALERT_CHECKPOINT);
357
- query->data.node_id = strdupz(sas->node_id);
358
- query->claim_id = strdupz(sas->claim_id);
351
+ query->data.node_id = sas->node_id; // Will be freed on query free
352
+ query->claim_id = sas->claim_id; // Will be freed on query free
353
query->version = 0; // force snapshot
354
aclk_add_job(query);
355
destroy_send_alarm_snapshot(sas);
src/aclk/schema-wrappers/alarm_stream.cc
-2
@@ -176,8 +176,6 @@ struct send_alarm_snapshot *parse_send_alarm_snapshot(const char *data, size_t l
176
177
void destroy_send_alarm_snapshot(struct send_alarm_snapshot *ptr)
178
{
179
- freez(ptr->claim_id);
180
- freez(ptr->node_id);
179
freez(ptr->snapshot_uuid);
180
freez(ptr);
181
}
src/database/sqlite/sqlite_aclk.c
+33
-18
@@ -373,8 +373,6 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
373
return;
374
}
375
376
- struct ctxs_checkpoint *cmd;
377
-
376
bool ok_to_send = true;
377
378
switch (query->type) {
@@ -387,20 +385,12 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
385
break;
386
case CTX_CHECKPOINT:;
387
worker_is_busy(UV_EVENT_CTX_CHECKPOINT);
390
- cmd = query->data.payload;
391
- rrdcontext_hub_checkpoint_command(cmd);
392
- freez(cmd->claim_id);
393
- freez(cmd->node_id);
394
- freez(cmd);
388
+ rrdcontext_hub_checkpoint_command(query->data.payload);
389
ok_to_send = false;
390
break;
391
case CTX_STOP_STREAMING:
392
worker_is_busy(UV_EVENT_CTX_STOP_STREAMING);
399
- cmd = query->data.payload;
400
- rrdcontext_hub_stop_streaming_command(cmd);
401
- freez(cmd->claim_id);
402
- freez(cmd->node_id);
403
- freez(cmd);
393
+ rrdcontext_hub_stop_streaming_command(query->data.payload);
394
ok_to_send = false;
395
break;
396
case SEND_NODE_INSTANCES:
@@ -411,21 +401,16 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
401
case ALERT_START_STREAMING:
402
worker_is_busy(UV_EVENT_ALERT_START_STREAMING);
403
aclk_start_alert_streaming(query->data.node_id, query->version);
414
- freez(query->data.node_id);
404
ok_to_send = false;
405
break;
406
case ALERT_CHECKPOINT:
407
worker_is_busy(UV_EVENT_ALERT_CHECKPOINT);
408
aclk_alert_version_check(query->data.node_id, query->claim_id, query->version);
420
- freez(query->data.node_id);
421
- freez(query->claim_id);
409
ok_to_send = false;
410
break;
411
case CREATE_NODE_INSTANCE:
412
worker_is_busy(UV_EVENT_CREATE_NODE_INSTANCE);
413
create_node_instance_result_job(query->machine_guid, query->data.node_id);
427
- freez(query->data.node_id);
428
- freez(query->machine_guid);
414
ok_to_send = false;
415
break;
416
@@ -611,6 +596,19 @@ int schedule_query_in_worker(uv_loop_t *loop, struct aclk_sync_config_s *config,
596
return rc;
597
}
598
599
+static void free_query_list(Pvoid_t JudyL)
600
+{
601
+ bool first = true;
602
+ Pvoid_t *Pvalue;
603
+ Word_t Index = 0;
604
+ aclk_query_t query;
605
+ while ((Pvalue = JudyLFirstThenNext(JudyL, &Index, &first))) {
606
+ if (!*Pvalue)
607
+ continue;
608
+ query = *Pvalue;
609
+ aclk_query_free(query);
610
+ }
611
+}
612
613
static void aclk_synchronization(void *arg)
614
{
@@ -650,8 +648,11 @@ static void aclk_synchronization(void *arg)
648
649
struct worker_data *data;
650
aclk_query_t query;
651
+
652
+ // This holds queries that need to be executed one by one
653
struct judy_list_t *aclk_query_batch = NULL;
654
- struct judy_list_t *aclk_query_execute = callocz(1, sizeof(*aclk_query_execute));;
654
+ // This holds queries that can be dispatched in parallel in ACLK QUERY worker threads
655
+ struct judy_list_t *aclk_query_execute = callocz(1, sizeof(*aclk_query_execute));
656
size_t pending_queries = 0;
657
658
Pvoid_t *Pvalue;
@@ -880,6 +881,20 @@ static void aclk_synchronization(void *arg)
881
882
(void) uv_loop_close(loop);
883
884
+ // Free execute commands / queries
885
+ if (pending_queries) {
886
+ free_query_list(aclk_query_execute->JudyL);
887
+ (void)JudyLFreeArray(&aclk_query_execute->JudyL, PJE0);
888
+ }
889
+ freez(aclk_query_execute);
890
+
891
+ // Free batch commands
892
+ if (aclk_query_batch) {
893
+ free_query_list(aclk_query_batch->JudyL);
894
+ (void)JudyLFreeArray(&aclk_query_batch->JudyL, PJE0);
895
+ freez(aclk_query_batch);
896
+ }
897
+
898
aral_by_size_release(config->ar);
899
900
worker_unregister();