Store alert config asynchronously (#19885)
Add support for running prepared SQL statements that store data, in the metadata thread
Stelios Fragkakis committed
Mar 20, 2025 at 09:20 UTC
70f4422287db1077c212bf3828f1666e991ccff0
5 files changed
+103
-41
src/daemon/libuv_workers.c
+1
@@ -53,6 +53,7 @@ static void register_libuv_worker_jobs_internal(void) {
53
worker_register_job_name(UV_EVENT_CTX_CLEANUP_SCHEDULE, "metadata ctx cleanup schedule");
54
worker_register_job_name(UV_EVENT_CTX_CLEANUP, "metadata ctx cleanup");
55
worker_register_job_name(UV_EVENT_STORE_ALERT_TRANSITIONS, "metadata store alert transitions");
56
+ worker_register_job_name(UV_EVENT_STORE_SQL_STATEMENTS, "metadata store sql statements");
57
worker_register_job_name(UV_EVENT_CHART_LABEL_CLEANUP, "metadata chart label cleanup");
58
worker_register_job_name(UV_EVENT_HEALTH_LOG_CLEANUP, "alert transitions cleanup");
59
worker_register_job_name(UV_EVENT_UUID_DELETION, "metadata dimension deletion");
src/daemon/libuv_workers.h
+1
@@ -54,6 +54,7 @@ enum event_loop_job {
54
UV_EVENT_STORE_CHART,
55
UV_EVENT_STORE_DIMENSION,
56
UV_EVENT_STORE_ALERT_TRANSITIONS,
57
+ UV_EVENT_STORE_SQL_STATEMENTS,
58
UV_EVENT_HEALTH_LOG_CLEANUP,
59
UV_EVENT_CHART_LABEL_CLEANUP,
60
UV_EVENT_UUID_DELETION,
src/database/sqlite/sqlite_health.c
+35
-35
@@ -12,6 +12,9 @@ extern __thread bool is_health_thread;
12
#define SQLITE3_BIND_STRING_OR_NULL(res, param, key) \
13
((key) ? sqlite3_bind_text((res), (param), string2str(key), -1, SQLITE_STATIC) : sqlite3_bind_null((res), (param)))
14
15
+#define SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, param, key) \
16
+ ((key) ? sqlite3_bind_text((res), (param), string2str(key), -1, SQLITE_TRANSIENT) : sqlite3_bind_null((res), (param)))
17
+
18
#define SQLITE3_COLUMN_STRINGDUP_OR_NULL(res, param) \
19
({ \
20
int _param = (param); \
@@ -814,37 +817,37 @@ done:
817
818
void sql_alert_store_config(RRD_ALERT_PROTOTYPE *ap)
819
{
817
- static __thread sqlite3_stmt *res = NULL;
820
+ sqlite3_stmt *res = NULL;
821
int param = 0;
822
820
- if (!PREPARE_COMPILED_STATEMENT(db_meta, SQL_STORE_ALERT_CONFIG_HASH, &res))
823
+ if (!PREPARE_STATEMENT(db_meta, SQL_STORE_ALERT_CONFIG_HASH, &res))
824
return;
825
826
CLEAN_BUFFER *buf = buffer_create(128, NULL);
827
828
SQLITE_BIND_FAIL(
826
- done, sqlite3_bind_blob(res, ++param, &ap->config.hash_id, sizeof(ap->config.hash_id), SQLITE_STATIC));
829
+ done, sqlite3_bind_blob(res, ++param, &ap->config.hash_id, sizeof(ap->config.hash_id), SQLITE_TRANSIENT));
830
831
if (ap->match.is_template) {
829
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, NULL));
830
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.name));
831
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->match.on.context));
832
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, NULL));
833
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.name));
834
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->match.on.context));
835
}
836
else {
834
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.name));
835
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, NULL));
836
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->match.on.chart));
837
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.name));
838
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, NULL));
839
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->match.on.chart));
840
}
841
839
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.classification));
840
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.component));
841
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.type));
842
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, NULL)); // lookup
842
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.classification));
843
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.component));
844
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.type));
845
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, NULL)); // lookup
846
SQLITE_BIND_FAIL(done, sqlite3_bind_int(res, ++param, ap->config.update_every));
844
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.units));
847
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.units));
848
849
if (ap->config.calculation)
847
- SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, expression_source(ap->config.calculation), -1, SQLITE_STATIC));
850
+ SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, expression_source(ap->config.calculation), -1, SQLITE_TRANSIENT));
851
else
852
SQLITE_BIND_FAIL(done,sqlite3_bind_null(res, ++param));
853
@@ -853,18 +856,18 @@ void sql_alert_store_config(RRD_ALERT_PROTOTYPE *ap)
856
SQLITE_BIND_FAIL(done, sqlite3_bind_double(res, ++param, nan_value));
857
858
if (ap->config.warning)
856
- SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, expression_source(ap->config.warning), -1, SQLITE_STATIC));
859
+ SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, expression_source(ap->config.warning), -1, SQLITE_TRANSIENT));
860
else
861
SQLITE_BIND_FAIL(done, sqlite3_bind_null(res, ++param));
862
863
if (ap->config.critical)
861
- SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, expression_source(ap->config.critical), -1, SQLITE_STATIC));
864
+ SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, expression_source(ap->config.critical), -1, SQLITE_TRANSIENT));
865
else
866
SQLITE_BIND_FAIL(done, sqlite3_bind_null(res, ++param));
867
865
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.exec));
866
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.recipient));
867
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.info));
868
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.exec));
869
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.recipient));
870
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.info));
871
872
if (ap->config.delay_up_duration)
873
buffer_sprintf(buf, "up %ds ", ap->config.delay_up_duration);
@@ -879,10 +882,10 @@ void sql_alert_store_config(RRD_ALERT_PROTOTYPE *ap)
882
buffer_sprintf(buf, "max %ds", ap->config.delay_max_duration);
883
884
// delay
882
- SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, buffer_tostring(buf), -1, SQLITE_STATIC));
885
+ SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, buffer_tostring(buf), -1, SQLITE_TRANSIENT));
886
887
if (ap->config.alert_action_options & ALERT_ACTION_OPTION_NO_CLEAR_NOTIFICATION)
885
- SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, "no-clear-notification", -1, SQLITE_STATIC));
888
+ SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, "no-clear-notification", -1, SQLITE_TRANSIENT));
889
else
890
SQLITE_BIND_FAIL(done, sqlite3_bind_null(res, ++param));
891
@@ -891,14 +894,14 @@ void sql_alert_store_config(RRD_ALERT_PROTOTYPE *ap)
894
SQLITE_BIND_FAIL(done, sqlite3_bind_null(res, ++param));
895
else {
896
snprintfz(repeat, sizeof(repeat) - 1, "warning %us critical %us", ap->config.warn_repeat_every, ap->config.crit_repeat_every);
894
- SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, repeat, -1, SQLITE_STATIC));
897
+ SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, repeat, -1, SQLITE_TRANSIENT));
898
}
899
897
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->match.host_labels));
900
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->match.host_labels));
901
902
if (ap->config.after) {
900
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.dimensions));
901
- SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, time_grouping_id2txt(ap->config.time_group), -1, SQLITE_STATIC));
903
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.dimensions));
904
+ SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, time_grouping_id2txt(ap->config.time_group), -1, SQLITE_TRANSIENT));
905
SQLITE_BIND_FAIL(done, sqlite3_bind_int(res, ++param, (int) RRDR_OPTIONS_REMOVE_OVERLAPPING(ap->config.options)));
906
SQLITE_BIND_FAIL(done, sqlite3_bind_int64(res, ++param, (int) ap->config.after));
907
SQLITE_BIND_FAIL(done, sqlite3_bind_int64(res, ++param, (int) ap->config.before));
@@ -911,23 +914,20 @@ void sql_alert_store_config(RRD_ALERT_PROTOTYPE *ap)
914
}
915
916
SQLITE_BIND_FAIL(done, sqlite3_bind_int(res, ++param, ap->config.update_every));
914
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.source));
915
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->match.chart_labels));
916
- SQLITE_BIND_FAIL(done, SQLITE3_BIND_STRING_OR_NULL(res, ++param, ap->config.summary));
917
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.source));
918
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->match.chart_labels));
919
+ SQLITE_BIND_FAIL(done, SQLITE3_BIND_TRANSIENT_STRING_OR_NULL(res, ++param, ap->config.summary));
920
921
SQLITE_BIND_FAIL(done, sqlite3_bind_int(res, ++param, ap->config.time_group_condition));
922
SQLITE_BIND_FAIL(done, sqlite3_bind_double(res, ++param, ap->config.time_group_value));
923
SQLITE_BIND_FAIL(done, sqlite3_bind_int(res, ++param, ap->config.dims_group));
924
SQLITE_BIND_FAIL(done, sqlite3_bind_int(res, ++param, ap->config.data_source));
925
923
- param = 0;
924
- int rc = sqlite3_step_monitored(res);
925
- if (unlikely(rc != SQLITE_DONE))
926
- error_report("Failed to store alert config, rc = %d", rc);
927
-
926
+ metadata_execute_store_statement(res);
927
+ return;
928
done:
929
REPORT_BIND_FAIL(res, param);
930
- SQLITE_RESET(res);
930
+ SQLITE_FINALIZE(res);
931
}
932
933
#define SQL_SELECT_HEALTH_LAST_EXECUTED_EVENT \
src/database/sqlite/sqlite_metadata.c
+65
-6
@@ -8,6 +8,12 @@
8
9
#define DB_METADATA_VERSION 18
10
11
+#define COMPUTE_DURATION(var_name, unit, start, end) \
12
+ char var_name[64]; \
13
+ duration_snprintf(var_name, sizeof(var_name), \
14
+ (int64_t)((end) - (start)), unit, true)
15
+
16
+
17
extern long long def_journal_size_limit;
18
19
const char *database_config[] = {
@@ -199,6 +205,7 @@ enum metadata_opcode {
205
METADATA_ADD_HOST_AE,
206
METADATA_DEL_HOST_AE,
207
METADATA_ADD_CTX_CLEANUP,
208
+ METADATA_EXECUTE_STORE_STATEMENT,
209
METADATA_MAINTENANCE,
210
METADATA_SYNC_SHUTDOWN,
211
METADATA_UNITTEST,
@@ -1821,6 +1828,7 @@ struct scan_metadata_payload {
1828
void *pending_alert_list;
1829
void *pending_ctx_cleanup_list;
1830
void *pending_uuid_deletion;
1831
+ void *pending_sql_statement;
1832
BUFFER *work_buffer;
1833
};
1834
@@ -2291,6 +2299,40 @@ static void store_alert_transitions(struct judy_list_t *pending_alert_list)
2299
worker_is_idle();
2300
}
2301
2302
+static void store_sql_statements(struct judy_list_t *pending_sql_statement)
2303
+{
2304
+ if (!pending_sql_statement)
2305
+ return;
2306
+
2307
+ worker_is_busy(METADATA_EXECUTE_STORE_STATEMENT);
2308
+
2309
+ usec_t started_ut = now_monotonic_usec();
2310
+
2311
+ size_t entries = pending_sql_statement->count;
2312
+ Word_t Index = 0;
2313
+ bool first = true;
2314
+ Pvoid_t *Pvalue;
2315
+ while ((Pvalue = JudyLFirstThenNext(pending_sql_statement->JudyL, &Index, &first))) {
2316
+ sqlite3_stmt *stmt = *Pvalue;
2317
+
2318
+ if (unlikely(!stmt))
2319
+ continue;
2320
+
2321
+ int rc = sqlite3_step_monitored(stmt);
2322
+ if (unlikely(rc != SQLITE_DONE))
2323
+ nd_log_daemon(NDLP_ERR, "Failed to execute sql statement, rc = %d", rc);
2324
+
2325
+ SQLITE_FINALIZE(stmt);
2326
+ }
2327
+ (void) JudyLFreeArray(&pending_sql_statement->JudyL, PJE0);
2328
+ freez(pending_sql_statement);
2329
+
2330
+ COMPUTE_DURATION(report_duration, "us", started_ut, now_monotonic_usec());
2331
+ nd_log_daemon(NDLP_DEBUG, "Stored and processed %zu sql statements in %s", entries, report_duration);
2332
+
2333
+ worker_is_idle();
2334
+}
2335
+
2336
static void meta_store_host_labels(RRDHOST *host, BUFFER *work_buffer, size_t *query_counter)
2337
{
2338
rrdhost_flag_clear(host, RRDHOST_FLAG_METADATA_LABELS);
@@ -2336,12 +2378,6 @@ static void store_host_claim_id(RRDHOST *host, size_t *query_counter)
2378
(*query_counter)++;
2379
}
2380
2339
-#define COMPUTE_DURATION(var_name, unit, start, end) \
2340
- char var_name[64]; \
2341
- duration_snprintf(var_name, sizeof(var_name), \
2342
- (int64_t)((end) - (start)), unit, true)
2343
-
2344
-
2381
void store_host_info_and_metadata(RRDHOST *host, BUFFER *work_buffer, size_t *query_counter)
2382
{
2383
// Store labels (if needed)
@@ -2370,6 +2406,7 @@ static void start_metadata_hosts(uv_work_t *req)
2406
BUFFER *work_buffer = data->work_buffer;
2407
usec_t all_started_ut = now_monotonic_usec();
2408
2409
+ store_sql_statements((struct judy_list_t *)data->pending_sql_statement);
2410
store_alert_transitions((struct judy_list_t *)data->pending_alert_list);
2411
store_ctx_cleanup_list(wc, (struct judy_list_t *)data->pending_ctx_cleanup_list);
2412
@@ -2437,6 +2474,7 @@ static void metadata_event_loop(void *arg)
2474
worker_register_job_name(METADATA_LOAD_HOST_CONTEXT, "host load context");
2475
worker_register_job_name(METADATA_ADD_HOST_AE, "add host alert entry");
2476
worker_register_job_name(METADATA_DEL_HOST_AE, "delete host alert entry");
2477
+ worker_register_job_name(METADATA_EXECUTE_STORE_STATEMENT, "add sql statement");
2478
2479
int ret;
2480
unsigned cmd_batch_size;
@@ -2486,11 +2524,13 @@ static void metadata_event_loop(void *arg)
2524
struct judy_list_t *pending_ae_list = NULL;
2525
struct judy_list_t *pending_ctx_cleanup_list = NULL;
2526
struct judy_list_t *pending_uuid_deletion = NULL;
2527
+ struct judy_list_t *pending_sql_statement = NULL;
2528
2529
while (shutdown == 0 || (wc->flags & METADATA_FLAG_PROCESSING)) {
2530
nd_uuid_t *uuid;
2531
RRDHOST *host = NULL;
2532
ALARM_ENTRY *ae = NULL;
2533
+ sqlite3_stmt *stmt;
2534
2535
worker_is_idle();
2536
uv_run(loop, UV_RUN_DEFAULT);
@@ -2565,11 +2605,13 @@ static void metadata_event_loop(void *arg)
2605
data->pending_alert_list = pending_ae_list;
2606
data->pending_ctx_cleanup_list = pending_ctx_cleanup_list;
2607
data->pending_uuid_deletion = pending_uuid_deletion;
2608
+ data->pending_sql_statement = pending_sql_statement;
2609
2610
data->work_buffer = work_buffer;
2611
pending_ae_list = NULL;
2612
pending_ctx_cleanup_list = NULL;
2613
pending_uuid_deletion = NULL;
2614
+ pending_sql_statement = NULL;
2615
2616
if (unlikely(cmd.completion))
2617
cmd.completion = NULL; // Do not complete after launching worker (worker will do)
@@ -2581,6 +2623,7 @@ static void metadata_event_loop(void *arg)
2623
pending_ae_list = data->pending_alert_list;
2624
pending_ctx_cleanup_list = data->pending_ctx_cleanup_list;
2625
pending_uuid_deletion = data->pending_uuid_deletion;
2626
+ pending_sql_statement = data->pending_sql_statement;
2627
freez(data);
2628
metadata_flag_clear(wc, METADATA_FLAG_PROCESSING);
2629
}
@@ -2614,6 +2657,15 @@ static void metadata_event_loop(void *arg)
2657
case METADATA_DEL_HOST_AE:
2658
(void) JudyLIns(&wc->ae_DelJudyL, (Word_t) (void *) cmd.param[0], PJE0);
2659
break;
2660
+ case METADATA_EXECUTE_STORE_STATEMENT:
2661
+ stmt = (sqlite3_stmt *) cmd.param[0];
2662
+ if (!pending_sql_statement)
2663
+ pending_sql_statement = callocz(1, sizeof(*pending_sql_statement));
2664
+
2665
+ Pvalue = JudyLIns(&pending_sql_statement->JudyL, ++pending_sql_statement->count, PJE0);
2666
+ if (Pvalue)
2667
+ *Pvalue = (void *)stmt;
2668
+ break;
2669
case METADATA_UNITTEST:;
2670
struct thread_unittest *tu = (struct thread_unittest *) cmd.param[0];
2671
sleep_usec(1000); // processing takes 1ms
@@ -2847,6 +2899,13 @@ void metadata_queue_ae_deletion(ALARM_ENTRY *ae)
2899
queue_metadata_cmd(METADATA_DEL_HOST_AE, ae, NULL);
2900
}
2901
2902
+void metadata_execute_store_statement(sqlite3_stmt *stmt)
2903
+{
2904
+ if (unlikely(!metasync_worker.loop))
2905
+ return;
2906
+ queue_metadata_cmd(METADATA_EXECUTE_STORE_STATEMENT, stmt, NULL);
2907
+}
2908
+
2909
void commit_alert_transitions(RRDHOST *host __maybe_unused)
2910
{
2911
if (unlikely(!metasync_worker.loop))
src/database/sqlite/sqlite_metadata.h
+1
@@ -61,6 +61,7 @@ void metadata_sync_shutdown_background(void);
61
void metadata_sync_shutdown_background_wait(void);
62
void metadata_queue_ctx_host_cleanup(nd_uuid_t *host_uuid, const char *context);
63
void store_host_info_and_metadata(RRDHOST *host, BUFFER *work_buffer, size_t *query_counter);
64
+void metadata_execute_store_statement(sqlite3_stmt *stmt);
65
66
// UNIT TEST
67
int metadata_unittest(void);