Fix crash during shutdown when there are pending messages to cloud (#20080)
* Notify the ACLK sync event loop that the MQTT connection is shutting down Prevent queueing messages if mqtt client has been destroyed / free resources * Refactor MQTT client usage in HTTP API V2 handling
Stelios Fragkakis committed
Apr 7, 2025 at 22:37 UTC
3f824362d6a4a5017a747fdc77580390c37f32fb
4 files changed
+44
-12
src/aclk/aclk.c
+7
-2
@@ -868,7 +868,8 @@ void *aclk_main(void *ptr)
868
// Keep reconnecting and talking until our time has come
869
// and the Grim Reaper (exit_initiated) calls
870
netdata_log_info("ACLK: Starting ACLK query event loop");
871
- aclk_query_init(mqttwss_client);
871
+ aclk_mqtt_client_set(mqttwss_client);
872
+ bool client_to_reset = true;
873
do {
874
worker_is_busy(WORKER_ACLK_CONNECT);
875
if (aclk_attempt_to_connect(mqttwss_client))
@@ -888,7 +889,9 @@ void *aclk_main(void *ptr)
889
nd_log(NDLS_ACCESS, NDLP_WARNING, "ACLK DISCONNECTED");
890
}
891
} while (service_running(SERVICE_ACLK));
891
-
892
+ aclk_mqtt_client_reset();
893
+ // No need to reset the client again when exiting
894
+ client_to_reset = false;
895
worker_is_busy(WORKER_ACLK_DISCONNECTED);
896
aclk_graceful_disconnect(mqttwss_client);
897
@@ -899,6 +902,8 @@ void *aclk_main(void *ptr)
902
903
exit_full:
904
free_topic_cache();
905
+ if (client_to_reset)
906
+ aclk_mqtt_client_reset();
907
mqtt_wss_destroy(mqttwss_client);
908
exit:
909
if (aclk_env) {
src/aclk/aclk_query.h
+2
-1
@@ -13,7 +13,8 @@ int mark_pending_req_cancelled(const char *msg_id);
13
void mark_pending_req_cancel_all();
14
15
void aclk_execute_query(aclk_query_t query);
16
-void aclk_query_init(mqtt_wss_client client);
16
+void aclk_mqtt_client_set(mqtt_wss_client client);
17
+void aclk_mqtt_client_reset();
18
int http_api_v2(mqtt_wss_client client, aclk_query_t query);
19
int send_bin_msg(mqtt_wss_client client, aclk_query_t query);
20
src/database/sqlite/sqlite_aclk.c
+33
-8
@@ -358,13 +358,15 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
358
}
359
360
bool ok_to_send = true;
361
+ mqtt_wss_client client = __atomic_load_n(&config->client, __ATOMIC_RELAXED);
362
363
switch (query->type) {
364
365
// Incoming : cloud -> agent
366
case HTTP_API_V2:
367
worker_is_busy(UV_EVENT_ACLK_QUERY_EXECUTE);
367
- http_api_v2(config->client, query);
368
+ if (client)
369
+ http_api_v2(client, query);
370
ok_to_send = false;
371
break;
372
case CTX_CHECKPOINT:;
@@ -429,8 +431,14 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
431
break;
432
}
433
432
- if (ok_to_send)
433
- send_bin_msg(config->client, query);
434
+ if (ok_to_send) {
435
+ if (client)
436
+ send_bin_msg(client, query);
437
+ else {
438
+ freez(query->data.bin_payload.payload);
439
+ nd_log_daemon(NDLP_ERR, "No client to send message %u", query->type);
440
+ }
441
+ }
442
443
aclk_query_free(query);
444
}
@@ -610,7 +618,8 @@ static void aclk_synchronization_event_loop(void *arg)
618
worker_register_job_name(ACLK_QUERY_EXECUTE_SYNC, "aclk query execute sync");
619
worker_register_job_name(ACLK_QUERY_BATCH_EXECUTE, "aclk batch execute");
620
worker_register_job_name(ACLK_QUERY_BATCH_ADD, "aclk batch add");
613
- worker_register_job_name(ACLK_MQTT_WSS_CLIENT, "config mqtt client");
621
+ worker_register_job_name(ACLK_MQTT_WSS_CLIENT_SET, "config mqtt client");
622
+ worker_register_job_name(ACLK_MQTT_WSS_CLIENT_RESET, "reset mqtt client");
623
worker_register_job_name(ACLK_DATABASE_NODE_UNREGISTER, "unregister node");
624
625
uv_loop_t *loop = &config->loop;
@@ -750,10 +759,14 @@ static void aclk_synchronization_event_loop(void *arg)
759
config->alert_push_running = false;
760
}
761
break;
753
- case ACLK_MQTT_WSS_CLIENT:
762
+ case ACLK_MQTT_WSS_CLIENT_SET:
763
config->client = (mqtt_wss_client)cmd.param[0];
764
break;
756
-
765
+ case ACLK_MQTT_WSS_CLIENT_RESET:
766
+ __atomic_store_n(&config->client, NULL, __ATOMIC_RELEASE);
767
+ struct completion *comp = cmd.param[0];
768
+ completion_mark_complete(comp);
769
+ break;
770
case ACLK_QUERY_EXECUTE:
771
query = (aclk_query_t)cmd.param[0];
772
@@ -1044,9 +1057,21 @@ void aclk_add_job(aclk_query_t query)
1057
queue_aclk_sync_cmd(ACLK_QUERY_BATCH_ADD, query, NULL);
1058
}
1059
1047
-void aclk_query_init(mqtt_wss_client client)
1060
+void aclk_mqtt_client_set(mqtt_wss_client client)
1061
{
1049
- queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT, client, NULL);
1062
+ queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT_SET, client, NULL);
1063
+}
1064
+
1065
+void aclk_mqtt_client_reset()
1066
+{
1067
+ if (!__atomic_load_n(&aclk_sync_config.client, __ATOMIC_RELAXED))
1068
+ return;
1069
+
1070
+ struct completion compl;
1071
+ completion_init(&compl);
1072
+ queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT_RESET, &compl, NULL);
1073
+ completion_wait_for(&compl);
1074
+ completion_destroy(&compl);
1075
}
1076
1077
void schedule_node_state_update(RRDHOST *host, uint64_t delay)
src/database/sqlite/sqlite_aclk.h
+2
-1
@@ -21,7 +21,8 @@ enum aclk_database_opcode {
21
ACLK_DATABASE_PUSH_ALERT,
22
ACLK_DATABASE_PUSH_ALERT_CONFIG,
23
ACLK_DATABASE_NODE_UNREGISTER,
24
- ACLK_MQTT_WSS_CLIENT,
24
+ ACLK_MQTT_WSS_CLIENT_SET,
25
+ ACLK_MQTT_WSS_CLIENT_RESET,
26
ACLK_QUERY_EXECUTE,
27
ACLK_QUERY_EXECUTE_SYNC,
28
ACLK_QUERY_BATCH_ADD,