Use worker when dispatching alert transitions to the cloud (#19397)
* Queue alert submission to worker * Correctly register aclk job execution * Correctly register aclk job execution (rename) * Remove job
Stelios Fragkakis committed
Jan 15, 2025 at 01:13 UTC
de14425b80aa99434cb76e3f12dd5e8e520b7338
5 files changed
+61
-12
src/daemon/libuv_workers.c
+5
@@ -56,6 +56,11 @@ void register_libuv_worker_jobs() {
56
worker_register_job_name(UV_EVENT_METADATA_CLEANUP, "metadata cleanup");
57
worker_register_job_name(UV_EVENT_METADATA_ML_LOAD, "metadata load ml models");
58
59
+ // aclk_sync
60
+ worker_register_job_name(UV_EVENT_ACLK_NODE_INFO, "aclk host node info");
61
+ worker_register_job_name(UV_EVENT_ACLK_ALERT_PUSH, "aclk alert push");
62
+ worker_register_job_name(UV_EVENT_ACLK_QUERY_EXECUTE, "aclk query execute");
63
+
64
// netdatacli
65
worker_register_job_name(UV_EVENT_SCHEDULE_CMD, "schedule command");
66
src/daemon/libuv_workers.h
+5
@@ -48,6 +48,11 @@ enum event_loop_job {
48
UV_EVENT_METADATA_CLEANUP,
49
UV_EVENT_METADATA_ML_LOAD,
50
51
+ // aclk_sync
52
+ UV_EVENT_ACLK_NODE_INFO,
53
+ UV_EVENT_ACLK_ALERT_PUSH,
54
+ UV_EVENT_ACLK_QUERY_EXECUTE,
55
+
56
// netdatacli
57
UV_EVENT_SCHEDULE_CMD,
58
};
src/database/sqlite/sqlite_aclk.c
+51
-7
@@ -20,6 +20,7 @@ struct aclk_sync_config_s {
20
bool initialized;
21
mqtt_wss_client client;
22
int aclk_queries_running;
23
+ bool alert_push_running;
24
SPINLOCK cmd_queue_lock;
25
struct aclk_database_cmd *cmd_base;
26
} aclk_sync_config = { 0 };
@@ -285,7 +286,6 @@ static void timer_cb(uv_timer_t *handle)
286
if (aclk_online_for_alerts()) {
287
cmd.opcode = ACLK_DATABASE_PUSH_ALERT;
288
aclk_database_enq_cmd(&cmd);
288
- aclk_check_node_info_and_collectors();
289
}
290
}
291
@@ -297,7 +297,6 @@ struct aclk_query_payload {
297
298
static void after_aclk_run_query_job(uv_work_t *req, int status __maybe_unused)
299
{
300
- worker_is_busy(ACLK_QUERY_EXECUTE);
300
struct aclk_query_payload *payload = req->data;
301
struct aclk_sync_config_s *config = payload->config;
302
config->aclk_queries_running--;
@@ -321,11 +320,16 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
320
321
static void aclk_run_query_job(uv_work_t *req)
322
{
323
+ register_libuv_worker_jobs();
324
+
325
+ worker_is_busy(UV_EVENT_ACLK_QUERY_EXECUTE);
326
+
327
struct aclk_query_payload *payload = req->data;
328
struct aclk_sync_config_s *config = payload->config;
329
aclk_query_t query = (aclk_query_t) payload->data;
330
331
aclk_run_query(config, query);
332
+ worker_is_idle();
333
}
334
335
static void node_update_timer_cb(uv_timer_t *handle)
@@ -346,6 +350,34 @@ static void close_callback(uv_handle_t *handle, void *data __maybe_unused)
350
uv_close(handle, NULL); // Automatically close and free the handle
351
}
352
353
+struct alert_push_data {
354
+ uv_work_t request;
355
+ struct aclk_sync_config_s *config;
356
+};
357
+
358
+static void after_start_alert_push(uv_work_t *req, int status __maybe_unused)
359
+{
360
+ struct alert_push_data *data = req->data;
361
+ struct aclk_sync_config_s *config = data->config;
362
+
363
+ config->alert_push_running = false;
364
+ freez(data);
365
+}
366
+
367
+// Worker thread to scan hosts for pending metadata to store
368
+static void start_alert_push(uv_work_t *req __maybe_unused)
369
+{
370
+ register_libuv_worker_jobs();
371
+
372
+ worker_is_busy(UV_EVENT_ACLK_NODE_INFO);
373
+ aclk_check_node_info_and_collectors();
374
+ worker_is_idle();
375
+
376
+ worker_is_busy(UV_EVENT_ACLK_ALERT_PUSH);
377
+ aclk_push_alert_events_for_all_hosts();
378
+ worker_is_idle();
379
+}
380
+
381
static void aclk_synchronization(void *arg)
382
{
383
struct aclk_sync_config_s *config = arg;
@@ -357,9 +389,7 @@ static void aclk_synchronization(void *arg)
389
worker_register_job_name(ACLK_DATABASE_NODE_STATE, "node state");
390
worker_register_job_name(ACLK_DATABASE_PUSH_ALERT, "alert push");
391
worker_register_job_name(ACLK_DATABASE_PUSH_ALERT_CONFIG, "alert conf push");
360
- worker_register_job_name(ACLK_QUERY_EXECUTE, "query execute");
361
- worker_register_job_name(ACLK_QUERY_EXECUTE_SYNC, "query execute sync");
362
- worker_register_job_name(ACLK_DATABASE_TIMER, "timer");
392
+ worker_register_job_name(ACLK_QUERY_EXECUTE_SYNC, "aclk query execute sync");
393
394
uv_loop_t *loop = &config->loop;
395
fatal_assert(0 == uv_loop_init(loop));
@@ -378,6 +408,8 @@ static void aclk_synchronization(void *arg)
408
int query_thread_count = netdata_conf_cloud_query_threads();
409
netdata_log_info("Starting ACLK synchronization thread with %d parallel query threads", query_thread_count);
410
411
+ struct alert_push_data *data;
412
+
413
while (likely(service_running(SERVICE_ACLK))) {
414
enum aclk_database_opcode opcode;
415
worker_is_idle();
@@ -441,9 +473,21 @@ static void aclk_synchronization(void *arg)
473
aclk_push_alert_config_event(cmd.param[0], cmd.param[1]);
474
break;
475
case ACLK_DATABASE_PUSH_ALERT:
444
- aclk_push_alert_events_for_all_hosts();
445
- break;
476
477
+ if (config->alert_push_running)
478
+ break;
479
+
480
+ config->alert_push_running = true;
481
+
482
+ data = mallocz(sizeof(*data));
483
+ data->request.data = data;
484
+ data->config = config;
485
+
486
+ if (uv_queue_work(loop, &data->request, start_alert_push, after_start_alert_push)) {
487
+ freez(data);
488
+ config->alert_push_running = false;
489
+ }
490
+ break;
491
case ACLK_MQTT_WSS_CLIENT:
492
config->client = (mqtt_wss_client) cmd.param[0];
493
break;
src/database/sqlite/sqlite_aclk.h
-1
@@ -24,7 +24,6 @@ enum aclk_database_opcode {
24
ACLK_MQTT_WSS_CLIENT,
25
ACLK_QUERY_EXECUTE,
26
ACLK_QUERY_EXECUTE_SYNC,
27
- ACLK_DATABASE_TIMER,
27
28
// leave this last
29
// we need it to check for worker utilization
src/database/sqlite/sqlite_aclk_alert.c
-4
@@ -594,10 +594,6 @@ void aclk_push_alert_events_for_all_hosts(void)
594
{
595
RRDHOST *host;
596
597
- // Checking if we shutting down
598
- if (!service_running(SERVICE_ACLK))
599
- return;
600
-
597
dfe_start_reentrant(rrdhost_root_index, host) {
598
if (!rrdhost_flag_check(host, RRDHOST_FLAG_ACLK_STREAM_ALERTS) ||
599
rrdhost_flag_check(host, RRDHOST_FLAG_PENDING_CONTEXT_LOAD))