@cryptotaxi247 / netdata-1 / commits / e79fc0224

Improve metadata event loop shutdown (#20132)

* Metadata shutdown in one step (not need to prepare and then shutdown) Terminate running jobs as soon as possible Terminate event loop and run a final metadata check Record progress during shutdown Anticipate early shutdown before ctx load is complete Allow 10 seconds for shutdown * Handle additional error cases / cleanup * Detect failure to load contexts in the background Detect Judy array insert failure and try to recover if possible * Fix counters for pending transitions to save * Handle failure to async save the alert entry

Stelios Fragkakis committed Apr 22, 2025 at 18:16 UTC e79fc02244f1eab5dbc42bb37a2b07088ca50e75
7 files changed +328 -387
src/daemon/daemon-shutdown-watcher.c
-2
@@ -138,7 +138,6 @@ void *watcher_main(void *arg)
138 watcher_wait_for_step(WATCHER_STEP_ID_STOP_ACLK_MQTT_THREAD, shutdown_start_time);
139 watcher_wait_for_step(WATCHER_STEP_ID_STOP_ALL_REMAINING_WORKER_THREADS, shutdown_start_time);
140 watcher_wait_for_step(WATCHER_STEP_ID_CANCEL_MAIN_THREADS, shutdown_start_time);
141 - watcher_wait_for_step(WATCHER_STEP_ID_PREPARE_METASYNC_SHUTDOWN, shutdown_start_time);
141 watcher_wait_for_step(WATCHER_STEP_ID_STOP_COLLECTION_FOR_ALL_HOSTS, shutdown_start_time);
142 watcher_wait_for_step(WATCHER_STEP_ID_WAIT_FOR_DBENGINE_COLLECTORS_TO_FINISH, shutdown_start_time);
143 watcher_wait_for_step(WATCHER_STEP_ID_STOP_DBENGINE_TIERS, shutdown_start_time);
@@ -171,7 +170,6 @@ void watcher_thread_start() {
170 "stop exporters, health and web servers threads";
171 watcher_steps[WATCHER_STEP_ID_STOP_COLLECTORS_AND_STREAMING_THREADS].msg = "stop collectors and streaming threads";
172 watcher_steps[WATCHER_STEP_ID_STOP_REPLICATION_THREADS].msg = "stop replication threads";
174 - watcher_steps[WATCHER_STEP_ID_PREPARE_METASYNC_SHUTDOWN].msg = "prepare metasync shutdown";
173 watcher_steps[WATCHER_STEP_ID_DISABLE_ML_DETEC_AND_TRAIN_THREADS].msg = "disable ML detection and training threads";
174 watcher_steps[WATCHER_STEP_ID_STOP_CONTEXT_THREAD].msg = "stop context thread";
175 watcher_steps[WATCHER_STEP_ID_CLEAR_WEB_CLIENT_CACHE].msg = "clear web client cache";
src/daemon/daemon-shutdown-watcher.h
-1
@@ -19,7 +19,6 @@ typedef enum {
19 WATCHER_STEP_ID_STOP_ACLK_MQTT_THREAD,
20 WATCHER_STEP_ID_STOP_ALL_REMAINING_WORKER_THREADS,
21 WATCHER_STEP_ID_CANCEL_MAIN_THREADS,
22 - WATCHER_STEP_ID_PREPARE_METASYNC_SHUTDOWN,
22 WATCHER_STEP_ID_STOP_COLLECTION_FOR_ALL_HOSTS,
23 WATCHER_STEP_ID_WAIT_FOR_DBENGINE_COLLECTORS_TO_FINISH,
24 WATCHER_STEP_ID_STOP_DBENGINE_TIERS,
src/daemon/daemon-shutdown.c
+1 -4
@@ -265,9 +265,6 @@ static void netdata_cleanup_and_exit(EXIT_REASON reason, bool abnormal, bool exi
265 cancel_main_threads();
266 watcher_step_complete(WATCHER_STEP_ID_CANCEL_MAIN_THREADS);
267
268 - metadata_sync_shutdown_background();
269 - watcher_step_complete(WATCHER_STEP_ID_PREPARE_METASYNC_SHUTDOWN);
270 -
268 if (abnormal) {
269 watcher_step_complete(WATCHER_STEP_ID_STOP_COLLECTION_FOR_ALL_HOSTS);
270 watcher_step_complete(WATCHER_STEP_ID_WAIT_FOR_DBENGINE_COLLECTORS_TO_FINISH);
@@ -310,7 +307,7 @@ static void netdata_cleanup_and_exit(EXIT_REASON reason, bool abnormal, bool exi
307 watcher_step_complete(WATCHER_STEP_ID_STOP_DBENGINE_TIERS);
308 #endif
309
313 - metadata_sync_shutdown_background_wait();
310 + metadata_sync_shutdown();
311 watcher_step_complete(WATCHER_STEP_ID_STOP_METASYNC_THREADS);
312 }
313
src/database/sqlite/sqlite_aclk.c
+5 -1
@@ -1034,7 +1034,11 @@ void aclk_synchronization_init(void)
1034
1035 nd_log_daemon(NDLP_INFO, "Created %d archived hosts", number_of_children);
1036 // Trigger host context load for hosts that have been created
1037 - metadata_queue_load_host_context();
1037 + if (unlikely(!metadata_queue_load_host_context())) {
1038 + nd_log_daemon(NDLP_WARNING, "Failed to queue command to load contexts for archived hosts");
1039 + // Reset context load flag so that contexts will be loaded on demand
1040 + reset_host_context_load_flag();
1041 + }
1042
1043 rc = sqlite3_exec_monitored(db_meta, SQL_FETCH_ALL_INSTANCES, aclk_config_parameters, NULL, &err_msg);
1044
src/database/sqlite/sqlite_metadata.c
+312 -371
@@ -13,6 +13,7 @@
13 duration_snprintf(var_name, sizeof(var_name), \
14 (int64_t)((end) - (start)), unit, true)
15
16 +#define SHUTDOWN_REQUESTED(config) (__atomic_load_n(&(config)->shutdown_requested, __ATOMIC_RELAXED))
17
18 extern long long def_journal_size_limit;
19
@@ -189,7 +190,7 @@ sqlite3 *db_meta = NULL;
190
191 #define METADATA_HOST_CHECK_FIRST_CHECK (5) // First check for pending metadata
192 #define METADATA_HOST_CHECK_INTERVAL (5) // Repeat check for pending metadata
192 -#define METADATA_MAX_BATCH_SIZE (512) // Maximum commands to execute before running the event loop
193 +#define METADATA_MAX_BATCH_SIZE (64) // Maximum commands to execute before running the event loop
194
195 #define DATABASE_VACUUM_FREQUENCY_SECONDS (60)
196 #define DATABASE_FREE_PAGES_THRESHOLD_PC (5) // Percentage of free pages to trigger vacuum
@@ -199,7 +200,7 @@ enum metadata_opcode {
200 METADATA_DATABASE_NOOP = 0,
201 METADATA_DEL_DIMENSION,
202 METADATA_STORE_CLAIM_ID,
202 - METADATA_SCAN_HOSTS,
203 + METADATA_STORE,
204 METADATA_LOAD_HOST_CONTEXT,
205 METADATA_DELETE_HOST_CHART_LABELS,
206 METADATA_ADD_HOST_AE,
@@ -217,36 +218,27 @@ enum metadata_opcode {
218 #define MAX_PARAM_LIST (2)
219 struct metadata_cmd {
220 enum metadata_opcode opcode;
220 - struct completion *completion;
221 const void *param[MAX_PARAM_LIST];
222 struct metadata_cmd *prev, *next;
223 };
224
225 -typedef enum {
226 - METADATA_FLAG_PROCESSING = (1 << 0), // store or cleanup
227 - METADATA_FLAG_SHUTDOWN = (1 << 1), // Shutting down
228 -} METADATA_FLAG;
229 -
230 -struct metadata_wc {
225 +struct meta_config_s {
226 uv_thread_t thread;
227 uv_loop_t loop;
228 uv_async_t async;
229 uv_timer_t timer_req;
230 time_t metadata_check_after;
231 Pvoid_t ae_DelJudyL;
237 - METADATA_FLAG flags;
232 bool initialized;
233 + bool ctx_load_running;
234 + bool metadata_running;
235 + bool shutdown_requested;
236 SPINLOCK cmd_queue_lock;
237 struct completion start_stop_complete;
241 - struct completion *scan_complete;
238 /* FIFO command queue */
239 struct metadata_cmd *cmd_base;
240 ARAL *ar;
245 -} metasync_worker;
246 -
247 -#define metadata_flag_check(target_flags, flag) (__atomic_load_n(&((target_flags)->flags), __ATOMIC_SEQ_CST) & (flag))
248 -#define metadata_flag_set(target_flags, flag) __atomic_or_fetch(&((target_flags)->flags), (flag), __ATOMIC_SEQ_CST)
249 -#define metadata_flag_clear(target_flags, flag) __atomic_and_fetch(&((target_flags)->flags), ~(flag), __ATOMIC_SEQ_CST)
241 +} meta_config;
242
243 //
244 // For unittest
@@ -262,9 +254,6 @@ int sql_metadata_cache_stats(int op)
254 {
255 int count, dummy;
256
265 - if (!REQUIRE_DB(db_meta))
266 - return 0;
267 -
257 sqlite3_db_status(db_meta, op, &count, &dummy, 0);
258 return count;
259 }
@@ -367,7 +356,7 @@ done:
356 }
357
358 #define SQL_SCHEDULE_HOST_CTX_CLEANUP \
370 - "INSERT INTO ctx_metadata_cleanup (host_id, context, date_created) " \
359 + "INSERT INTO ctx_metadata_cleanup (host_id, context, date_created) " \
360 "VALUES (@host_id, @context, UNIXEPOCH()) ON CONFLICT DO UPDATE SET date_created = excluded.date_created; END"
361
362 // Schedule context cleanup for host
@@ -414,7 +403,7 @@ bool sql_set_host_label(nd_uuid_t *host_id, const char *label_key, const char *l
403 SQLITE_BIND_FAIL(done, sqlite3_bind_text(res, ++param, label_value, -1, SQLITE_STATIC));
404
405 param = 0;
417 - int rc = execute_insert(res);
406 + int rc = sqlite3_step_monitored(res);
407 status = (rc == SQLITE_DONE);
408 if (false == status)
409 error_report("Failed to store node instance information, rc = %d", rc);
@@ -439,9 +428,6 @@ void sql_update_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id)
428 set_host_node_id(host, node_id);
429 rrd_wrunlock();
430
442 - if (!REQUIRE_DB(db_meta))
443 - return;
444 -
431 if (!PREPARE_STATEMENT(db_meta, SQL_UPDATE_NODE_ID, &res))
432 return;
433
@@ -459,32 +445,6 @@ done:
445 SQLITE_FINALIZE(res);
446 }
447
462 -#define SQL_SELECT_NODE_ID "SELECT node_id FROM node_instance WHERE host_id = @host_id AND node_id IS NOT NULL"
463 -
464 -int get_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id)
465 -{
466 - sqlite3_stmt *res = NULL;
467 -
468 - if (!REQUIRE_DB(db_meta))
469 - return 1;
470 -
471 - if (!PREPARE_STATEMENT(db_meta, SQL_SELECT_NODE_ID, &res))
472 - return 1;
473 -
474 - int param = 0, rc = 0;
475 - SQLITE_BIND_FAIL(done, sqlite3_bind_blob(res, ++param, host_id, sizeof(*host_id), SQLITE_STATIC));
476 -
477 - param = 0;
478 - rc = sqlite3_step_monitored(res);
479 - if (likely(rc == SQLITE_ROW && node_id))
480 - uuid_copy(*node_id, *((nd_uuid_t *) sqlite3_column_blob(res, 0)));
481 -
482 -done:
483 - REPORT_BIND_FAIL(res, param);
484 - SQLITE_FINALIZE(res);
485 - return (rc == SQLITE_ROW) ? 0 : -1;
486 -}
487 -
448 #define SQL_INVALIDATE_NODE_INSTANCES \
449 "UPDATE node_instance SET node_id = NULL WHERE EXISTS " \
450 "(SELECT host_id FROM node_instance WHERE host_id = @host_id AND (@claim_id IS NULL OR claim_id <> @claim_id))"
@@ -493,9 +453,6 @@ void invalidate_node_instances(nd_uuid_t *host_id, nd_uuid_t *claim_id)
453 {
454 sqlite3_stmt *res = NULL;
455
496 - if (!REQUIRE_DB(db_meta))
497 - return;
498 -
456 if (!PREPARE_STATEMENT(db_meta, SQL_INVALIDATE_NODE_INSTANCES, &res))
457 return;
458
@@ -508,7 +465,7 @@ void invalidate_node_instances(nd_uuid_t *host_id, nd_uuid_t *claim_id)
465 SQLITE_BIND_FAIL(done, sqlite3_bind_null(res, ++param));
466
467 param = 0;
511 - int rc = execute_insert(res);
468 + int rc = sqlite3_step_monitored(res);
469 if (unlikely(rc != SQLITE_DONE))
470 error_report("Failed to invalidate node instance information, rc = %d", rc);
471
@@ -523,9 +480,6 @@ void sql_load_node_id(RRDHOST *host)
480 {
481 sqlite3_stmt *res = NULL;
482
526 - if (!REQUIRE_DB(db_meta))
527 - return;
528 -
483 if (!PREPARE_STATEMENT(db_meta, SQL_GET_HOST_NODE_ID, &res))
484 return;
485
@@ -614,7 +568,7 @@ static int exec_statement_with_uuid(const char *sql, nd_uuid_t *uuid)
568 SQLITE_BIND_FAIL(done, sqlite3_bind_blob(res, ++param, uuid, sizeof(*uuid), SQLITE_STATIC));
569
570 param = 0;
617 - int rc = execute_insert(res);
571 + int rc = sqlite3_step_monitored(res);
572 if (likely(rc == SQLITE_DONE))
573 result = SQLITE_OK;
574 else
@@ -918,9 +872,6 @@ static int store_claim_id(nd_uuid_t *host_id, nd_uuid_t *claim_id)
872 sqlite3_stmt *res = NULL;
873 int rc = 0;
874
921 - if (!REQUIRE_DB(db_meta))
922 - return 1;
923 -
875 if (!PREPARE_STATEMENT(db_meta, SQL_STORE_CLAIM_ID, &res))
876 return 1;
877
@@ -1036,9 +987,6 @@ static int add_host_sysinfo_key_value(const char *name, const char *value, nd_uu
987 {
988 sqlite3_stmt *res = NULL;
989
1039 - if (!REQUIRE_DB(db_meta))
1040 - return 0;
1041 -
990 if (!PREPARE_STATEMENT(db_meta, SQL_STORE_HOST_SYSTEM_INFO_VALUES, &res))
991 return 0;
992
@@ -1180,7 +1128,7 @@ static bool dimension_can_be_deleted(nd_uuid_t *dim_uuid __maybe_unused, sqlite3
1128
1129 static bool run_cleanup_loop(
1130 sqlite3_stmt *res,
1183 - struct metadata_wc *wc,
1131 + struct meta_config_s *config,
1132 bool (*check_cb)(nd_uuid_t *, sqlite3_stmt **, bool),
1133 void (*action_cb)(nd_uuid_t *, sqlite3_stmt **, bool),
1134 uint32_t *total_checked,
@@ -1191,7 +1139,7 @@ static bool run_cleanup_loop(
1139 bool check_flag,
1140 bool action_flag)
1141 {
1194 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)))
1142 + if (unlikely(SHUTDOWN_REQUESTED(config)))
1143 return true;
1144
1145 int rc = sqlite3_bind_int64(res, 1, (sqlite3_int64) *row_id);
@@ -1204,7 +1152,7 @@ static bool run_cleanup_loop(
1152 uint32_t l_checked = 0;
1153 uint32_t l_deleted = 0;
1154 while (!time_expired && sqlite3_step_monitored(res) == SQLITE_ROW) {
1207 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)))
1155 + if (unlikely(SHUTDOWN_REQUESTED(config)))
1156 break;
1157
1158 *row_id = sqlite3_column_int64(res, 1);
@@ -1317,7 +1265,7 @@ static uint64_t get_rowid_from_statement(const char *sql)
1265
1266 #define SQL_GET_MAX_DIM_ROW_ID "SELECT MAX(rowid) FROM dimension"
1267
1320 -static bool check_dimension_metadata(struct metadata_wc *wc)
1268 +static bool check_dimension_metadata(struct meta_config_s *wc)
1269 {
1270 static time_t next_execution_t = 0;
1271 static uint64_t last_row_id = 0;
@@ -1383,7 +1331,7 @@ static bool check_dimension_metadata(struct metadata_wc *wc)
1331
1332 #define SQL_GET_MAX_CHART_ROW_ID "SELECT MAX(rowid) FROM chart"
1333
1386 -static bool check_chart_metadata(struct metadata_wc *wc)
1334 +static bool check_chart_metadata(struct meta_config_s *wc)
1335 {
1336 static time_t next_execution_t = 0;
1337 static uint64_t last_row_id = 0;
@@ -1455,7 +1403,7 @@ static bool check_chart_metadata(struct metadata_wc *wc)
1403
1404 #define SQL_GET_MAX_CHART_LABEL_ROW_ID "SELECT MAX(rowid) FROM chart_label"
1405
1458 -static bool check_label_metadata(struct metadata_wc *wc)
1406 +static bool check_label_metadata(struct meta_config_s *wc)
1407 {
1408 static time_t next_execution_t = 0;
1409 static uint64_t last_row_id = 0;
@@ -1529,7 +1477,7 @@ static bool check_label_metadata(struct metadata_wc *wc)
1477 return false;
1478 }
1479
1532 -static void cleanup_health_log(struct metadata_wc *wc)
1480 +static void cleanup_health_log(struct meta_config_s *config)
1481 {
1482 static time_t next_execution_t = 0;
1483
@@ -1549,12 +1497,12 @@ static void cleanup_health_log(struct metadata_wc *wc)
1497 dfe_start_reentrant(rrdhost_root_index, host)
1498 {
1499 sql_health_alarm_log_cleanup(host);
1552 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)))
1500 + if (unlikely(SHUTDOWN_REQUESTED(config)))
1501 break;
1502 }
1503 dfe_done(host);
1504
1557 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN))) {
1505 + if (unlikely(SHUTDOWN_REQUESTED(config))) {
1506 worker_is_idle();
1507 return;
1508 }
@@ -1569,7 +1517,7 @@ static void cleanup_health_log(struct metadata_wc *wc)
1517 // EVENT LOOP STARTS HERE
1518 //
1519
1572 -static void metadata_free_cmd_queue(struct metadata_wc *wc)
1520 +static void metadata_free_cmd_queue(struct meta_config_s *wc)
1521 {
1522 spinlock_lock(&wc->cmd_queue_lock);
1523 while(wc->cmd_base) {
@@ -1580,32 +1528,24 @@ static void metadata_free_cmd_queue(struct metadata_wc *wc)
1528 spinlock_unlock(&wc->cmd_queue_lock);
1529 }
1530
1583 -static void metadata_enq_cmd(struct metadata_wc *wc, struct metadata_cmd *cmd)
1531 +static bool metadata_enq_cmd(struct metadata_cmd *cmd)
1532 {
1585 - if(unlikely(!wc->initialized))
1586 - return;
1587 -
1588 - if (cmd->opcode == METADATA_SYNC_SHUTDOWN) {
1589 - metadata_flag_set(wc, METADATA_FLAG_SHUTDOWN);
1590 - goto wakeup_event_loop;
1591 - }
1592 -
1593 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)))
1594 - goto wakeup_event_loop;
1533 + if(unlikely(!__atomic_load_n(&meta_config.initialized, __ATOMIC_RELAXED)))
1534 + return false;
1535
1596 - struct metadata_cmd *t = aral_mallocz(wc->ar);
1536 + struct metadata_cmd *t = aral_mallocz(meta_config.ar);
1537 *t = *cmd;
1538 t->prev = t->next = NULL;
1539
1600 - spinlock_lock(&wc->cmd_queue_lock);
1601 - DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(wc->cmd_base, t, prev, next);
1602 - spinlock_unlock(&wc->cmd_queue_lock);
1540 + spinlock_lock(&meta_config.cmd_queue_lock);
1541 + DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(meta_config.cmd_base, t, prev, next);
1542 + spinlock_unlock(&meta_config.cmd_queue_lock);
1543
1604 -wakeup_event_loop:
1605 - (void) uv_async_send(&wc->async);
1544 + (void) uv_async_send(&meta_config.async);
1545 + return true;
1546 }
1547
1608 -static struct metadata_cmd metadata_deq_cmd(struct metadata_wc *wc)
1548 +static struct metadata_cmd metadata_deq_cmd(struct meta_config_s *wc)
1549 {
1550 struct metadata_cmd ret, *to_free = NULL;
1551
@@ -1616,10 +1556,8 @@ static struct metadata_cmd metadata_deq_cmd(struct metadata_wc *wc)
1556 ret = *t;
1557 to_free = t;
1558 }
1619 - else {
1559 + else
1560 ret.opcode = METADATA_DATABASE_NOOP;
1621 - ret.completion = NULL;
1622 - }
1561 spinlock_unlock(&wc->cmd_queue_lock);
1562
1563 aral_freez(wc->ar, to_free);
@@ -1641,13 +1579,13 @@ static void timer_cb(uv_timer_t* handle)
1579 uv_stop(handle->loop);
1580 uv_update_time(handle->loop);
1581
1644 - struct metadata_wc *wc = handle->data;
1582 + struct meta_config_s *wc = handle->data;
1583
1584 if (wc->metadata_check_after < now_realtime_sec()) {
1585 struct metadata_cmd cmd;
1586 memset(&cmd, 0, sizeof(cmd));
1649 - cmd.opcode = METADATA_SCAN_HOSTS;
1650 - metadata_enq_cmd(wc, &cmd);
1587 + cmd.opcode = METADATA_STORE;
1588 + (void) metadata_enq_cmd(&cmd);
1589 }
1590 }
1591
@@ -1685,7 +1623,7 @@ void vacuum_database(sqlite3 *database, const char *db_alias, int threshold, int
1623
1624 static bool clean_host_chart_dimensions(sqlite3_stmt **res, int64_t chart_row_id, size_t *checked, size_t *deleted)
1625 {
1688 - struct metadata_wc *wc = &metasync_worker;
1626 + struct meta_config_s *config = &meta_config;
1627
1628 bool can_continue = false;
1629
@@ -1712,7 +1650,7 @@ static bool clean_host_chart_dimensions(sqlite3_stmt **res, int64_t chart_row_id
1650 (*deleted)++;
1651 }
1652 (*checked)++;
1715 - can_continue = (!metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)) && sql_metadata_wal_size_acceptable();
1653 + can_continue = (!SHUTDOWN_REQUESTED(config)) && sql_metadata_wal_size_acceptable();
1654 }
1655 SQLITE_FINALIZE(dim_del_stmt);
1656
@@ -1729,7 +1667,7 @@ static void cleanup_host_context_metadata(Pvoid_t CTX_JudyL, void *data)
1667 if (!CTX_JudyL || !data)
1668 return;
1669
1732 - struct metadata_wc *wc = &metasync_worker;
1670 + struct meta_config_s *config = &meta_config;
1671
1672 RRDHOST *host = data;
1673
@@ -1765,8 +1703,7 @@ static void cleanup_host_context_metadata(Pvoid_t CTX_JudyL, void *data)
1703 ctx_delete_metadata_cleanup_context(&context_res, &host->host_id.uuid, context);
1704 }
1705 string_freez(ctx);
1768 - can_continue =
1769 - can_continue && (!metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)) && sql_metadata_wal_size_acceptable();
1706 + can_continue = can_continue && (!SHUTDOWN_REQUESTED(config)) && sql_metadata_wal_size_acceptable();
1707 }
1708 SQLITE_FINALIZE(dimension_res);
1709 SQLITE_FINALIZE(context_res);
@@ -1783,7 +1720,7 @@ done:
1720 SQLITE_FINALIZE(res);
1721 }
1722
1786 -void run_metadata_cleanup(struct metadata_wc *wc)
1723 +void run_metadata_cleanup(struct meta_config_s *config)
1724 {
1725 static time_t next_context_list_cleanup = 0;
1726
@@ -1792,15 +1729,12 @@ void run_metadata_cleanup(struct metadata_wc *wc)
1729 if (!next_context_list_cleanup)
1730 next_context_list_cleanup = now + 5;
1731
1795 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)))
1796 - return;
1797 -
1732 if (next_context_list_cleanup < now && sql_metadata_wal_size_acceptable()) {
1733 RRDHOST *host;
1734 worker_is_busy(UV_EVENT_CTX_CLEANUP);
1735 dfe_start_reentrant(rrdhost_root_index, host) {
1736 ctx_get_context_list_to_cleanup(&host->host_id.uuid, cleanup_host_context_metadata, host);
1803 - if (metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN) || false == sql_metadata_wal_size_acceptable())
1737 + if (SHUTDOWN_REQUESTED(config) || false == sql_metadata_wal_size_acceptable())
1738 break;
1739 }
1740 dfe_done(host);
@@ -1808,26 +1742,26 @@ void run_metadata_cleanup(struct metadata_wc *wc)
1742 next_context_list_cleanup = now_realtime_sec() + METADATA_MAINTENANCE_CTX_CLEAN_REPEAT;
1743 }
1744
1811 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)))
1745 + if (unlikely(SHUTDOWN_REQUESTED(config)))
1746 return;
1747
1814 - if (check_dimension_metadata(wc))
1815 - if (check_chart_metadata(wc))
1816 - check_label_metadata(wc);
1748 + if (check_dimension_metadata(config))
1749 + if (check_chart_metadata(config))
1750 + check_label_metadata(config);
1751
1818 - cleanup_health_log(wc);
1752 + cleanup_health_log(config);
1753
1820 - if (unlikely(metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN)))
1821 - return;
1754 + if (unlikely(SHUTDOWN_REQUESTED(config)))
1755 + return;
1756
1757 vacuum_database(db_meta, "METADATA", DATABASE_FREE_PAGES_THRESHOLD_PC, DATABASE_FREE_PAGES_VACUUM_PC);
1758
1759 (void) sqlite3_wal_checkpoint(db_meta, NULL);
1760 }
1761
1828 -struct scan_metadata_payload {
1762 +struct work_payload {
1763 uv_work_t request;
1830 - struct metadata_wc *wc;
1764 + struct meta_config_s *config;
1765 void *pending_alert_list;
1766 void *pending_ctx_cleanup_list;
1767 void *pending_uuid_deletion;
@@ -1904,9 +1838,11 @@ static void restore_host_context(void *arg)
1838 }
1839
1840 // Callback after scan of hosts is done
1907 -static void after_start_host_load_context(uv_work_t *req, int status __maybe_unused)
1841 +static void after_ctx_hosts_load(uv_work_t *req, int status __maybe_unused)
1842 {
1909 - struct scan_metadata_payload *data = req->data;
1843 + struct work_payload *data = req->data;
1844 + struct meta_config_s *config = data->config;
1845 + config->ctx_load_running = false;
1846 freez(data);
1847 }
1848
@@ -1947,12 +1883,22 @@ static bool cleanup_finished_threads(struct host_context_load_thread *hclt, size
1883 return found_slot || wait;
1884 }
1885
1950 -static void start_all_host_load_context(uv_work_t *req __maybe_unused)
1886 +void reset_host_context_load_flag()
1887 +{
1888 + RRDHOST *host;
1889 + dfe_start_reentrant(rrdhost_root_index, host)
1890 + {
1891 + rrdhost_flag_set(host, RRDHOST_FLAG_PENDING_CONTEXT_LOAD);
1892 + }
1893 + dfe_done(host);
1894 +}
1895 +
1896 +static void ctx_hosts_load(uv_work_t *req)
1897 {
1898 register_libuv_worker_jobs();
1899
1954 - struct scan_metadata_payload *data = req->data;
1955 - struct metadata_wc *wc = data->wc;
1900 + struct work_payload *data = req->data;
1901 + struct meta_config_s *config = data->config;
1902
1903 worker_is_busy(UV_EVENT_HOST_CONTEXT_LOAD);
1904 usec_t started_ut = now_monotonic_usec(); (void)started_ut;
@@ -1975,7 +1921,7 @@ static void start_all_host_load_context(uv_work_t *req __maybe_unused)
1921 if (!rrdhost_flag_check(host, RRDHOST_FLAG_PENDING_CONTEXT_LOAD))
1922 continue;
1923
1978 - if (metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN))
1924 + if (unlikely(SHUTDOWN_REQUESTED(config)))
1925 break;
1926
1927 nd_log_daemon(NDLP_DEBUG, "Loading context for host %s", rrdhost_hostname(host));
@@ -2033,34 +1979,29 @@ static void start_all_host_load_context(uv_work_t *req __maybe_unused)
1979 db_meta_thread = NULL;
1980 db_context_thread = NULL;
1981 }
2036 -
1982 worker_is_idle();
1983 }
1984
1985 // Callback after scan of hosts is done
1986 static void after_metadata_hosts(uv_work_t *req, int status __maybe_unused)
1987 {
2043 - struct scan_metadata_payload *data = req->data;
2044 - struct metadata_wc *wc = data->wc;
1988 + struct work_payload *data = req->data;
1989 + struct meta_config_s *config = data->config;
1990
1991 bool first = true;
1992 Word_t Index = 0;
1993 Pvoid_t *Pvalue;
2049 - while ((Pvalue = JudyLFirstThenNext(wc->ae_DelJudyL, &Index, &first))) {
1994 + while ((Pvalue = JudyLFirstThenNext(config->ae_DelJudyL, &Index, &first))) {
1995 ALARM_ENTRY *ae = (ALARM_ENTRY *) Index;
1996 if(!__atomic_load_n(&ae->pending_save_count, __ATOMIC_RELAXED)) {
1997 health_alarm_log_free_one_nochecks_nounlink(ae);
2053 - (void) JudyLDel(&wc->ae_DelJudyL, Index, PJE0);
1998 + (void) JudyLDel(&config->ae_DelJudyL, Index, PJE0);
1999 first = true;
2000 Index = 0;
2001 }
2002 }
2003
2059 - metadata_flag_clear(wc, METADATA_FLAG_PROCESSING);
2060 -
2061 - if (unlikely(wc->scan_complete))
2062 - completion_mark_complete(wc->scan_complete);
2063 -
2004 + config->metadata_running = false;
2005 freez(data);
2006 }
2007
@@ -2110,7 +2051,7 @@ size_t populate_metrics_from_database(void *mrg, void (*populate_cb)(void *mrg,
2051 return count;
2052 }
2053
2113 -static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool shutting_down)
2054 +static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool is_worker)
2055 {
2056 static bool skip_models = false;
2057 RRDSET *st;
@@ -2120,6 +2061,8 @@ static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool shutting
2061 sqlite3_stmt *store_dimension = NULL;
2062 sqlite3_stmt *store_chart = NULL;
2063
2064 + bool load_ml_models = is_worker;
2065 +
2066 bool host_need_recheck = false;
2067 (void)db_execute(db_meta, "BEGIN TRANSACTION");
2068
@@ -2131,7 +2074,9 @@ static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool shutting
2074
2075 buffer_flush(work_buffer);
2076
2134 - worker_is_busy(UV_EVENT_STORE_CHART);
2077 + if (is_worker)
2078 + worker_is_busy(UV_EVENT_STORE_CHART);
2079 +
2080 rc = check_and_update_chart_labels(st, work_buffer);
2081 if (unlikely(rc))
2082 error_report("METADATA: 'host:%s': Failed to update labels for chart %s", rrdhost_hostname(host), rrdset_name(st));
@@ -2145,18 +2090,24 @@ static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool shutting
2090 rrdhost_hostname(host),
2091 rrdset_name(st));
2092 }
2148 - worker_is_idle();
2093 + if (is_worker)
2094 + worker_is_idle();
2095 }
2096
2097 RRDDIM *rd;
2098 rrddim_foreach_read(rd, st) {
2099 + if (load_ml_models) {
2100 + if (rrddim_flag_check(rd, RRDDIM_FLAG_ML_MODEL_LOAD)) {
2101 + rrddim_flag_clear(rd, RRDDIM_FLAG_ML_MODEL_LOAD);
2102 + if (likely(!skip_models)) {
2103 + if (is_worker)
2104 + worker_is_busy(UV_EVENT_METADATA_ML_LOAD);
2105 +
2106 + skip_models = ml_dimension_load_models(rd, &ml_load_stmt);
2107
2154 - if (rrddim_flag_check(rd, RRDDIM_FLAG_ML_MODEL_LOAD)) {
2155 - rrddim_flag_clear(rd, RRDDIM_FLAG_ML_MODEL_LOAD);
2156 - if (likely(!skip_models && !shutting_down)) {
2157 - worker_is_busy(UV_EVENT_METADATA_ML_LOAD);
2158 - skip_models = ml_dimension_load_models(rd, &ml_load_stmt);
2159 - worker_is_idle();
2108 + if (is_worker)
2109 + worker_is_idle();
2110 + }
2111 }
2112 }
2113
@@ -2170,7 +2121,9 @@ static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool shutting
2121 else
2122 rrddim_flag_clear(rd, RRDDIM_FLAG_META_HIDDEN);
2123
2173 - worker_is_busy(UV_EVENT_STORE_DIMENSION);
2124 + if (is_worker)
2125 + worker_is_busy(UV_EVENT_STORE_DIMENSION);
2126 +
2127 rc = store_dimension_metadata(rd, &store_dimension);
2128 if (unlikely(rc)) {
2129 host_need_recheck = true;
@@ -2182,7 +2135,8 @@ static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool shutting
2135 rrddim_name(rd));
2136 }
2137
2185 - worker_is_idle();
2138 + if (is_worker)
2139 + worker_is_idle();
2140 }
2141 rrddim_foreach_done(rd);
2142 }
@@ -2220,7 +2174,7 @@ struct judy_list_t {
2174 Word_t count;
2175 };
2176
2223 -static void do_pending_uuid_deletion(struct metadata_wc *wc, struct judy_list_t *pending_uuid_deletion)
2177 +static void do_pending_uuid_deletion(struct meta_config_s *config, struct judy_list_t *pending_uuid_deletion)
2178 {
2179 if (!pending_uuid_deletion)
2180 return;
@@ -2237,12 +2191,11 @@ static void do_pending_uuid_deletion(struct metadata_wc *wc, struct judy_list_t
2191 if (!*Pvalue)
2192 continue;
2193
2240 - if (metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN))
2241 - break;
2242 -
2194 nd_uuid_t *uuid = *Pvalue;
2244 - if (dimension_can_be_deleted(uuid, NULL, false))
2245 - delete_dimension_uuid(uuid, NULL, false);
2195 + if (likely(!SHUTDOWN_REQUESTED(config))) {
2196 + if (dimension_can_be_deleted(uuid, NULL, false))
2197 + delete_dimension_uuid(uuid, NULL, false);
2198 + }
2199
2200 freez(uuid);
2201 }
@@ -2259,7 +2212,7 @@ static void do_pending_uuid_deletion(struct metadata_wc *wc, struct judy_list_t
2212 worker_is_idle();
2213 }
2214
2262 -static void store_ctx_cleanup_list(struct metadata_wc *wc, struct judy_list_t *pending_ctx_cleanup_list)
2215 +static void store_ctx_cleanup_list(struct meta_config_s *config, struct judy_list_t *pending_ctx_cleanup_list)
2216 {
2217 if (!pending_ctx_cleanup_list)
2218 return;
@@ -2277,11 +2230,11 @@ static void store_ctx_cleanup_list(struct metadata_wc *wc, struct judy_list_t *p
2230 if (!*Pvalue)
2231 continue;
2232
2280 - if (metadata_flag_check(wc, METADATA_FLAG_SHUTDOWN))
2281 - break;
2282 -
2233 struct host_ctx_cleanup_s *ctx_cleanup = *Pvalue;
2284 - sql_schedule_host_ctx_cleanup(&res, &ctx_cleanup->host_uuid, string2str(ctx_cleanup->context));
2234 +
2235 + if (likely(!SHUTDOWN_REQUESTED(config)))
2236 + sql_schedule_host_ctx_cleanup(&res, &ctx_cleanup->host_uuid, string2str(ctx_cleanup->context));
2237 +
2238 string_freez(ctx_cleanup->context);
2239 freez(ctx_cleanup);
2240 }
@@ -2299,12 +2252,13 @@ static void store_ctx_cleanup_list(struct metadata_wc *wc, struct judy_list_t *p
2252 worker_is_idle();
2253 }
2254
2302 -static void store_alert_transitions(struct judy_list_t *pending_alert_list)
2255 +static void store_alert_transitions(struct judy_list_t *pending_alert_list, bool is_worker)
2256 {
2257 if (!pending_alert_list)
2258 return;
2259
2307 - worker_is_busy(UV_EVENT_STORE_ALERT_TRANSITIONS);
2260 + if (is_worker)
2261 + worker_is_busy(UV_EVENT_STORE_ALERT_TRANSITIONS);
2262
2263 usec_t started_ut = now_monotonic_usec(); (void)started_ut;
2264
@@ -2334,15 +2288,31 @@ static void store_alert_transitions(struct judy_list_t *pending_alert_list)
2288 entries,
2289 (double)(ended_ut - started_ut) / USEC_PER_MS);
2290
2337 - worker_is_idle();
2291 + if (is_worker)
2292 + worker_is_idle();
2293 }
2294
2340 -static void store_sql_statements(struct judy_list_t *pending_sql_statement)
2295 +static int execute_statement(sqlite3_stmt *stmt)
2296 +{
2297 + if (!stmt)
2298 + return SQLITE_OK;
2299 +
2300 + int rc = sqlite3_step_monitored(stmt);
2301 + if (unlikely(rc != SQLITE_DONE))
2302 + nd_log_daemon(NDLP_ERR, "Failed to execute sql statement, rc = %d", rc);
2303 +
2304 + SQLITE_FINALIZE(stmt);
2305 +
2306 + return rc;
2307 +}
2308 +
2309 +static void store_sql_statements(struct judy_list_t *pending_sql_statement, bool is_worker)
2310 {
2311 if (!pending_sql_statement)
2312 return;
2313
2345 - worker_is_busy(METADATA_EXECUTE_STORE_STATEMENT);
2314 + if (is_worker)
2315 + worker_is_busy(METADATA_EXECUTE_STORE_STATEMENT);
2316
2317 usec_t started_ut = now_monotonic_usec();
2318
@@ -2352,15 +2322,7 @@ static void store_sql_statements(struct judy_list_t *pending_sql_statement)
2322 Pvoid_t *Pvalue;
2323 while ((Pvalue = JudyLFirstThenNext(pending_sql_statement->JudyL, &Index, &first))) {
2324 sqlite3_stmt *stmt = *Pvalue;
2355 -
2356 - if (unlikely(!stmt))
2357 - continue;
2358 -
2359 - int rc = sqlite3_step_monitored(stmt);
2360 - if (unlikely(rc != SQLITE_DONE))
2361 - nd_log_daemon(NDLP_ERR, "Failed to execute sql statement, rc = %d", rc);
2362 -
2363 - SQLITE_FINALIZE(stmt);
2325 + execute_statement(stmt);
2326 }
2327 (void) JudyLFreeArray(&pending_sql_statement->JudyL, PJE0);
2328 freez(pending_sql_statement);
@@ -2368,7 +2330,8 @@ static void store_sql_statements(struct judy_list_t *pending_sql_statement)
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
2371 - worker_is_idle();
2333 + if (is_worker)
2334 + worker_is_idle();
2335 }
2336
2337 static void meta_store_host_labels(RRDHOST *host, BUFFER *work_buffer)
@@ -2412,8 +2375,6 @@ static void store_host_claim_id(RRDHOST *host)
2375 rrdhost_flag_set(host, RRDHOST_FLAG_METADATA_CLAIMID | RRDHOST_FLAG_METADATA_UPDATE);
2376 }
2377
2415 -
2416 -
2378 void store_host_info_and_metadata(RRDHOST *host, BUFFER *work_buffer)
2379 {
2380 // Store labels (if needed)
@@ -2429,63 +2390,94 @@ void store_host_info_and_metadata(RRDHOST *host, BUFFER *work_buffer)
2390 store_host_and_system_info(host);
2391 }
2392
2432 -// Worker thread to scan hosts for pending metadata to store
2433 -static void start_metadata_hosts(uv_work_t *req)
2393 +static void store_hosts_metadata(BUFFER *work_buffer, bool is_worker)
2394 {
2435 - register_libuv_worker_jobs();
2436 -
2437 - struct scan_metadata_payload *data = req->data;
2438 - struct metadata_wc *wc = data->wc;
2439 -
2440 - bool shutting_down = (!wc->scan_complete);
2441 -
2442 - BUFFER *work_buffer = data->work_buffer;
2443 - usec_t all_started_ut = now_monotonic_usec();
2444 -
2445 - store_sql_statements((struct judy_list_t *)data->pending_sql_statement);
2446 - store_alert_transitions((struct judy_list_t *)data->pending_alert_list);
2447 - store_ctx_cleanup_list(wc, (struct judy_list_t *)data->pending_ctx_cleanup_list);
2448 -
2449 - worker_is_busy(UV_EVENT_METADATA_STORE);
2450 -
2395 RRDHOST *host;
2452 - dfe_start_reentrant(rrdhost_root_index, host) {
2396 + size_t host_count = 0;
2397 + usec_t started_ut;
2398 + if (!is_worker) {
2399 + started_ut = now_monotonic_usec();
2400 + dfe_start_reentrant(rrdhost_root_index, host) {
2401 + host_count++;
2402 + }
2403 + dfe_done(host);
2404 + }
2405
2406 + size_t count = 0;
2407 + dfe_start_reentrant(rrdhost_root_index, host)
2408 + {
2409 + count++;
2410 if (rrdhost_flag_check(host, RRDHOST_FLAG_ARCHIVED) || !rrdhost_flag_check(host, RRDHOST_FLAG_METADATA_UPDATE))
2411 continue;
2412
2457 - usec_t started_ut = now_monotonic_usec();
2458 - rrdhost_flag_clear(host,RRDHOST_FLAG_METADATA_UPDATE);
2459 - worker_is_busy(UV_EVENT_STORE_HOST);
2413 + rrdhost_flag_clear(host, RRDHOST_FLAG_METADATA_UPDATE);
2414 + if (is_worker)
2415 + worker_is_busy(UV_EVENT_STORE_HOST);
2416
2417 // store labels, claim_id, host and system info (if needed)
2418 store_host_info_and_metadata(host, work_buffer);
2463 - worker_is_idle();
2419 + if (is_worker)
2420 + worker_is_idle();
2421
2465 - metadata_scan_host(host, work_buffer, shutting_down);
2422 + metadata_scan_host(host, work_buffer, true);
2423
2467 - COMPUTE_DURATION(report_duration, "us", started_ut, now_monotonic_usec());
2468 - nd_log_daemon(NDLP_DEBUG, "Host %s saved metadata in %s", rrdhost_hostname(host), report_duration);
2424 + if (!is_worker)
2425 + nd_log_daemon(NDLP_INFO, "METADATA: Progress of metadata storage: %6.2f%% completed", (100.0 * count / host_count));
2426 }
2427 dfe_done(host);
2428
2429 + if (!is_worker) {
2430 + COMPUTE_DURATION(report_duration, "us", started_ut, now_monotonic_usec());
2431 + nd_log_daemon(
2432 + NDLP_INFO,
2433 + "METADATA: Progress of metadata storage: %6.2f%% completed in %s",
2434 + (100.0 * count / host_count),
2435 + report_duration);
2436 + }
2437 +}
2438 +
2439 +// Worker thread to scan hosts for pending metadata to store
2440 +static void start_metadata_hosts(uv_work_t *req)
2441 +{
2442 + register_libuv_worker_jobs();
2443 +
2444 + struct work_payload *data = req->data;
2445 + struct meta_config_s *config = data->config;
2446 +
2447 + BUFFER *work_buffer = data->work_buffer;
2448 + usec_t all_started_ut = now_monotonic_usec();
2449 +
2450 + store_sql_statements((struct judy_list_t *)data->pending_sql_statement, true);
2451 +
2452 + store_alert_transitions((struct judy_list_t *)data->pending_alert_list, true);
2453 +
2454 + if (!SHUTDOWN_REQUESTED(config))
2455 + store_ctx_cleanup_list(config, (struct judy_list_t *)data->pending_ctx_cleanup_list);
2456 +
2457 + worker_is_busy(UV_EVENT_METADATA_STORE);
2458 +
2459 + store_hosts_metadata(work_buffer, true);
2460
2461 COMPUTE_DURATION(report_duration, "us", all_started_ut, now_monotonic_usec());
2462 nd_log_daemon(NDLP_DEBUG, "Checking all hosts completed in %s", report_duration);
2463
2476 - do_pending_uuid_deletion(wc, (struct judy_list_t *)data->pending_uuid_deletion);
2477 -
2478 - run_metadata_cleanup(wc);
2464 + if (!SHUTDOWN_REQUESTED(config)) {
2465 + do_pending_uuid_deletion(config, (struct judy_list_t *)data->pending_uuid_deletion);
2466 + run_metadata_cleanup(config);
2467 + }
2468
2480 - wc->metadata_check_after = now_realtime_sec() + METADATA_HOST_CHECK_INTERVAL;
2469 + config->metadata_check_after = now_realtime_sec() + METADATA_HOST_CHECK_INTERVAL;
2470 worker_is_idle();
2471 }
2472
2473 #define EVENT_LOOP_NAME "METASYNC"
2474
2475 +#define MAX_SHUTDOWN_TIMEOUT_SECONDS (10)
2476 +#define SHUTDOWN_SLEEP_INTERVAL_MS (100)
2477 +
2478 static void metadata_event_loop(void *arg)
2479 {
2488 - struct metadata_wc *config = arg;
2480 + struct meta_config_s *config = arg;
2481 uv_thread_set_name_np(EVENT_LOOP_NAME);
2482 worker_register(EVENT_LOOP_NAME);
2483
@@ -2495,14 +2487,12 @@ static void metadata_event_loop(void *arg)
2487 worker_register_job_name(METADATA_DEL_DIMENSION, "delete dimension");
2488 worker_register_job_name(METADATA_STORE_CLAIM_ID, "add claim id");
2489 worker_register_job_name(METADATA_ADD_CTX_CLEANUP, "host ctx cleanup");
2498 - worker_register_job_name(METADATA_SCAN_HOSTS, "host metadata store");
2490 + worker_register_job_name(METADATA_STORE, "host metadata store");
2491 worker_register_job_name(METADATA_LOAD_HOST_CONTEXT, "host load context");
2492 worker_register_job_name(METADATA_ADD_HOST_AE, "add host alert entry");
2493 worker_register_job_name(METADATA_DEL_HOST_AE, "delete host alert entry");
2494 worker_register_job_name(METADATA_EXECUTE_STORE_STATEMENT, "add sql statement");
2495
2504 - unsigned cmd_batch_size;
2505 -
2496 uv_loop_t *loop = &config->loop;
2497 fatal_assert(0 == uv_loop_init(loop));
2498 fatal_assert(0 == uv_async_init(loop, &config->async, async_cb));
@@ -2513,25 +2503,22 @@ static void metadata_event_loop(void *arg)
2503 config->timer_req.data = config;
2504
2505 nd_log(NDLS_DAEMON, NDLP_DEBUG, "Starting metadata sync thread");
2516 -
2517 - struct metadata_cmd cmd;
2518 - memset(&cmd, 0, sizeof(cmd));
2519 - metadata_flag_clear(config, METADATA_FLAG_PROCESSING);
2506 config->metadata_check_after = now_realtime_sec() + METADATA_HOST_CHECK_FIRST_CHECK;
2521 - completion_mark_complete(&config->start_stop_complete);
2507
2508 BUFFER *work_buffer = buffer_create(1024, &netdata_buffers_statistics.buffers_sqlite);
2524 - struct scan_metadata_payload *data;
2509 + struct work_payload *data;
2510 Pvoid_t *Pvalue;
2526 - struct judy_list_t *pending_ae_list = NULL;
2511 + struct judy_list_t *pending_alert_list = NULL;
2512 struct judy_list_t *pending_ctx_cleanup_list = NULL;
2513 struct judy_list_t *pending_uuid_deletion = NULL;
2514 struct judy_list_t *pending_sql_statement = NULL;
2515
2531 - int shutdown = 0;
2516 config->initialized = true;
2533 - while (shutdown == 0 || (config->flags & METADATA_FLAG_PROCESSING)) {
2534 - nd_uuid_t *uuid;
2517 + nd_log_daemon(NDLP_INFO, "METADATA: Synchronization thread is up and running");
2518 + completion_mark_complete(&config->start_stop_complete);
2519 +
2520 + while (likely(config->shutdown_requested == false)) {
2521 + nd_uuid_t *uuid;
2522 RRDHOST *host = NULL;
2523 ALARM_ENTRY *ae = NULL;
2524 sqlite3_stmt *stmt;
@@ -2541,21 +2528,14 @@ static void metadata_event_loop(void *arg)
2528 uv_run(loop, UV_RUN_DEFAULT);
2529
2530 /* wait for commands */
2544 - cmd_batch_size = 0;
2531 + unsigned cmd_batch_size = 0;
2532 do {
2546 - if (unlikely(cmd_batch_size >= METADATA_MAX_BATCH_SIZE))
2533 + if (unlikely(++cmd_batch_size >= METADATA_MAX_BATCH_SIZE))
2534 break;
2535
2549 - cmd = metadata_deq_cmd(config);
2536 + struct metadata_cmd cmd = metadata_deq_cmd(config);
2537 opcode = cmd.opcode;
2538
2552 - if (unlikely(opcode == METADATA_DATABASE_NOOP && metadata_flag_check(config, METADATA_FLAG_SHUTDOWN))) {
2553 - shutdown = 1;
2554 - continue;
2555 - }
2556 -
2557 - ++cmd_batch_size;
2558 -
2539 if (likely(opcode != METADATA_DATABASE_NOOP))
2540 worker_is_busy(opcode);
2541
@@ -2563,7 +2543,7 @@ static void metadata_event_loop(void *arg)
2543 case METADATA_DATABASE_NOOP:
2544 break;
2545 case METADATA_DEL_DIMENSION:
2566 - uuid = (nd_uuid_t *) cmd.param[0];
2546 + uuid = (nd_uuid_t *)cmd.param[0];
2547 if (!pending_uuid_deletion)
2548 pending_uuid_deletion = callocz(1, sizeof(*pending_uuid_deletion));
2549
@@ -2577,18 +2557,18 @@ static void metadata_event_loop(void *arg)
2557 }
2558 break;
2559 case METADATA_STORE_CLAIM_ID:
2580 - store_claim_id((nd_uuid_t *) cmd.param[0], (nd_uuid_t *) cmd.param[1]);
2581 - freez((void *) cmd.param[0]);
2582 - freez((void *) cmd.param[1]);
2560 + store_claim_id((nd_uuid_t *)cmd.param[0], (nd_uuid_t *)cmd.param[1]);
2561 + freez((void *)cmd.param[0]);
2562 + freez((void *)cmd.param[1]);
2563 break;
2564
2565 case METADATA_ADD_CTX_CLEANUP:
2566 if (!pending_ctx_cleanup_list)
2567 pending_ctx_cleanup_list = callocz(1, sizeof(*pending_ctx_cleanup_list));
2568
2589 - struct host_ctx_cleanup_s *ctx_cleanup = (struct host_ctx_cleanup_s *) cmd.param[0];
2569 + struct host_ctx_cleanup_s *ctx_cleanup = (struct host_ctx_cleanup_s *)cmd.param[0];
2570 Pvalue = JudyLIns(&pending_ctx_cleanup_list->JudyL, ++pending_ctx_cleanup_list->count, PJE0);
2591 - if (Pvalue != PJERR)
2571 + if (Pvalue && Pvalue != PJERR)
2572 *Pvalue = ctx_cleanup;
2573 else {
2574 // Failure in Judy, attempt to continue running anyway
@@ -2597,124 +2577,120 @@ static void metadata_event_loop(void *arg)
2577 freez(ctx_cleanup);
2578 }
2579 break;
2600 - case METADATA_SCAN_HOSTS:
2601 - if (unlikely(metadata_flag_check(config, METADATA_FLAG_PROCESSING)))
2602 - break;
2603 -
2604 - if (unittest_running)
2580 + case METADATA_STORE:
2581 + if (config->metadata_running || unittest_running)
2582 break;
2583
2584 data = mallocz(sizeof(*data));
2585 data->request.data = data;
2609 - data->wc = config;
2610 - data->pending_alert_list = pending_ae_list;
2586 + data->config = config;
2587 + data->pending_alert_list = pending_alert_list;
2588 data->pending_ctx_cleanup_list = pending_ctx_cleanup_list;
2589 data->pending_uuid_deletion = pending_uuid_deletion;
2590 data->pending_sql_statement = pending_sql_statement;
2591
2592 data->work_buffer = work_buffer;
2616 - pending_ae_list = NULL;
2593 + pending_alert_list = NULL;
2594 pending_ctx_cleanup_list = NULL;
2595 pending_uuid_deletion = NULL;
2596 pending_sql_statement = NULL;
2620 -
2621 - if (unlikely(cmd.completion))
2622 - cmd.completion = NULL; // Do not complete after launching worker (worker will do)
2623 -
2624 - metadata_flag_set(config, METADATA_FLAG_PROCESSING);
2597 + config->metadata_running = true;
2598 if (uv_queue_work(loop, &data->request, start_metadata_hosts, after_metadata_hosts)) {
2626 - // Failed to launch worker -- let the event loop handle completion
2627 - cmd.completion = config->scan_complete;
2628 - pending_ae_list = data->pending_alert_list;
2599 + pending_alert_list = data->pending_alert_list;
2600 pending_ctx_cleanup_list = data->pending_ctx_cleanup_list;
2601 pending_uuid_deletion = data->pending_uuid_deletion;
2602 pending_sql_statement = data->pending_sql_statement;
2603 freez(data);
2633 - metadata_flag_clear(config, METADATA_FLAG_PROCESSING);
2604 + config->metadata_running = false;
2605 }
2606 break;
2607 case METADATA_LOAD_HOST_CONTEXT:
2637 - if (unittest_running)
2608 + if (config->ctx_load_running || unittest_running)
2609 break;
2610
2640 - data = callocz(1,sizeof(*data));
2611 + config->ctx_load_running = true;
2612 + data = callocz(1, sizeof(*data));
2613 data->request.data = data;
2642 - data->wc = config;
2643 - if (uv_queue_work(loop, &data->request, start_all_host_load_context, after_start_host_load_context)) {
2614 + data->config = config;
2615 + if (uv_queue_work(loop, &data->request, ctx_hosts_load, after_ctx_hosts_load)) {
2616 freez(data);
2617 + config->ctx_load_running = false;
2618 + // Fallback reset context so hosts will load on demand
2619 + reset_host_context_load_flag();
2620 }
2621 break;
2622 case METADATA_ADD_HOST_AE:
2648 - host = (RRDHOST *) cmd.param[0];
2649 - ae = (ALARM_ENTRY *) cmd.param[1];
2623 + host = (RRDHOST *)cmd.param[0];
2624 + ae = (ALARM_ENTRY *)cmd.param[1];
2625
2651 - if (!pending_ae_list)
2652 - pending_ae_list = callocz(1, sizeof(*pending_ae_list));
2626 + if (!pending_alert_list)
2627 + pending_alert_list = callocz(1, sizeof(*pending_alert_list));
2628
2654 - Pvalue = JudyLIns(&pending_ae_list->JudyL, ++pending_ae_list->count, PJE0);
2655 - if (Pvalue)
2656 - *Pvalue = (void *)host;
2629 + Pvalue = JudyLIns(&pending_alert_list->JudyL, ++pending_alert_list->count, PJE0);
2630 + if (!Pvalue || Pvalue == PJERR)
2631 + fatal("METASYNC: Corrupted pending_alert_list Judy array");
2632 + *Pvalue = (void *)host;
2633
2658 - Pvalue = JudyLIns(&pending_ae_list->JudyL, ++pending_ae_list->count, PJE0);
2659 - if (Pvalue)
2660 - *Pvalue = (void *)ae;
2634 + Pvalue = JudyLIns(&pending_alert_list->JudyL, ++pending_alert_list->count, PJE0);
2635 + if (!Pvalue || Pvalue == PJERR)
2636 + fatal("METASYNC: Corrupted pending_alert_list Judy array");
2637 + *Pvalue = (void *)ae;
2638 break;
2639 case METADATA_DEL_HOST_AE:
2663 - (void) JudyLIns(&config->ae_DelJudyL, (Word_t) (void *) cmd.param[0], PJE0);
2640 + (void)JudyLIns(&config->ae_DelJudyL, (Word_t)(void *)cmd.param[0], PJE0);
2641 break;
2642 case METADATA_EXECUTE_STORE_STATEMENT:
2666 - stmt = (sqlite3_stmt *) cmd.param[0];
2643 + stmt = (sqlite3_stmt *)cmd.param[0];
2644 if (!pending_sql_statement)
2645 pending_sql_statement = callocz(1, sizeof(*pending_sql_statement));
2646
2647 Pvalue = JudyLIns(&pending_sql_statement->JudyL, ++pending_sql_statement->count, PJE0);
2671 - if (Pvalue)
2648 + if (Pvalue && Pvalue != PJERR)
2649 *Pvalue = (void *)stmt;
2650 + else {
2651 + // Fallback execute immediately
2652 + execute_statement(stmt);
2653 + }
2654 + break;
2655 + case METADATA_SYNC_SHUTDOWN:
2656 + config->shutdown_requested = true;
2657 break;
2658 case METADATA_UNITTEST:;
2675 - struct thread_unittest *tu = (struct thread_unittest *) cmd.param[0];
2659 + struct thread_unittest *tu = (struct thread_unittest *)cmd.param[0];
2660 sleep_usec(1000); // processing takes 1ms
2661 __atomic_fetch_add(&tu->processed, 1, __ATOMIC_SEQ_CST);
2662 break;
2663 default:
2664 break;
2665 }
2682 -
2683 - if (cmd.completion)
2684 - completion_mark_complete(cmd.completion);
2666 } while (opcode != METADATA_DATABASE_NOOP);
2667 }
2668 config->initialized = false;
2669
2670 + if (!uv_timer_stop(&config->timer_req))
2671 + uv_close((uv_handle_t *)&config->timer_req, NULL);
2672 +
2673 + uv_close((uv_handle_t *)&config->async, NULL);
2674 uv_walk(loop, libuv_close_callback, NULL);
2690 - uv_run(loop, UV_RUN_NOWAIT);
2675
2692 - int rc;
2693 - do {
2694 - rc = uv_loop_close(loop);
2695 - } while (rc != UV_EBUSY);
2676 + size_t loop_count = (MAX_SHUTDOWN_TIMEOUT_SECONDS * USEC_PER_MS) / SHUTDOWN_SLEEP_INTERVAL_MS;
2677 + while ((config->metadata_running || config->ctx_load_running) && --loop_count) {
2678 + if (!uv_run(loop, UV_RUN_NOWAIT))
2679 + break; // No pending callbacks
2680 + uv_sleep(SHUTDOWN_SLEEP_INTERVAL_MS);
2681 + }
2682
2697 - buffer_free(work_buffer);
2698 - worker_unregister();
2683 + (void)uv_loop_close(loop);
2684
2700 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "Shutting down metadata thread");
2701 - completion_mark_complete(&config->start_stop_complete);
2702 - if (config->scan_complete) {
2703 - completion_destroy(config->scan_complete);
2704 - freez(config->scan_complete);
2705 - }
2685 + store_hosts_metadata(work_buffer, false);
2686
2707 - Word_t Index;
2708 - bool first;
2687 + store_alert_transitions(pending_alert_list, false);
2688
2710 - if (pending_ae_list) {
2711 - (void)JudyLFreeArray(&pending_ae_list->JudyL, PJE0);
2712 - freez(pending_ae_list);
2713 - }
2689 + store_sql_statements(pending_sql_statement, false);
2690
2691 if (pending_ctx_cleanup_list) {
2716 - Index = 0;
2717 - first = true;
2692 + Word_t Index = 0;
2693 + bool first = true;
2694 while ((Pvalue = JudyLFirstThenNext(pending_ctx_cleanup_list->JudyL, &Index, &first))) {
2695 if (!*Pvalue)
2696 continue;
@@ -2726,91 +2702,43 @@ static void metadata_event_loop(void *arg)
2702 freez(pending_ctx_cleanup_list);
2703 }
2704
2705 + buffer_free(work_buffer);
2706 metadata_free_cmd_queue(config);
2707 aral_by_size_release(config->ar);
2708 worker_unregister();
2709 +
2710 + completion_mark_complete(&config->start_stop_complete);
2711 }
2712
2713 void metadata_sync_shutdown(void)
2714 {
2736 - completion_init(&metasync_worker.start_stop_complete);
2737 -
2715 struct metadata_cmd cmd;
2716 memset(&cmd, 0, sizeof(cmd));
2717 nd_log_daemon(NDLP_DEBUG, "METADATA: Sending a shutdown command");
2718 cmd.opcode = METADATA_SYNC_SHUTDOWN;
2742 - metadata_enq_cmd(&metasync_worker, &cmd);
2719 + metadata_enq_cmd(&cmd);
2720
2721 /* wait for metadata thread to shut down */
2722 nd_log_daemon(NDLP_DEBUG, "METADATA: Waiting for shutdown ACK");
2746 - completion_wait_for(&metasync_worker.start_stop_complete);
2747 - completion_destroy(&metasync_worker.start_stop_complete);
2748 - int rc = uv_thread_join(&metasync_worker.thread);
2723 + completion_wait_for(&meta_config.start_stop_complete);
2724 + completion_destroy(&meta_config.start_stop_complete);
2725 + int rc = uv_thread_join(&meta_config.thread);
2726 if (rc)
2727 nd_log_daemon(NDLP_ERR, "METADATA: Failed to join synchronization thread, error %s", uv_err_name(rc));
2728 else
2729 nd_log_daemon(NDLP_INFO, "METADATA: synchronization thread shutdown completed");
2730 }
2731
2755 -void metadata_sync_shutdown_prepare(void)
2756 -{
2757 - static bool running = false;
2758 - if (unlikely(!metasync_worker.initialized || running))
2759 - return;
2760 -
2761 - running = true;
2762 -
2763 - struct metadata_cmd cmd;
2764 - memset(&cmd, 0, sizeof(cmd));
2765 -
2766 - struct metadata_wc *wc = &metasync_worker;
2767 -
2768 - struct completion *compl = mallocz(sizeof(*compl));
2769 - completion_init(compl);
2770 - __atomic_store_n(&wc->scan_complete, compl, __ATOMIC_RELAXED);
2771 -
2772 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "METADATA: Sending a scan host command");
2773 - uint32_t max_wait_iterations = 2000;
2774 - while (unlikely(metadata_flag_check(&metasync_worker, METADATA_FLAG_PROCESSING)) && max_wait_iterations--) {
2775 - if (max_wait_iterations == 1999)
2776 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "METADATA: Current worker is running; waiting to finish");
2777 - sleep_usec(1000);
2778 - }
2779 -
2780 - cmd.opcode = METADATA_SCAN_HOSTS;
2781 - metadata_enq_cmd(&metasync_worker, &cmd);
2782 -
2783 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "METADATA: Waiting for host scan completion");
2784 - completion_wait_for(wc->scan_complete);
2785 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "METADATA: Host scan complete; can continue with shutdown");
2786 -}
2787 -
2788 -void *metadata_sync_shutdown_thread(void *ptr __maybe_unused) {
2789 - metadata_sync_shutdown_prepare();
2790 - return NULL;
2791 -}
2792 -
2793 -static ND_THREAD *metdata_sync_shutdown_background_wait_thread = NULL;
2794 -void metadata_sync_shutdown_background(void) {
2795 - metdata_sync_shutdown_background_wait_thread = nd_thread_create(
2796 - "METASYNC-SHUTDOWN", NETDATA_THREAD_OPTION_JOINABLE, metadata_sync_shutdown_thread, NULL);
2797 -}
2798 -
2799 -void metadata_sync_shutdown_background_wait(void) {
2800 - nd_thread_join(metdata_sync_shutdown_background_wait_thread);
2801 - metadata_sync_shutdown();
2802 -}
2803 -
2732 // -------------------------------------------------------------
2733 // Init function called on agent startup
2734
2735 void metadata_sync_init(void)
2736 {
2809 - memset(&metasync_worker, 0, sizeof(metasync_worker));
2810 - completion_init(&metasync_worker.start_stop_complete);
2737 + memset(&meta_config, 0, sizeof(meta_config));
2738 + completion_init(&meta_config.start_stop_complete);
2739
2740 int retries = 0;
2813 - int create_uv_thread_rc = create_uv_thread(&metasync_worker.thread, metadata_event_loop, &metasync_worker, &retries);
2741 + int create_uv_thread_rc = create_uv_thread(&meta_config.thread, metadata_event_loop, &meta_config, &retries);
2742 if (create_uv_thread_rc)
2743 nd_log_daemon(NDLP_ERR, "Failed to create SQLite metadata sync thread, error %s, after %d retries", uv_err_name(create_uv_thread_rc), retries);
2744
@@ -2819,21 +2747,22 @@ void metadata_sync_init(void)
2747 if (retries)
2748 nd_log_daemon(NDLP_WARNING, "SQLite metadata sync thread was created after %d attempts", retries);
2749
2822 - completion_wait_for(&metasync_worker.start_stop_complete);
2823 - completion_destroy(&metasync_worker.start_stop_complete);
2824 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "SQLite metadata sync initialization complete");
2750 + // Wait for initialization
2751 + completion_wait_for(&meta_config.start_stop_complete);
2752 +
2753 + // Reset the completion, we will use it again during shutdown
2754 + completion_reset(&meta_config.start_stop_complete);
2755 }
2756
2757 // Helpers
2758
2829 -static inline void queue_metadata_cmd(enum metadata_opcode opcode, const void *param0, const void *param1)
2759 +static inline bool queue_metadata_cmd(enum metadata_opcode opcode, const void *param0, const void *param1)
2760 {
2761 struct metadata_cmd cmd;
2762 cmd.opcode = opcode;
2763 cmd.param[0] = param0;
2764 cmd.param[1] = param1;
2835 - cmd.completion = NULL;
2836 - metadata_enq_cmd(&metasync_worker, &cmd);
2765 + return metadata_enq_cmd(&cmd);
2766 }
2767
2768 // Public
@@ -2844,7 +2773,8 @@ void metaqueue_delete_dimension_uuid(nd_uuid_t *uuid)
2773
2774 nd_uuid_t *use_uuid = mallocz(sizeof(*uuid));
2775 uuid_copy(*use_uuid, *uuid);
2847 - queue_metadata_cmd(METADATA_DEL_DIMENSION, use_uuid, NULL);
2776 + if (!queue_metadata_cmd(METADATA_DEL_DIMENSION, use_uuid, NULL))
2777 + freez(use_uuid);
2778 }
2779
2780 void metaqueue_store_claim_id(nd_uuid_t *host_uuid, nd_uuid_t *claim_uuid)
@@ -2860,7 +2790,10 @@ void metaqueue_store_claim_id(nd_uuid_t *host_uuid, nd_uuid_t *claim_uuid)
2790 local_claim_uuid = mallocz(sizeof(*claim_uuid));
2791 uuid_copy(*local_claim_uuid, *claim_uuid);
2792 }
2863 - queue_metadata_cmd(METADATA_STORE_CLAIM_ID, local_host_uuid, local_claim_uuid);
2793 + if (unlikely(!queue_metadata_cmd(METADATA_STORE_CLAIM_ID, local_host_uuid, local_claim_uuid))) {
2794 + freez(local_host_uuid);
2795 + freez(local_claim_uuid);
2796 + }
2797 }
2798
2799 void metaqueue_ml_load_models(RRDDIM *rd)
@@ -2868,10 +2801,9 @@ void metaqueue_ml_load_models(RRDDIM *rd)
2801 rrddim_flag_set(rd, RRDDIM_FLAG_ML_MODEL_LOAD);
2802 }
2803
2871 -void metadata_queue_load_host_context()
2804 +bool metadata_queue_load_host_context()
2805 {
2873 - queue_metadata_cmd(METADATA_LOAD_HOST_CONTEXT, NULL, NULL);
2874 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "Queued command to load host contexts");
2806 + return queue_metadata_cmd(METADATA_LOAD_HOST_CONTEXT, NULL, NULL);
2807 }
2808
2809 void metadata_queue_ctx_host_cleanup(nd_uuid_t *host_uuid, const char *context)
@@ -2884,17 +2816,27 @@ void metadata_queue_ctx_host_cleanup(nd_uuid_t *host_uuid, const char *context)
2816 uuid_copy(ctx_cleanup->host_uuid, *host_uuid);
2817 ctx_cleanup->context = string_strdupz(context);
2818
2887 - queue_metadata_cmd(METADATA_ADD_CTX_CLEANUP, ctx_cleanup, NULL);
2819 + if (unlikely(!queue_metadata_cmd(METADATA_ADD_CTX_CLEANUP, ctx_cleanup, NULL))) {
2820 + string_freez(ctx_cleanup->context);
2821 + freez(ctx_cleanup);
2822 + }
2823 }
2824
2890 -void metadata_queue_ae_save(RRDHOST *host, ALARM_ENTRY *ae)
2825 +bool metadata_queue_ae_save(RRDHOST *host, ALARM_ENTRY *ae)
2826 {
2827 if (unlikely(!host || !ae))
2893 - return;
2828 + return true;
2829
2830 __atomic_add_fetch(&host->health.pending_transitions, 1, __ATOMIC_RELAXED);
2831 __atomic_add_fetch(&ae->pending_save_count, 1, __ATOMIC_RELAXED);
2897 - queue_metadata_cmd(METADATA_ADD_HOST_AE, host, ae);
2832 +
2833 + if (unlikely(!queue_metadata_cmd(METADATA_ADD_HOST_AE, host, ae))) {
2834 + // Failed to queue, reset counters
2835 + __atomic_sub_fetch(&host->health.pending_transitions, 1, __ATOMIC_RELAXED);
2836 + __atomic_sub_fetch(&ae->pending_save_count, 1, __ATOMIC_RELAXED);
2837 + return false;
2838 + }
2839 + return true;
2840 }
2841
2842 void metadata_queue_ae_deletion(ALARM_ENTRY *ae)
@@ -2902,7 +2844,7 @@ void metadata_queue_ae_deletion(ALARM_ENTRY *ae)
2844 if (unlikely(!ae))
2845 return;
2846
2905 - queue_metadata_cmd(METADATA_DEL_HOST_AE, ae, NULL);
2847 + (void) queue_metadata_cmd(METADATA_DEL_HOST_AE, ae, NULL);
2848 }
2849
2850 void metadata_execute_store_statement(sqlite3_stmt *stmt)
@@ -2910,12 +2852,12 @@ void metadata_execute_store_statement(sqlite3_stmt *stmt)
2852 if (unlikely(!stmt))
2853 return;
2854
2913 - queue_metadata_cmd(METADATA_EXECUTE_STORE_STATEMENT, stmt, NULL);
2855 + (void) queue_metadata_cmd(METADATA_EXECUTE_STORE_STATEMENT, stmt, NULL);
2856 }
2857
2858 void commit_alert_transitions(RRDHOST *host __maybe_unused)
2859 {
2918 - queue_metadata_cmd(METADATA_SCAN_HOSTS, NULL, NULL);
2860 + (void) queue_metadata_cmd(METADATA_STORE, NULL, NULL);
2861 }
2862
2863 uint64_t sqlite_get_meta_space(void)
@@ -3007,12 +2949,11 @@ static void *unittest_queue_metadata(void *arg) {
2949 cmd.opcode = METADATA_UNITTEST;
2950 cmd.param[0] = tu;
2951 cmd.param[1] = NULL;
3010 - cmd.completion = NULL;
3011 - metadata_enq_cmd(&metasync_worker, &cmd);
2952 + metadata_enq_cmd(&cmd);
2953
2954 do {
2955 __atomic_fetch_add(&tu->added, 1, __ATOMIC_SEQ_CST);
3015 - metadata_enq_cmd(&metasync_worker, &cmd);
2956 + metadata_enq_cmd(&cmd);
2957 sleep_usec(10000);
2958 } while (!__atomic_load_n(&tu->join, __ATOMIC_RELAXED));
2959 return arg;
@@ -3050,7 +2991,7 @@ static void *metadata_unittest_threads(void)
2991 unittest_queue_metadata,
2992 &tu);
2993 }
3053 - (void) uv_async_send(&metasync_worker.async);
2994 + (void) uv_async_send(&meta_config.async);
2995 sleep_usec(seconds_to_run * USEC_PER_SEC);
2996
2997 __atomic_store_n(&tu.join, 1, __ATOMIC_RELAXED);
src/database/sqlite/sqlite_metadata.h
+5 -6
@@ -26,18 +26,17 @@ typedef enum db_check_action_type {
26 // To initialize and shutdown
27 void metadata_sync_init(void);
28 void metadata_sync_shutdown(void);
29 -void metadata_sync_shutdown_prepare(void);
29
30 void metaqueue_delete_dimension_uuid(nd_uuid_t *uuid);
31 void metaqueue_store_claim_id(nd_uuid_t *host_uuid, nd_uuid_t *claim_uuid);
32 void metaqueue_ml_load_models(RRDDIM *rd);
33 void detect_machine_guid_change(nd_uuid_t *host_uuid);
35 -void metadata_queue_load_host_context();
34 +bool metadata_queue_load_host_context();
35 +void reset_host_context_load_flag();
36 void vacuum_database(sqlite3 *database, const char *db_alias, int threshold, int vacuum_pc);
37
38 int sql_metadata_cache_stats(int op);
39
40 -int get_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id);
40 void sql_update_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id);
41 void sql_load_node_id(RRDHOST *host);
42
@@ -53,12 +52,12 @@ int sql_init_meta_database(db_check_action_type_t rebuild, int memory);
52 void cleanup_agent_event_log(void);
53 void add_agent_event(event_log_type_t event_id, int64_t value);
54 usec_t get_agent_event_time_median(event_log_type_t event_id);
56 -void metadata_queue_ae_save(RRDHOST *host, ALARM_ENTRY *ae);
55 +bool metadata_queue_ae_save(RRDHOST *host, ALARM_ENTRY *ae);
56 void metadata_queue_ae_deletion(ALARM_ENTRY *ae);
57 void commit_alert_transitions(RRDHOST *host);
58
60 -void metadata_sync_shutdown_background(void);
61 -void metadata_sync_shutdown_background_wait(void);
59 +//void metadata_sync_shutdown_background(void);
60 +//void metadata_sync_shutdown_background_wait(void);
61 void metadata_queue_ctx_host_cleanup(nd_uuid_t *host_uuid, const char *context);
62 void store_host_info_and_metadata(RRDHOST *host, BUFFER *work_buffer);
63 void metadata_execute_store_statement(sqlite3_stmt *stmt);
src/health/health_log.c
+5 -2
@@ -43,9 +43,12 @@ void health_alarm_entry_destroy(ALARM_ENTRY *ae) {
43
44 inline void health_alarm_log_save(RRDHOST *host, ALARM_ENTRY *ae, bool async)
45 {
46 + bool saved = false;
47 if (async)
47 - metadata_queue_ae_save(host, ae);
48 - else
48 + saved = metadata_queue_ae_save(host, ae);
49 +
50 + // if not async or async failed to queue, do it now
51 + if (false == saved)
52 sql_health_alarm_log_save(host, ae);
53 }
54