Use aral in ACLK (#19459)
Code cleanup Use ARAL aclk_run_query always runs from worker
Stelios Fragkakis committed
Jan 22, 2025 at 23:50 UTC
1dc75d58d0a4e7322db5fe311baeb53914c86912
1 file changed
+31
-39
src/database/sqlite/sqlite_aclk.c
+31
-39
@@ -74,30 +74,33 @@ struct aclk_sync_config_s {
74
SPINLOCK cmd_queue_lock;
75
uint32_t aclk_jobs_pending;
76
struct aclk_database_cmd *cmd_base;
77
+ ARAL *ar;
78
} aclk_sync_config = { 0 };
79
80
static struct aclk_database_cmd aclk_database_deq_cmd(void)
81
{
82
struct aclk_database_cmd ret = { 0 };
83
+ struct aclk_database_cmd *to_free = NULL;
84
85
spinlock_lock(&aclk_sync_config.cmd_queue_lock);
86
if(aclk_sync_config.cmd_base) {
87
struct aclk_database_cmd *t = aclk_sync_config.cmd_base;
88
DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(aclk_sync_config.cmd_base, t, prev, next);
89
ret = *t;
88
- freez(t);
90
+ to_free = t;
91
}
92
else {
93
ret.opcode = ACLK_DATABASE_NOOP;
94
}
95
spinlock_unlock(&aclk_sync_config.cmd_queue_lock);
96
+ aral_freez(aclk_sync_config.ar, to_free);
97
98
return ret;
99
}
100
101
static void aclk_database_enq_cmd(struct aclk_database_cmd *cmd)
102
{
100
- struct aclk_database_cmd *t = mallocz(sizeof(*t));
103
+ struct aclk_database_cmd *t = aral_mallocz(aclk_sync_config.ar);
104
*t = *cmd;
105
t->prev = t->next = NULL;
106
@@ -367,7 +370,7 @@ static void after_aclk_run_query_job(uv_work_t *req, int status __maybe_unused)
370
freez(payload);
371
}
372
370
-static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query, bool is_worker)
373
+static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query)
374
{
375
if (query->type == UNKNOWN || query->type >= ACLK_QUERY_TYPE_COUNT) {
376
error_report("Unknown query in query queue. %u", query->type);
@@ -382,14 +385,12 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
385
386
// Incoming : cloud -> agent
387
case HTTP_API_V2:
385
- if (is_worker)
386
- worker_is_busy(UV_EVENT_ACLK_QUERY_EXECUTE);
388
+ worker_is_busy(UV_EVENT_ACLK_QUERY_EXECUTE);
389
http_api_v2(config->client, query);
390
ok_to_send = false;
391
break;
392
case CTX_CHECKPOINT:;
391
- if (is_worker)
392
- worker_is_busy(UV_EVENT_CTX_CHECKPOINT);
393
+ worker_is_busy(UV_EVENT_CTX_CHECKPOINT);
394
cmd = query->data.payload;
395
rrdcontext_hub_checkpoint_command(cmd);
396
freez(cmd->claim_id);
@@ -398,8 +399,7 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
399
ok_to_send = false;
400
break;
401
case CTX_STOP_STREAMING:
401
- if (is_worker)
402
- worker_is_busy(UV_EVENT_CTX_STOP_STREAMING);
402
+ worker_is_busy(UV_EVENT_CTX_STOP_STREAMING);
403
cmd = query->data.payload;
404
rrdcontext_hub_stop_streaming_command(cmd);
405
freez(cmd->claim_id);
@@ -408,29 +408,25 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
408
ok_to_send = false;
409
break;
410
case SEND_NODE_INSTANCES:
411
- if (is_worker)
412
- worker_is_busy(UV_EVENT_SEND_NODE_INSTANCES);
411
+ worker_is_busy(UV_EVENT_SEND_NODE_INSTANCES);
412
aclk_send_node_instances();
413
ok_to_send = false;
414
break;
415
case ALERT_START_STREAMING:
417
- if (is_worker)
418
- worker_is_busy(UV_EVENT_ALERT_START_STREAMING);
416
+ worker_is_busy(UV_EVENT_ALERT_START_STREAMING);
417
aclk_start_alert_streaming(query->data.node_id, query->version);
418
freez(query->data.node_id);
419
ok_to_send = false;
420
break;
421
case ALERT_CHECKPOINT:
424
- if (is_worker)
425
- worker_is_busy(UV_EVENT_ALERT_CHECKPOINT);
422
+ worker_is_busy(UV_EVENT_ALERT_CHECKPOINT);
423
aclk_alert_version_check(query->data.node_id, query->claim_id, query->version);
424
freez(query->data.node_id);
425
freez(query->claim_id);
426
ok_to_send = false;
427
break;
428
case CREATE_NODE_INSTANCE:
432
- if (is_worker)
433
- worker_is_busy(UV_EVENT_CREATE_NODE_INSTANCE);
429
+ worker_is_busy(UV_EVENT_CREATE_NODE_INSTANCE);
430
create_node_instance_result_job(query->machine_guid, query->data.node_id);
431
freez(query->data.node_id);
432
freez(query->machine_guid);
@@ -439,36 +435,28 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
435
436
// Outgoing: agent -> cloud
437
case ALARM_PROVIDE_CFG:
442
- if (is_worker)
443
- worker_is_busy(UV_EVENT_ALARM_PROVIDE_CFG);
438
+ worker_is_busy(UV_EVENT_ALARM_PROVIDE_CFG);
439
break;
440
case ALARM_SNAPSHOT:
446
- if (is_worker)
447
- worker_is_busy(UV_EVENT_ALARM_SNAPSHOT);
441
+ worker_is_busy(UV_EVENT_ALARM_SNAPSHOT);
442
break;
443
case REGISTER_NODE:
450
- if (is_worker)
451
- worker_is_busy(UV_EVENT_REGISTER_NODE);
444
+ worker_is_busy(UV_EVENT_REGISTER_NODE);
445
break;
446
case UPDATE_NODE_COLLECTORS:
454
- if (is_worker)
455
- worker_is_busy(UV_EVENT_UPDATE_NODE_COLLECTORS);
447
+ worker_is_busy(UV_EVENT_UPDATE_NODE_COLLECTORS);
448
break;
449
case UPDATE_NODE_INFO:
458
- if (is_worker)
459
- worker_is_busy(UV_EVENT_UPDATE_NODE_INFO);
450
+ worker_is_busy(UV_EVENT_UPDATE_NODE_INFO);
451
break;
452
case CTX_SEND_SNAPSHOT:
462
- if (is_worker)
463
- worker_is_busy(UV_EVENT_CTX_SEND_SNAPSHOT);
453
+ worker_is_busy(UV_EVENT_CTX_SEND_SNAPSHOT);
454
break;
455
case CTX_SEND_SNAPSHOT_UPD:
466
- if (is_worker)
467
- worker_is_busy(UV_EVENT_CTX_SEND_SNAPSHOT_UPD);
456
+ worker_is_busy(UV_EVENT_CTX_SEND_SNAPSHOT_UPD);
457
break;
458
case NODE_STATE_UPDATE:
470
- if (is_worker)
471
- worker_is_busy(UV_EVENT_NODE_STATE_UPDATE);
459
+ worker_is_busy(UV_EVENT_NODE_STATE_UPDATE);
460
break;
461
default:
462
nd_log_daemon(NDLP_ERR, "Unknown msg type %u; ignoring", query->type);
@@ -490,7 +478,7 @@ static void aclk_run_query_job(uv_work_t *req)
478
struct aclk_sync_config_s *config = payload->config;
479
aclk_query_t query = (aclk_query_t) payload->data;
480
493
- aclk_run_query(config, query, true);
481
+ aclk_run_query(config, query);
482
worker_is_idle();
483
}
484
@@ -518,13 +506,13 @@ static void aclk_execute_batch(uv_work_t *req)
506
size_t entries = aclk_query_batch->count;
507
Word_t Index = 0;
508
bool first = true;
521
- Pvoid_t *PValue;
522
- while ((PValue = JudyLFirstThenNext(aclk_query_batch->JudyL, &Index, &first))) {
523
- if (!*PValue)
509
+ Pvoid_t *Pvalue;
510
+ while ((Pvalue = JudyLFirstThenNext(aclk_query_batch->JudyL, &Index, &first))) {
511
+ if (!*Pvalue)
512
continue;
513
526
- aclk_query_t query = *PValue;
527
- aclk_run_query(config, query, true);
514
+ aclk_query_t query = *Pvalue;
515
+ aclk_run_query(config, query);
516
}
517
518
(void) JudyLFreeArray(&aclk_query_batch->JudyL, PJE0);
@@ -613,6 +601,8 @@ static void aclk_synchronization(void *arg)
601
{
602
struct aclk_sync_config_s *config = arg;
603
uv_thread_set_name_np("ACLKSYNC");
604
+ config->ar = aral_by_size_acquire(sizeof(struct aclk_database_cmd));
605
+
606
worker_register("ACLKSYNC");
607
service_register(SERVICE_THREAD_TYPE_EVENT_LOOP, NULL, NULL, NULL, true);
608
@@ -867,6 +857,8 @@ static void aclk_synchronization(void *arg)
857
858
(void) uv_loop_close(loop);
859
860
+ aral_by_size_release(config->ar);
861
+
862
worker_unregister();
863
service_exits();
864
netdata_log_info("ACLK SYNC: Shutting down ACLK synchronization event loop");