Further improve ACLK synchronization shutdown (#20105)
Improve ACLK Sync and ACLK MQTT shutdown
Stelios Fragkakis committed
Apr 10, 2025 at 00:10 UTC
ca0183d4b8b74833606b554e2cf6ad0cc9aa5dc6
1 file changed
+20
-16
src/database/sqlite/sqlite_aclk.c
+20
-16
@@ -74,10 +74,10 @@ static struct aclk_database_cmd aclk_database_deq_cmd(void)
74
return ret;
75
}
76
77
-static void aclk_database_enq_cmd(struct aclk_database_cmd *cmd)
77
+static bool aclk_database_enq_cmd(struct aclk_database_cmd *cmd)
78
{
79
if(unlikely(!__atomic_load_n(&aclk_sync_config.initialized, __ATOMIC_RELAXED)))
80
- return;
80
+ return false;
81
82
struct aclk_database_cmd *t = aral_mallocz(aclk_sync_config.ar);
83
*t = *cmd;
@@ -88,6 +88,7 @@ static void aclk_database_enq_cmd(struct aclk_database_cmd *cmd)
88
spinlock_unlock(&aclk_sync_config.cmd_queue_lock);
89
90
(void) uv_async_send(&aclk_sync_config.async);
91
+ return true;
92
}
93
94
enum {
@@ -1047,22 +1048,24 @@ void aclk_synchronization_init(void)
1048
nd_log_daemon(NDLP_INFO, "ACLK sync initialization completed");
1049
}
1050
1050
-static inline void queue_aclk_sync_cmd(enum aclk_database_opcode opcode, const void *param0, const void *param1)
1051
+static inline bool queue_aclk_sync_cmd(enum aclk_database_opcode opcode, const void *param0, const void *param1)
1052
{
1053
struct aclk_database_cmd cmd;
1054
cmd.opcode = opcode;
1055
cmd.param[0] = (void *) param0;
1056
cmd.param[1] = (void *) param1;
1056
- aclk_database_enq_cmd(&cmd);
1057
+ return aclk_database_enq_cmd(&cmd);
1058
}
1059
1060
void aclk_synchronization_shutdown(void)
1061
{
1062
// Send shutdown command, not that the completion is initialized
1063
// on init and still valid
1063
- queue_aclk_sync_cmd(ACLK_SYNC_SHUTDOWN, NULL, NULL);
1064
+ aclk_mqtt_client_reset();
1065
+
1066
+ if (queue_aclk_sync_cmd(ACLK_SYNC_SHUTDOWN, NULL, NULL))
1067
+ completion_wait_for(&aclk_sync_config.start_stop_complete);
1068
1065
- completion_wait_for(&aclk_sync_config.start_stop_complete);
1069
completion_destroy(&aclk_sync_config.start_stop_complete);
1070
int rc = uv_thread_join(&aclk_sync_config.thread);
1071
if (rc)
@@ -1085,7 +1088,7 @@ void aclk_execute_query(aclk_query_t query)
1088
if (unlikely(!query))
1089
return;
1090
1088
- queue_aclk_sync_cmd(ACLK_QUERY_EXECUTE, query, NULL);
1091
+ (void) queue_aclk_sync_cmd(ACLK_QUERY_EXECUTE, query, NULL);
1092
}
1093
1094
void aclk_add_job(aclk_query_t query)
@@ -1093,12 +1096,12 @@ void aclk_add_job(aclk_query_t query)
1096
if (unlikely(!query))
1097
return;
1098
1096
- queue_aclk_sync_cmd(ACLK_QUERY_BATCH_ADD, query, NULL);
1099
+ (void) queue_aclk_sync_cmd(ACLK_QUERY_BATCH_ADD, query, NULL);
1100
}
1101
1102
void aclk_mqtt_client_set(mqtt_wss_client client)
1103
{
1101
- queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT_SET, client, NULL);
1104
+ (void) queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT_SET, client, NULL);
1105
}
1106
1107
void aclk_mqtt_client_reset()
@@ -1108,8 +1111,8 @@ void aclk_mqtt_client_reset()
1111
1112
struct completion compl;
1113
completion_init(&compl);
1111
- queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT_RESET, &compl, NULL);
1112
- completion_wait_for(&compl);
1114
+ if (queue_aclk_sync_cmd(ACLK_MQTT_WSS_CLIENT_RESET, &compl, NULL))
1115
+ completion_wait_for(&compl);
1116
completion_destroy(&compl);
1117
}
1118
@@ -1118,14 +1121,15 @@ void schedule_node_state_update(RRDHOST *host, uint64_t delay)
1121
if (unlikely(!host))
1122
return;
1123
1121
- queue_aclk_sync_cmd(ACLK_DATABASE_NODE_STATE, host, (void *)(uintptr_t)delay);
1124
+ (void) queue_aclk_sync_cmd(ACLK_DATABASE_NODE_STATE, host, (void *)(uintptr_t)delay);
1125
}
1126
1127
void unregister_node(const char *machine_guid)
1128
{
1129
if (unlikely(!machine_guid))
1130
return;
1128
- queue_aclk_sync_cmd(ACLK_DATABASE_NODE_UNREGISTER, strdupz(machine_guid), NULL);
1131
+
1132
+ (void) queue_aclk_sync_cmd(ACLK_DATABASE_NODE_UNREGISTER, strdupz(machine_guid), NULL);
1133
}
1134
1135
void destroy_aclk_config(RRDHOST *host)
@@ -1141,8 +1145,8 @@ void destroy_aclk_config(RRDHOST *host)
1145
struct completion compl;
1146
completion_init(&compl);
1147
1144
- queue_aclk_sync_cmd(ACLK_CANCEL_NODE_UPDATE_TIMER, (void *)host, (void *)&compl);
1145
- completion_wait_for(&compl);
1148
+ if (queue_aclk_sync_cmd(ACLK_CANCEL_NODE_UPDATE_TIMER, (void *)host, (void *)&compl))
1149
+ completion_wait_for(&compl);
1150
completion_destroy(&compl);
1151
}
1152
@@ -1152,5 +1156,5 @@ void destroy_aclk_config(RRDHOST *host)
1156
1157
void aclk_queue_node_info(RRDHOST *host, bool immediate)
1158
{
1155
- queue_aclk_sync_cmd(ACLK_QUEUE_NODE_INFO, (void *)host, (void *)(uintptr_t)immediate);
1159
+ (void) queue_aclk_sync_cmd(ACLK_QUEUE_NODE_INFO, (void *)host, (void *)(uintptr_t)immediate);
1160
}