Code cleanup and improvements (#20323)
* Refactor worker data variable for clarity and consistency Replaced "worker_data" with "worker" across functions for improved readability and uniformity. Removed unnecessary comments and unused code related to worker structures. * Change get_worker to initialize request data * No need to set worker.request.data No need to check for worker * Fix shutdown loop timeout calculation (unit conversion was misleading although correct)
Stelios Fragkakis committed
May 23, 2025 at 19:44 UTC
ac1aa30b37c28273027a594c2fd8e27f63217300
3 files changed
+64
-84
src/daemon/libuv_workers.c
+9
-5
@@ -129,15 +129,19 @@ void init_worker_pool(WorkerPool *pool) {
129
pool->top = MAX_ACTIVE_WORKERS; // All workers are initially free
130
}
131
132
-// Get a worker (reuse if available, NULL if pool exhausted)
132
+// Get a worker from the pool
133
+// Needs to be called from the uv event loop thread
134
worker_data_t *get_worker(WorkerPool *pool) {
135
+ worker_data_t *worker;
136
if (pool->top == 0) {
135
- worker_data_t *worker = callocz(1, sizeof(worker_data_t));
137
+ worker = callocz(1, sizeof(worker_data_t));
138
worker->allocated = true; // Mark as allocated
137
- return worker;
139
+ } else {
140
+ int index = pool->free_stack[--pool->top]; // Pop from stack
141
+ worker = &pool->workers[index];
142
}
139
- int index = pool->free_stack[--pool->top]; // Pop from stack
140
- return &pool->workers[index];
143
+ worker->request.data = worker;
144
+ return worker;
145
}
146
147
// Return a worker for reuse
src/database/sqlite/sqlite_aclk.c
+43
-66
@@ -301,10 +301,10 @@ static void async_cb(uv_async_t *handle)
301
302
static void after_aclk_run_query_job(uv_work_t *req, int status __maybe_unused)
303
{
304
- worker_data_t *worker_data = req->data;
305
- struct aclk_sync_config_s *config = worker_data->config;
304
+ worker_data_t *worker = req->data;
305
+ struct aclk_sync_config_s *config = worker->config;
306
config->aclk_queries_running--;
307
- return_worker(&config->worker_pool, worker_data);
307
+ return_worker(&config->worker_pool, worker);
308
}
309
310
static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t *query)
@@ -404,9 +404,9 @@ static void aclk_run_query_job(uv_work_t *req)
404
{
405
register_libuv_worker_jobs();
406
407
- worker_data_t *worker_data = req->data;
408
- struct aclk_sync_config_s *config = worker_data->config;
409
- aclk_query_t *query = (aclk_query_t *) worker_data->payload;
407
+ worker_data_t *worker = req->data;
408
+ struct aclk_sync_config_s *config = worker->config;
409
+ aclk_query_t *query = (aclk_query_t *)worker->payload;
410
411
aclk_run_query(config, query);
412
worker_is_idle();
@@ -414,19 +414,19 @@ static void aclk_run_query_job(uv_work_t *req)
414
415
static void after_aclk_execute_batch(uv_work_t *req, int status __maybe_unused)
416
{
417
- worker_data_t *worker_data = req->data;
418
- struct aclk_sync_config_s *config = worker_data->config;
417
+ worker_data_t *worker = req->data;
418
+ struct aclk_sync_config_s *config = worker->config;
419
config->aclk_batch_job_is_running = false;
420
- return_worker(&config->worker_pool, worker_data);
420
+ return_worker(&config->worker_pool, worker);
421
}
422
423
static void aclk_execute_batch(uv_work_t *req)
424
{
425
register_libuv_worker_jobs();
426
427
- worker_data_t *worker_data = req->data;
428
- struct aclk_sync_config_s *config = worker_data->config;
429
- struct judy_list_t *aclk_query_batch = worker_data->payload;
427
+ worker_data_t *worker = req->data;
428
+ struct aclk_sync_config_s *config = worker->config;
429
+ struct judy_list_t *aclk_query_batch = worker->payload;
430
431
if (!aclk_query_batch)
432
return;
@@ -456,12 +456,6 @@ static void aclk_execute_batch(uv_work_t *req)
456
worker_is_idle();
457
}
458
459
-//struct worker_data {
460
-// uv_work_t request;
461
-// void *payload;
462
-// struct aclk_sync_config_s *config;
463
-//};
464
-
459
struct notify_timer_cb_data {
460
void *payload;
461
struct completion *completion;
@@ -469,20 +463,20 @@ struct notify_timer_cb_data {
463
464
static void after_do_unregister_node(uv_work_t *req, int status __maybe_unused)
465
{
472
- worker_data_t *worker_data = req->data;
473
- struct aclk_sync_config_s *config = worker_data->config;
474
- return_worker(&config->worker_pool, worker_data);
466
+ worker_data_t *worker = req->data;
467
+ struct aclk_sync_config_s *config = worker->config;
468
+ return_worker(&config->worker_pool, worker);
469
}
470
471
static void do_unregister_node(uv_work_t *req)
472
{
473
register_libuv_worker_jobs();
474
481
- worker_data_t *worker_data = req->data;
475
+ worker_data_t *worker = req->data;
476
477
worker_is_busy(UV_EVENT_UNREGISTER_NODE);
478
485
- sql_unregister_node(worker_data->payload);
479
+ sql_unregister_node(worker->payload);
480
481
worker_is_idle();
482
}
@@ -510,11 +504,11 @@ static void node_update_timer_cb(uv_timer_t *handle)
504
505
static void after_start_alert_push(uv_work_t *req, int status __maybe_unused)
506
{
513
- struct worker_data *data = req->data;
514
- struct aclk_sync_config_s *config = data->config;
507
+ struct worker_data *worker = req->data;
508
+ struct aclk_sync_config_s *config = worker->config;
509
510
config->alert_push_running = false;
517
- return_worker(&config->worker_pool, data);
511
+ return_worker(&config->worker_pool, worker);
512
}
513
514
// Worker thread to scan hosts for pending metadata to store
@@ -542,20 +536,15 @@ static void start_alert_push(uv_work_t *req __maybe_unused)
536
537
int schedule_query_in_worker(uv_loop_t *loop, struct aclk_sync_config_s *config, aclk_query_t *query) {
538
545
- worker_data_t *worker_data = get_worker(&config->worker_pool);
546
- if (!worker_data) {
547
- netdata_log_error("Failed to get worker data");
548
- return -1;
549
- }
550
- worker_data->request.data = worker_data;
551
- worker_data->payload = query;
552
- worker_data->config = config;
539
+ worker_data_t *worker = get_worker(&config->worker_pool);
540
+ worker->payload = query;
541
+ worker->config = config;
542
543
config->aclk_queries_running++;
555
- int rc = uv_queue_work(loop, &worker_data->request, aclk_run_query_job, after_aclk_run_query_job);
544
+ int rc = uv_queue_work(loop, &worker->request, aclk_run_query_job, after_aclk_run_query_job);
545
if (rc) {
546
config->aclk_queries_running--;
558
- return_worker(&config->worker_pool, worker_data);
547
+ return_worker(&config->worker_pool, worker);
548
}
549
return rc;
550
}
@@ -581,14 +570,14 @@ static void timer_cb(uv_timer_t *handle)
570
struct aclk_sync_config_s *config = handle->data;
571
572
if (aclk_online_for_alerts()) {
584
- worker_data_t *worker_data;
585
- if (!config->alert_push_running && (worker_data = get_worker(&config->worker_pool))) {
586
- worker_data->request.data = worker_data;
587
- worker_data->config = config;
573
+ worker_data_t *worker;
574
+ if (!config->alert_push_running) {
575
+ worker = get_worker(&config->worker_pool);
576
+ worker->config = config;
577
config->alert_push_running = true;
589
- if (uv_queue_work(handle->loop, &worker_data->request, start_alert_push, after_start_alert_push)) {
578
+ if (uv_queue_work(handle->loop, &worker->request, start_alert_push, after_start_alert_push)) {
579
config->alert_push_running = false;
591
- return_worker(&config->worker_pool, worker_data);
580
+ return_worker(&config->worker_pool, worker);
581
}
582
}
583
}
@@ -642,7 +631,6 @@ static void *aclk_synchronization_event_loop(void *arg)
631
int query_thread_count = (int) netdata_conf_cloud_query_threads();
632
netdata_log_info("Starting ACLK synchronization thread with %d parallel query threads", query_thread_count);
633
645
- //struct worker_data *worker_datadata;
634
struct notify_timer_cb_data *timer_cb_data;
635
aclk_query_t *query;
636
@@ -653,7 +641,7 @@ static void *aclk_synchronization_event_loop(void *arg)
641
size_t pending_queries = 0;
642
643
Pvoid_t *Pvalue;
656
- worker_data_t *worker_data;
644
+ worker_data_t *worker;
645
646
unsigned cmd_batch_size;
647
@@ -772,19 +760,13 @@ static void *aclk_synchronization_event_loop(void *arg)
760
break;
761
762
case ACLK_DATABASE_NODE_UNREGISTER:
775
- worker_data = get_worker(&config->worker_pool);
776
- if (!worker_data) {
777
- nd_log_daemon(NDLP_ERR, "Failed to get worker for unregister node");
778
- freez(cmd.param[0]);
779
- break;
780
- }
781
- worker_data->request.data = worker_data;
782
- worker_data->config = config;
783
- worker_data->payload = cmd.param[0];
763
+ worker = get_worker(&config->worker_pool);
764
+ worker->config = config;
765
+ worker->payload = cmd.param[0];
766
785
- if (uv_queue_work(loop, &worker_data->request, do_unregister_node, after_do_unregister_node)) {
767
+ if (uv_queue_work(loop, &worker->request, do_unregister_node, after_do_unregister_node)) {
768
freez(cmd.param[0]);
787
- return_worker(&config->worker_pool, worker_data);
769
+ return_worker(&config->worker_pool, worker);
770
}
771
break;
772
case ACLK_DATABASE_PUSH_ALERT_CONFIG:
@@ -884,23 +866,18 @@ static void *aclk_synchronization_event_loop(void *arg)
866
if (!aclk_query_batch || config->aclk_batch_job_is_running)
867
break;
868
887
- worker_data = get_worker(&config->worker_pool);
888
- if (!worker_data) {
889
- nd_log_daemon(NDLP_ERR, "Failed to get worker for ACLK batch job");
890
- break;
891
- }
892
- worker_data->request.data = worker_data;
893
- worker_data->config = config;
894
- worker_data->payload = aclk_query_batch;
869
+ worker = get_worker(&config->worker_pool);
870
+ worker->config = config;
871
+ worker->payload = aclk_query_batch;
872
873
config->aclk_batch_job_is_running = true;
874
config->aclk_jobs_pending -= aclk_query_batch->count;
875
aclk_query_batch = NULL;
876
900
- if (uv_queue_work(loop, &worker_data->request, aclk_execute_batch, after_aclk_execute_batch)) {
901
- aclk_query_batch = worker_data->payload;
877
+ if (uv_queue_work(loop, &worker->request, aclk_execute_batch, after_aclk_execute_batch)) {
878
+ aclk_query_batch = worker->payload;
879
config->aclk_jobs_pending += aclk_query_batch->count;
903
- return_worker(&config->worker_pool, worker_data);
880
+ return_worker(&config->worker_pool, worker);
881
config->aclk_batch_job_is_running = false;
882
}
883
break;
src/database/sqlite/sqlite_metadata.c
+12
-13
@@ -1847,8 +1847,8 @@ static void ctx_hosts_load(uv_work_t *req)
1847
{
1848
register_libuv_worker_jobs();
1849
1850
- worker_data_t *data = req->data;
1851
- struct meta_config_s *config = data->config;
1850
+ worker_data_t *worker = req->data;
1851
+ struct meta_config_s *config = worker->config;
1852
1853
worker_is_busy(UV_EVENT_HOST_CONTEXT_LOAD);
1854
usec_t started_ut = now_monotonic_usec(); (void)started_ut;
@@ -2394,18 +2394,18 @@ static void start_metadata_hosts(uv_work_t *req)
2394
{
2395
register_libuv_worker_jobs();
2396
2397
- worker_data_t *data = req->data;
2398
- struct meta_config_s *config = data->config;
2397
+ worker_data_t *worker = req->data;
2398
+ struct meta_config_s *config = worker->config;
2399
2400
- BUFFER *work_buffer = data->work_buffer;
2400
+ BUFFER *work_buffer = worker->work_buffer;
2401
usec_t all_started_ut = now_monotonic_usec();
2402
2403
- store_sql_statements((struct judy_list_t *)data->pending_sql_statement, true);
2403
+ store_sql_statements((struct judy_list_t *)worker->pending_sql_statement, true);
2404
2405
- store_alert_transitions((struct judy_list_t *)data->pending_alert_list, true);
2405
+ store_alert_transitions((struct judy_list_t *)worker->pending_alert_list, true);
2406
2407
if (!SHUTDOWN_REQUESTED(config))
2408
- store_ctx_cleanup_list(config, (struct judy_list_t *)data->pending_ctx_cleanup_list);
2408
+ store_ctx_cleanup_list(config, (struct judy_list_t *)worker->pending_ctx_cleanup_list);
2409
2410
worker_is_busy(UV_EVENT_METADATA_STORE);
2411
@@ -2415,7 +2415,7 @@ static void start_metadata_hosts(uv_work_t *req)
2415
nd_log_daemon(NDLP_DEBUG, "Checking all hosts completed in %s", report_duration);
2416
2417
if (!SHUTDOWN_REQUESTED(config)) {
2418
- do_pending_uuid_deletion(config, (struct judy_list_t *)data->pending_uuid_deletion);
2418
+ do_pending_uuid_deletion(config, (struct judy_list_t *)worker->pending_uuid_deletion);
2419
run_metadata_cleanup(config);
2420
}
2421
@@ -2544,7 +2544,6 @@ static void *metadata_event_loop(void *arg)
2544
break;
2545
2546
worker = get_worker(&config->worker_pool);
2547
- worker->request.data = worker;
2547
worker->config = config;
2548
worker->pending_alert_list = pending_alert_list;
2549
worker->pending_ctx_cleanup_list = pending_ctx_cleanup_list;
@@ -2572,7 +2571,6 @@ static void *metadata_event_loop(void *arg)
2571
2572
worker = get_worker(&config->worker_pool);
2573
config->ctx_load_running = true;
2575
- worker->request.data = worker;
2574
worker->config = config;
2575
if (uv_queue_work(loop, &worker->request, ctx_hosts_load, after_ctx_hosts_load)) {
2576
config->ctx_load_running = false;
@@ -2635,11 +2633,12 @@ static void *metadata_event_loop(void *arg)
2633
uv_close((uv_handle_t *)&config->async, NULL);
2634
uv_walk(loop, libuv_close_callback, NULL);
2635
2638
- size_t loop_count = (MAX_SHUTDOWN_TIMEOUT_SECONDS * USEC_PER_MS) / SHUTDOWN_SLEEP_INTERVAL_MS;
2639
- while ((config->metadata_running || config->ctx_load_running) && --loop_count) {
2636
+ size_t loop_count = (MAX_SHUTDOWN_TIMEOUT_SECONDS * MSEC_PER_SEC) / SHUTDOWN_SLEEP_INTERVAL_MS;
2637
+ while ((config->metadata_running || config->ctx_load_running) && loop_count > 0) {
2638
if (!uv_run(loop, UV_RUN_NOWAIT))
2639
break; // No pending callbacks
2640
sleep_usec(SHUTDOWN_SLEEP_INTERVAL_MS * USEC_PER_MS);
2641
+ loop_count--;
2642
}
2643
2644
(void)uv_loop_close(loop);