Event loop cleanup (#21091)
event loop cleanup
Stelios Fragkakis committed
Oct 5, 2025 at 10:10 UTC
e1ef92522ffd326eb52111792830931d5a8ca3e7
3 files changed
+17
-36
src/database/engine/rrdengine.c
+7
-23
@@ -270,9 +270,6 @@ static void work_standard_worker(uv_work_t *req) {
270
271
__atomic_sub_fetch(&rrdeng_main.work_cmd.atomics.dispatched, 1, __ATOMIC_RELAXED);
272
__atomic_sub_fetch(&rrdeng_main.work_cmd.atomics.executing, 1, __ATOMIC_RELAXED);
273
-
274
- // signal the event loop a worker is available
275
- rrdeng_async_wakeup();
273
}
274
275
static void after_work_standard_callback(uv_work_t* req, int status) {
@@ -1890,7 +1887,6 @@ void async_cb(uv_async_t *handle)
1887
1888
#define TIMER_PERIOD_MS (1000)
1889
1893
-
1890
static void *extent_read_tp_worker(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t *uv_work_req __maybe_unused) {
1891
EPDL *epdl = data;
1892
epdl_find_extent_and_populate_pages(ctx, epdl, true);
@@ -2162,13 +2158,12 @@ bool rrdeng_ctx_tier_cap_exceeded(struct rrdengine_instance *ctx)
2158
return false;
2159
}
2160
2165
-static void retention_timer_cb(uv_timer_t *handle) {
2161
+static void retention_timer_cb(uv_timer_t *handle __maybe_unused)
2162
+{
2163
if (!localhost)
2164
return;
2165
2166
worker_is_busy(RRDENG_RETENTION_TIMER_CB);
2170
- uv_stop(handle->loop);
2171
- uv_update_time(handle->loop);
2167
2168
for (size_t tier = 0; tier < nd_profile.storage_tiers; tier++) {
2169
STORAGE_ENGINE *eng = localhost->db[tier].eng;
@@ -2180,18 +2175,14 @@ static void retention_timer_cb(uv_timer_t *handle) {
2175
worker_is_idle();
2176
}
2177
2183
-static void timer_per_sec_cb(uv_timer_t* handle) {
2178
+static void timer_per_sec_cb(uv_timer_t *handle __maybe_unused)
2179
+{
2180
worker_is_busy(RRDENG_TIMER_CB);
2185
- uv_stop(handle->loop);
2186
- uv_update_time(handle->loop);
2181
2182
worker_set_metric(RRDENG_OPCODES_WAITING, (NETDATA_DOUBLE)rrdeng_main.cmd_queue.unsafe.waiting);
2183
worker_set_metric(RRDENG_WORKS_DISPATCHED, (NETDATA_DOUBLE)__atomic_load_n(&rrdeng_main.work_cmd.atomics.dispatched, __ATOMIC_RELAXED));
2184
worker_set_metric(RRDENG_WORKS_EXECUTING, (NETDATA_DOUBLE)__atomic_load_n(&rrdeng_main.work_cmd.atomics.executing, __ATOMIC_RELAXED));
2185
2192
- // rrdeng_enq_cmd(NULL, RRDENG_OPCODE_EVICT_MAIN, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
2193
- // rrdeng_enq_cmd(NULL, RRDENG_OPCODE_EVICT_OPEN, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
2194
- // rrdeng_enq_cmd(NULL, RRDENG_OPCODE_EVICT_EXTENT, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
2186
rrdeng_enq_cmd(NULL, RRDENG_OPCODE_FLUSH_MAIN, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
2187
rrdeng_enq_cmd(NULL, RRDENG_OPCODE_CLEANUP, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
2188
@@ -2412,18 +2403,9 @@ void dbengine_event_loop(void* arg) {
2403
2404
while (likely(!shutdown)) {
2405
worker_is_idle();
2415
- uv_run(&main->loop, UV_RUN_DEFAULT);
2406
+ uv_run(&main->loop, UV_RUN_ONCE);
2407
2417
- /* wait for commands */
2418
- size_t count = 0;
2408
do {
2420
- count++;
2421
-
2422
- if(count % 100 == 0) {
2423
- worker_is_idle();
2424
- uv_run(&main->loop, UV_RUN_NOWAIT);
2425
- }
2426
-
2409
worker_is_busy(RRDENG_OPCODE_MAX);
2410
cmd = rrdeng_deq_cmd(RRDENG_OPCODE_NOOP);
2411
opcode = cmd.opcode;
@@ -2586,6 +2568,8 @@ void dbengine_event_loop(void* arg) {
2568
break;
2569
}
2570
}
2571
+ if (opcode != RRDENG_OPCODE_NOOP)
2572
+ uv_run(&main->loop, UV_RUN_NOWAIT);
2573
2574
} while (opcode != RRDENG_OPCODE_NOOP);
2575
}
src/database/sqlite/sqlite_aclk.c
+5
-8
@@ -650,8 +650,6 @@ static void aclk_synchronization_event_loop(void *arg)
650
Pvoid_t *Pvalue;
651
worker_data_t *worker;
652
653
- unsigned cmd_batch_size;
654
-
653
config->shutdown_requested = false;
654
config->initialized = true;
655
completion_mark_complete(&config->start_stop_complete);
@@ -662,16 +660,11 @@ static void aclk_synchronization_event_loop(void *arg)
660
struct aclk_sync_cfg_t *aclk_host_config;
661
aclk_query_t *query;
662
worker_is_idle();
665
- uv_run(loop, UV_RUN_DEFAULT);
663
+ uv_run(loop, UV_RUN_ONCE);
664
667
- /* wait for commands */
668
- cmd_batch_size = 0;
665
do {
666
cmd_data_t cmd;
667
672
- if (unlikely(++cmd_batch_size >= MAX_BATCH_SIZE))
673
- break;
674
-
668
if (config->run_query_batch) {
669
opcode = ACLK_QUERY_BATCH_EXECUTE;
670
config->run_query_batch = false;
@@ -886,6 +879,10 @@ static void aclk_synchronization_event_loop(void *arg)
879
default:
880
break;
881
}
882
+
883
+ if (opcode != ACLK_DATABASE_NOOP)
884
+ uv_run(loop, UV_RUN_NOWAIT);
885
+
886
} while (opcode != ACLK_DATABASE_NOOP);
887
}
888
config->initialized = false;
src/database/sqlite/sqlite_metadata.c
+5
-5
@@ -2521,14 +2521,10 @@ static void metadata_event_loop(void *arg)
2521
enum metadata_opcode opcode;
2522
2523
worker_is_idle();
2524
- uv_run(loop, UV_RUN_DEFAULT);
2524
+ uv_run(loop, UV_RUN_ONCE);
2525
2526
/* wait for commands */
2527
- unsigned cmd_batch_size = 0;
2527
do {
2529
- if (unlikely(++cmd_batch_size >= METADATA_MAX_BATCH_SIZE))
2530
- break;
2531
-
2528
cmd_data_t cmd;
2529
if (config->store_metadata && !config->metadata_running) {
2530
config->store_metadata = false;
@@ -2664,6 +2660,10 @@ static void metadata_event_loop(void *arg)
2660
default:
2661
break;
2662
}
2663
+
2664
+ if (likely(opcode != METADATA_DATABASE_NOOP))
2665
+ uv_run(loop, UV_RUN_NOWAIT);
2666
+
2667
} while (opcode != METADATA_DATABASE_NOOP);
2668
}
2669
config->initialized = false;