Make sure ACLK sync thread completes initialization (#19916)
* Use completion to make sure the aclk sync thread has finished initialization * Validate inputs in aclk functions to prevent null pointer dereferences
Stelios Fragkakis committed
Mar 20, 2025 at 18:41 UTC
448a0812e94f4b7e44705741a876d1581206d337
1 file changed
+12
-6
src/database/sqlite/sqlite_aclk.c
+12
-6
@@ -48,6 +48,7 @@ struct aclk_sync_config_s {
48
bool aclk_batch_job_is_running;
49
SPINLOCK cmd_queue_lock;
50
uint32_t aclk_jobs_pending;
51
+ struct completion start_stop_complete;
52
struct aclk_database_cmd *cmd_base;
53
ARAL *ar;
54
} aclk_sync_config = { 0 };
@@ -642,6 +643,8 @@ static void aclk_synchronization(void *arg)
643
struct aclk_query_payload *payload;
644
645
unsigned cmd_batch_size;
646
+
647
+ completion_mark_complete(&config->start_stop_complete);
648
while (likely(service_running(SERVICE_ACLK))) {
649
enum aclk_database_opcode opcode;
650
worker_is_idle();
@@ -886,7 +889,10 @@ static void aclk_synchronization(void *arg)
889
static void aclk_synchronization_init(void)
890
{
891
memset(&aclk_sync_config, 0, sizeof(aclk_sync_config));
892
+ completion_init(&aclk_sync_config.start_stop_complete);
893
fatal_assert(0 == uv_thread_create(&aclk_sync_config.thread, aclk_synchronization, &aclk_sync_config));
894
+ completion_wait_for(&aclk_sync_config.start_stop_complete);
895
+ completion_destroy(&aclk_sync_config.start_stop_complete);
896
}
897
898
// -------------------------------------------------------------
@@ -980,7 +986,7 @@ static inline void queue_aclk_sync_cmd(enum aclk_database_opcode opcode, const v
986
// Public
987
void aclk_push_alert_config(const char *node_id, const char *config_hash)
988
{
983
- if (unlikely(!aclk_sync_config.initialized))
989
+ if (unlikely(!node_id || !config_hash))
990
return;
991
992
queue_aclk_sync_cmd(ACLK_DATABASE_PUSH_ALERT_CONFIG, strdupz(node_id), strdupz(config_hash));
@@ -988,7 +994,7 @@ void aclk_push_alert_config(const char *node_id, const char *config_hash)
994
995
void aclk_execute_query(aclk_query_t query)
996
{
991
- if (unlikely(!aclk_sync_config.initialized))
997
+ if (unlikely(!query))
998
return;
999
1000
queue_aclk_sync_cmd(ACLK_QUERY_EXECUTE, query, NULL);
@@ -996,20 +1002,20 @@ void aclk_execute_query(aclk_query_t query)
1002
1003
void aclk_add_job(aclk_query_t query)
1004
{
999
- if (unlikely(!aclk_sync_config.initialized))
1005
+ if (unlikely(!query))
1006
return;
1007
1008
queue_aclk_sync_cmd(ACLK_QUERY_BATCH_ADD, query, NULL);
1009
}
1010
1005
-void aclk_query_init(mqtt_wss_client client) {
1006
-
1011
+void aclk_query_init(mqtt_wss_client client)
1012
+{
1013
queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT, client, NULL);
1014
}
1015
1016
void schedule_node_state_update(RRDHOST *host, uint64_t delay)
1017
{
1012
- if (unlikely(!aclk_sync_config.initialized || !host))
1018
+ if (unlikely(!host))
1019
return;
1020
1021
queue_aclk_sync_cmd(ACLK_DATABASE_NODE_STATE, host, (void *)(uintptr_t)delay);