@cryptotaxi247 / netdata-1 / commits / 70c345729

Reduce memory allocations in event loops (#20306)

* Add Worker pool to metadata/ aclk Add cmd pool for metadata / aclk event loops Do not queue commands from the event loop Cleanup ACLK sync commands Add unittest test for cmd_pool * Remove unused code and declarations * Move unittest code Align cmd_data_t * Remove emoticon (Unicode ) characters from unittest report * fix compilation warning * Increase pool size * Initialize worker allocation state in worker pool setup * Fix query pool initialization

Stelios Fragkakis committed May 21, 2025 at 08:47 UTC 70c345729620ad5d7c4a68568b96c1a336dcf804
14 files changed +497 -295
src/aclk/aclk.c
+2 -2
@@ -947,7 +947,7 @@ void aclk_create_node_instance_job(RRDHOST *host)
947 if (!claim_id_is_set(claim_id))
948 return;
949
950 - aclk_query_t query = aclk_query_new(REGISTER_NODE);
950 + aclk_query_t *query = aclk_query_new(REGISTER_NODE);
951 int32_t hops = rrdhost_ingestion_hops(host);
952 node_instance_creation_t node_instance_creation = {
953 .hops = hops,
@@ -974,7 +974,7 @@ void aclk_update_node_instance_job(RRDHOST *host, int live, int queryable)
974 if (!claim_id_is_set(claim_id))
975 return;
976
977 - aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
977 + aclk_query_t *query = aclk_query_new(NODE_STATE_UPDATE);
978
979 int32_t hops = rrdhost_ingestion_hops(host);
980 node_instance_connection_t node_state_update = {
src/aclk/aclk_alarm_api.c
+2 -2
@@ -18,7 +18,7 @@ void aclk_send_alarm_log_entry(struct alarm_log_entry *log_entry)
18
19 void aclk_send_provide_alarm_cfg(struct provide_alarm_configuration *cfg)
20 {
21 - aclk_query_t query = aclk_query_new(ALARM_PROVIDE_CFG);
21 + aclk_query_t *query = aclk_query_new(ALARM_PROVIDE_CFG);
22 query->data.bin_payload.payload = generate_provide_alarm_configuration(&query->data.bin_payload.size, cfg);
23 query->data.bin_payload.topic = ACLK_TOPICID_ALARM_CONFIG;
24 query->data.bin_payload.msg_name = "ProvideAlarmConfiguration";
@@ -27,7 +27,7 @@ void aclk_send_provide_alarm_cfg(struct provide_alarm_configuration *cfg)
27
28 void aclk_send_alarm_snapshot(alarm_snapshot_proto_ptr_t snapshot)
29 {
30 - aclk_query_t query = aclk_query_new(ALARM_SNAPSHOT);
30 + aclk_query_t *query = aclk_query_new(ALARM_SNAPSHOT);
31 query->data.bin_payload.payload = generate_alarm_snapshot_bin(&query->data.bin_payload.size, snapshot);
32 query->data.bin_payload.topic = ACLK_TOPICID_ALARM_SNAPSHOT;
33 query->data.bin_payload.msg_name = "AlarmSnapshot";
src/aclk/aclk_contexts_api.c
+4 -4
@@ -6,7 +6,7 @@
6
7 void aclk_send_contexts_snapshot(contexts_snapshot_t data)
8 {
9 - aclk_query_t query = aclk_query_new(CTX_SEND_SNAPSHOT);
9 + aclk_query_t *query = aclk_query_new(CTX_SEND_SNAPSHOT);
10 query->data.bin_payload.topic = ACLK_TOPICID_CTXS_SNAPSHOT;
11 query->data.bin_payload.payload = contexts_snapshot_2bin(data, &query->data.bin_payload.size);
12 query->data.bin_payload.msg_name = "ContextsSnapshot";
@@ -15,7 +15,7 @@ void aclk_send_contexts_snapshot(contexts_snapshot_t data)
15
16 void aclk_send_contexts_updated(contexts_updated_t data)
17 {
18 - aclk_query_t query = aclk_query_new(CTX_SEND_SNAPSHOT_UPD);
18 + aclk_query_t *query = aclk_query_new(CTX_SEND_SNAPSHOT_UPD);
19 query->data.bin_payload.topic = ACLK_TOPICID_CTXS_UPDATED;
20 query->data.bin_payload.payload = contexts_updated_2bin(data, &query->data.bin_payload.size);
21 query->data.bin_payload.msg_name = "ContextsUpdated";
@@ -24,7 +24,7 @@ void aclk_send_contexts_updated(contexts_updated_t data)
24
25 void aclk_update_node_collectors(struct update_node_collectors *collectors)
26 {
27 - aclk_query_t query = aclk_query_new(UPDATE_NODE_COLLECTORS);
27 + aclk_query_t *query = aclk_query_new(UPDATE_NODE_COLLECTORS);
28 query->data.bin_payload.topic = ACLK_TOPICID_NODE_COLLECTORS;
29 query->data.bin_payload.payload = generate_update_node_collectors_message(&query->data.bin_payload.size, collectors);
30 query->data.bin_payload.msg_name = "UpdateNodeCollectors";
@@ -33,7 +33,7 @@ void aclk_update_node_collectors(struct update_node_collectors *collectors)
33
34 void aclk_update_node_info(struct update_node_info *info)
35 {
36 - aclk_query_t query = aclk_query_new(UPDATE_NODE_INFO);
36 + aclk_query_t *query = aclk_query_new(UPDATE_NODE_INFO);
37 query->data.bin_payload.topic = ACLK_TOPICID_NODE_INFO;
38 query->data.bin_payload.payload = generate_update_node_info_message(&query->data.bin_payload.size, info);
39 query->data.bin_payload.msg_name = "UpdateNodeInfo";
src/aclk/aclk_query.c
+2 -2
@@ -103,7 +103,7 @@ static bool aclk_web_client_interrupt_cb(struct web_client *w __maybe_unused, vo
103 return req->canceled;
104 }
105
106 -int http_api_v2(mqtt_wss_client client, aclk_query_t query)
106 +int http_api_v2(mqtt_wss_client client, aclk_query_t *query)
107 {
108 ND_LOG_STACK lgs[] = {
109 ND_LOG_FIELD_TXT(NDF_SRC_TRANSPORT, "aclk"),
@@ -227,7 +227,7 @@ cleanup:
227 return retval;
228 }
229
230 -int send_bin_msg(mqtt_wss_client client, aclk_query_t query)
230 +int send_bin_msg(mqtt_wss_client client, aclk_query_t *query)
231 {
232 // this will be simplified when legacy support is removed
233 aclk_send_bin_message_subtopic_pid(
src/aclk/aclk_query.h
+3 -3
@@ -12,10 +12,10 @@
12 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);
15 +void aclk_execute_query(aclk_query_t *query);
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);
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
21 #endif //NETDATA_AGENT_CLOUD_LINK_H
src/aclk/aclk_query_queue.c
+57 -9
@@ -2,14 +2,66 @@
2
3 #include "aclk_query_queue.h"
4
5 -aclk_query_t aclk_query_new(aclk_query_type_t type)
5 +#define MAX_QUERY_ENTRIES (512)
6 +
7 +struct {
8 + aclk_query_t query_workers[MAX_QUERY_ENTRIES];
9 + int free_stack[MAX_QUERY_ENTRIES];
10 + int top;
11 + SPINLOCK spinlock;
12 +} queryPool;
13 +
14 +// Initialize the query pool
15 +__attribute__((constructor)) void init_query_pool()
16 +{
17 + spinlock_init(&queryPool.spinlock);
18 + for (int i = 0; i < MAX_QUERY_ENTRIES; i++) {
19 + queryPool.free_stack[i] = i;
20 + queryPool.query_workers[i].allocated = false;
21 + }
22 + queryPool.top = MAX_QUERY_ENTRIES;
23 +}
24 +
25 +static aclk_query_t *get_query()
26 +{
27 + spinlock_lock(&queryPool.spinlock);
28 + if (queryPool.top == 0) {
29 + spinlock_unlock(&queryPool.spinlock);
30 + aclk_query_t *query = callocz(1, sizeof(aclk_query_t));
31 + query->allocated = true;
32 + return query;
33 + }
34 + int index = queryPool.free_stack[--queryPool.top];
35 + memset(&queryPool.query_workers[index], 0, sizeof(aclk_query_t));
36 + spinlock_unlock(&queryPool.spinlock);
37 + return &queryPool.query_workers[index];
38 +}
39 +
40 +static void return_query(aclk_query_t *query)
41 {
7 - aclk_query_t query = callocz(1, sizeof(struct aclk_query));
42 + if (unlikely(query->allocated)) {
43 + freez(query);
44 + return;
45 + }
46 + spinlock_lock(&queryPool.spinlock);
47 + int index = (int) (query - queryPool.query_workers);
48 + if (index < 0 || index >= MAX_QUERY_ENTRIES) {
49 + spinlock_unlock(&queryPool.spinlock);
50 + return; // Invalid (should not happen)
51 + }
52 + queryPool.free_stack[queryPool.top++] = index;
53 + memset(query, 0, sizeof(aclk_query_t));
54 + spinlock_unlock(&queryPool.spinlock);
55 +}
56 +
57 +aclk_query_t *aclk_query_new(aclk_query_type_t type)
58 +{
59 + aclk_query_t *query = get_query();
60 query->type = type;
61 return query;
62 }
63
12 -void aclk_query_free(aclk_query_t query)
64 +void aclk_query_free(aclk_query_t *query)
65 {
66 struct ctxs_checkpoint *cmd;
67 switch (query->type) {
@@ -29,12 +81,8 @@ void aclk_query_free(aclk_query_t query)
81 freez(query->data.node_id);
82 freez(query->machine_guid);
83 break;
84 + // keep following cases together
85 case CTX_STOP_STREAMING:
33 - cmd = query->data.payload;
34 - freez(cmd->claim_id);
35 - freez(cmd->node_id);
36 - freez(cmd);
37 - break;
86 case CTX_CHECKPOINT:
87 cmd = query->data.payload;
88 freez(cmd->claim_id);
@@ -49,5 +97,5 @@ void aclk_query_free(aclk_query_t query)
97 freez(query->dedup_id);
98 freez(query->callback_topic);
99 freez(query->msg_id);
52 - freez(query);
100 + return_query(query);
101 }
src/aclk/aclk_query_queue.h
+7 -7
@@ -40,9 +40,9 @@ struct aclk_bin_payload {
40 const char *msg_name;
41 };
42
43 -typedef struct aclk_query *aclk_query_t;
44 -struct aclk_query {
43 +typedef struct {
44 aclk_query_type_t type;
45 + bool allocated;
46
47 // dedup_id is used to deduplicate queries in the list
48 // if type and dedup_id is the same message is deduplicated
@@ -68,13 +68,13 @@ struct aclk_query {
68 void *payload;
69 char *node_id;
70 } data;
71 -};
71 +} aclk_query_t;
72
73 -aclk_query_t aclk_query_new(aclk_query_type_t type);
74 -void aclk_query_free(aclk_query_t query);
73 +aclk_query_t *aclk_query_new(aclk_query_type_t type);
74 +void aclk_query_free(aclk_query_t *query);
75
76 -void aclk_execute_query(aclk_query_t query);
77 -void aclk_add_job(aclk_query_t query);
76 +void aclk_execute_query(aclk_query_t *query);
77 +void aclk_add_job(aclk_query_t *query);
78
79 #define QUEUE_IF_PAYLOAD_PRESENT(query) \
80 do { \
src/aclk/aclk_rx_msgs.c
+9 -10
@@ -125,8 +125,6 @@ static inline int aclk_v2_payload_get_query(const char *payload, char **query_ur
125
126 static int aclk_handle_cloud_http_request_v2(struct aclk_request *cloud_to_agent, char *raw_payload)
127 {
128 - aclk_query_t query;
129 -
128 errno_clear();
129 if (cloud_to_agent->version < ACLK_V_COMPRESSION) {
130 netdata_log_error(
@@ -136,7 +134,7 @@ static int aclk_handle_cloud_http_request_v2(struct aclk_request *cloud_to_agent
134 return 1;
135 }
136
139 - query = aclk_query_new(HTTP_API_V2);
137 + aclk_query_t *query = aclk_query_new(HTTP_API_V2);
138
139 if (unlikely(aclk_extract_v2_data(raw_payload, &query->data.http_api_v2.payload))) {
140 netdata_log_error("Error extracting payload expected after the JSON dictionary.");
@@ -255,7 +253,8 @@ int create_node_instance_result(const char *msg, size_t msg_len)
253
254 netdata_log_debug(D_ACLK, "CreateNodeInstanceResult: guid:%s nodeid:%s", res.machine_guid, res.node_id);
255
258 - aclk_query_t query = aclk_query_new(CREATE_NODE_INSTANCE);
256 + aclk_query_t *query = aclk_query_new(CREATE_NODE_INSTANCE);
257 +
258 query->data.node_id = res.node_id; // Will be freed on query free
259 query->machine_guid = res.machine_guid; // Will be freed on query free
260 aclk_add_job(query);
@@ -266,7 +265,7 @@ int send_node_instances(const char *msg, size_t msg_len)
265 {
266 UNUSED(msg);
267 UNUSED(msg_len);
269 - aclk_query_t query = aclk_query_new(SEND_NODE_INSTANCES);
268 + aclk_query_t *query = aclk_query_new(SEND_NODE_INSTANCES);
269 aclk_add_job(query);
270 return 0;
271 }
@@ -302,7 +301,7 @@ int start_alarm_streaming(const char *msg, size_t msg_len)
301 netdata_log_error("Error parsing StartAlarmStreaming");
302 return 1;
303 }
305 - aclk_query_t query = aclk_query_new(ALERT_START_STREAMING);
304 + aclk_query_t *query = aclk_query_new(ALERT_START_STREAMING);
305 query->data.node_id = res.node_id; // Will be freed on query free
306 query->version = res.version;
307 aclk_add_job(query);
@@ -318,7 +317,7 @@ int send_alarm_checkpoint(const char *msg, size_t msg_len)
317 freez(sac.claim_id);
318 return 1;
319 }
321 - aclk_query_t query = aclk_query_new(ALERT_CHECKPOINT);
320 + aclk_query_t *query = aclk_query_new(ALERT_CHECKPOINT);
321 query->data.node_id = sac.node_id; // Will be freed on query free
322 query->claim_id = sac.claim_id;
323 query->version = sac.version;
@@ -347,7 +346,7 @@ int send_alarm_snapshot(const char *msg, size_t msg_len)
346 destroy_send_alarm_snapshot(sas);
347 return 1;
348 }
350 - aclk_query_t query = aclk_query_new(ALERT_CHECKPOINT);
349 + aclk_query_t *query = aclk_query_new(ALERT_CHECKPOINT);
350 query->data.node_id = sas->node_id; // Will be freed on query free
351 query->claim_id = sas->claim_id; // Will be freed on query free
352 query->version = 0; // force snapshot
@@ -386,7 +385,7 @@ int contexts_checkpoint(const char *msg, size_t msg_len)
385 if (!cmd)
386 return 1;
387
389 - aclk_query_t query = aclk_query_new(CTX_CHECKPOINT);
388 + aclk_query_t *query = aclk_query_new(CTX_CHECKPOINT);
389 query->data.payload = cmd;
390 aclk_add_job(query);
391 return 0;
@@ -403,7 +402,7 @@ int stop_streaming_contexts(const char *msg, size_t msg_len)
402 if (!cmd)
403 return 1;
404
406 - aclk_query_t query = aclk_query_new(CTX_STOP_STREAMING);
405 + aclk_query_t *query = aclk_query_new(CTX_STOP_STREAMING);
406 query->data.payload = cmd;
407 aclk_add_job(query);
408 return 0;
src/daemon/libuv_workers.c
+163
@@ -119,3 +119,166 @@ void libuv_close_callback(uv_handle_t *handle, void *data __maybe_unused)
119 uv_close(handle, NULL);
120 }
121 }
122 +
123 +// Initialize the worker pool
124 +void init_worker_pool(WorkerPool *pool) {
125 + for (int i = 0; i < MAX_ACTIVE_WORKERS; i++) {
126 + pool->workers[i].allocated = false;
127 + pool->free_stack[i] = i; // Fill the stack with indices
128 + }
129 + pool->top = MAX_ACTIVE_WORKERS; // All workers are initially free
130 +}
131 +
132 +// Get a worker (reuse if available, NULL if pool exhausted)
133 +worker_data_t *get_worker(WorkerPool *pool) {
134 + if (pool->top == 0) {
135 + worker_data_t *worker = callocz(1, sizeof(worker_data_t));
136 + worker->allocated = true; // Mark as allocated
137 + return worker;
138 + }
139 + int index = pool->free_stack[--pool->top]; // Pop from stack
140 + return &pool->workers[index];
141 +}
142 +
143 +// Return a worker for reuse
144 +void return_worker(WorkerPool *pool, worker_data_t *worker) {
145 + if (unlikely(worker->allocated)) {
146 + freez(worker);
147 + return;
148 + }
149 +
150 + int index = (int) (worker - pool->workers);
151 + if (index < 0 || index >= MAX_ACTIVE_WORKERS) {
152 + return; // Invalid worker (should not happen)
153 + }
154 + pool->free_stack[pool->top++] = index; // Push index back to stack
155 +}
156 +
157 +// Initialize the command pool
158 +void init_cmd_pool(CmdPool *pool, int size) {
159 + pool->buffer = mallocz(sizeof(cmd_data_t) * size);
160 +
161 + pool->size = size;
162 + pool->head = 0;
163 + pool->tail = 0;
164 + pool->count = 0;
165 +
166 + uv_mutex_init(&pool->lock);
167 + uv_cond_init(&pool->not_full);
168 +;}
169 +
170 +bool push_cmd(CmdPool *pool, const cmd_data_t *cmd, bool wait_on_full)
171 +{
172 + uv_mutex_lock(&pool->lock);
173 +
174 + while (pool->count == pool->size) {
175 + if (wait_on_full)
176 + uv_cond_wait(&pool->not_full, &pool->lock);
177 + else {
178 + uv_mutex_unlock(&pool->lock); // No space, return
179 + return false;
180 + }
181 + }
182 +
183 + pool->buffer[pool->tail] = *cmd;
184 + pool->tail = (pool->tail + 1) % pool->size;
185 + pool->count++;
186 +
187 + uv_mutex_unlock(&pool->lock);
188 + return true;
189 +}
190 +
191 +bool pop_cmd(CmdPool *pool, cmd_data_t *out_cmd) {
192 + uv_mutex_lock(&pool->lock);
193 + if (pool->count == 0) {
194 + uv_mutex_unlock(&pool->lock); // No commands to pop
195 + return false;
196 + }
197 + *out_cmd = pool->buffer[pool->head];
198 + pool->head = (pool->head + 1) % pool->size;
199 + pool->count--;
200 +
201 + uv_cond_signal(&pool->not_full);
202 + uv_mutex_unlock(&pool->lock);
203 + return true;
204 +}
205 +
206 +void release_cmd_pool(CmdPool *pool) {
207 + if (pool->buffer) {
208 + free(pool->buffer);
209 + pool->buffer = NULL;
210 + }
211 + uv_mutex_destroy(&pool->lock);
212 + uv_cond_destroy(&pool->not_full);
213 +}
214 +
215 +/// Test
216 +
217 +typedef struct {
218 + CmdPool *pool;
219 + int total;
220 + int failed;
221 +} ThreadArgs;
222 +
223 +void push_thread(void *arg) {
224 + ThreadArgs *args = (ThreadArgs *)arg;
225 + CmdPool *pool = args->pool;
226 + for (int i = 0; i < args->total; ++i) {
227 + cmd_data_t cmd;
228 + snprintf(cmd.data, sizeof(cmd.data), "cmd-%d", i);
229 + push_cmd(pool, &cmd, true);
230 + }
231 + fprintf(stderr, "PUSHED: %d commands\n", args->total);
232 +}
233 +
234 +void pop_thread(void *arg) {
235 + ThreadArgs *args = (ThreadArgs *)arg;
236 + CmdPool *pool = args->pool;
237 +
238 + cmd_data_t cmd;
239 + for (int i = 0; i < args->total; ) {
240 + bool got = pop_cmd(pool, &cmd);
241 + if (got) {
242 + char expected[64];
243 + snprintf(expected, sizeof(expected), "cmd-%d", i);
244 + if (strcmp(cmd.data, expected) != 0) {
245 + fprintf(stderr, "POPPED: %s --- EXPECTED %s FAILED\n", cmd.data, expected);
246 + args->failed++;
247 + }
248 + i++;
249 + } else {
250 + uv_sleep(1); // avoid busy spin
251 + }
252 + }
253 + fprintf(stderr, "POPPED: %d commands\n", args->total);
254 +}
255 +
256 +int test_cmd_pool_fifo()
257 +{
258 + CmdPool pool;
259 +
260 + int pool_sizes[] = {32, 64, 128, 256};
261 +
262 + for (size_t i = 0; i < sizeof(pool_sizes) / sizeof(pool_sizes[0]); ++i) {
263 + int pool_size = pool_sizes[i];
264 + init_cmd_pool(&pool, pool_size);
265 +
266 + ThreadArgs args = {.pool = &pool, .total = 1000, .failed = 0};
267 + uv_thread_t producer, consumer;
268 + fprintf(stderr, "Testing pool size %d\n", pool_size);
269 +
270 + uv_thread_create(&producer, push_thread, &args);
271 + uv_thread_create(&consumer, pop_thread, &args);
272 +
273 + uv_thread_join(&producer);
274 + uv_thread_join(&consumer);
275 +
276 + release_cmd_pool(&pool);
277 + if (args.failed) {
278 + fprintf(stderr, "Multithreaded FIFO test failed with %d errors.\n", args.failed);
279 + return 1;
280 + }
281 + }
282 + fprintf(stderr, "Multithreaded FIFO test passed.\n");
283 + return 0;
284 +}
src/daemon/libuv_workers.h
+53
@@ -86,7 +86,60 @@ enum event_loop_job {
86 UV_EVENT_SCHEDULE_CMD,
87 };
88
89 +#define MAX_ACTIVE_WORKERS (256)
90 +
91 +typedef struct worker_data {
92 + uv_work_t request;
93 + void *config;
94 + void *pending_alert_list;
95 + void *pending_ctx_cleanup_list;
96 + void *pending_uuid_deletion;
97 + void *pending_sql_statement;
98 + union {
99 + void *payload;
100 + void *work_buffer;
101 + };
102 + bool allocated;
103 +} worker_data_t;
104 +
105 +typedef struct {
106 + worker_data_t workers[MAX_ACTIVE_WORKERS]; // Preallocated worker data pool
107 + int free_stack[MAX_ACTIVE_WORKERS]; // Stack of available worker data indices
108 + int top; // Stack pointer
109 +} WorkerPool;
110 +
111 +typedef struct {
112 + uint8_t opcode;
113 + uint8_t padding[sizeof(void *) - sizeof(uint8_t)]; // Padding to align the union
114 + union {
115 + void *param[2];
116 + char data[sizeof(void *) * 2];
117 + };
118 +} cmd_data_t;
119 +
120 +typedef struct {
121 + cmd_data_t *buffer;
122 + int size;
123 + int head;
124 + int tail;
125 + int count;
126 +
127 + uv_mutex_t lock;
128 + uv_cond_t not_full;
129 +} CmdPool;
130 +
131 +
132 void register_libuv_worker_jobs();
133 void libuv_close_callback(uv_handle_t *handle, void *data __maybe_unused);
134
135 +void init_worker_pool(WorkerPool *pool);
136 +worker_data_t *get_worker(WorkerPool *pool);
137 +void return_worker(WorkerPool *pool, worker_data_t *worker);
138 +
139 +void init_cmd_pool(CmdPool *pool, int size);
140 +bool push_cmd(CmdPool *pool, const cmd_data_t *cmd, bool wait_on_full);
141 +bool pop_cmd(CmdPool *pool, cmd_data_t *out_cmd);
142 +void release_cmd_pool(CmdPool *pool);
143 +int test_cmd_pool_fifo();
144 +
145 #endif //NETDATA_EVENT_LOOP_H
src/daemon/main.c
+4
@@ -452,6 +452,10 @@ int netdata_main(int argc, char **argv) {
452 unittest_running = true;
453 return buffer_unittest();
454 }
455 + else if(strcmp(optarg, "test_cmd_pool_fifo") == 0) {
456 + unittest_running = true;
457 + return test_cmd_pool_fifo();
458 + }
459 else if(strcmp(optarg, "uuidtest") == 0) {
460 unittest_running = true;
461 return uuid_unittest();
src/database/sqlite/sqlite_aclk.c
+126 -136
@@ -11,7 +11,6 @@ void sanity_check(void) {
11 #include "sqlite_aclk_node.h"
12 #include "aclk/aclk_query_queue.h"
13 #include "aclk/aclk_query.h"
14 -#include "aclk/aclk_capas.h"
14
15 static void create_node_instance_result_job(const char *machine_guid, const char *node_id)
16 {
@@ -44,51 +43,32 @@ struct aclk_sync_config_s {
43 bool initialized;
44 mqtt_wss_client client;
45 int aclk_queries_running;
46 + bool run_query_batch;
47 bool alert_push_running;
48 bool aclk_batch_job_is_running;
49 - SPINLOCK cmd_queue_lock;
49 uint32_t aclk_jobs_pending;
50 struct completion start_stop_complete;
52 - struct aclk_database_cmd *cmd_base;
53 - ARAL *ar;
51 + CmdPool cmd_pool;
52 + WorkerPool worker_pool;
53 } aclk_sync_config = { 0 };
54
56 -static struct aclk_database_cmd aclk_database_deq_cmd(void)
55 +static cmd_data_t aclk_database_deq_cmd(void)
56 {
58 - struct aclk_database_cmd ret = { 0 };
59 - struct aclk_database_cmd *to_free = NULL;
60 -
61 - spinlock_lock(&aclk_sync_config.cmd_queue_lock);
62 - if(aclk_sync_config.cmd_base) {
63 - struct aclk_database_cmd *t = aclk_sync_config.cmd_base;
64 - DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(aclk_sync_config.cmd_base, t, prev, next);
65 - ret = *t;
66 - to_free = t;
67 - }
68 - else {
69 - ret.opcode = ACLK_DATABASE_NOOP;
70 - }
71 - spinlock_unlock(&aclk_sync_config.cmd_queue_lock);
72 - aral_freez(aclk_sync_config.ar, to_free);
73 -
57 + cmd_data_t ret = { 0 };
58 + ret.opcode = ACLK_DATABASE_NOOP;
59 + (void) pop_cmd(&aclk_sync_config.cmd_pool, (cmd_data_t *) &ret);
60 return ret;
61 }
62
77 -static bool aclk_database_enq_cmd(struct aclk_database_cmd *cmd)
63 +static bool aclk_database_enq_cmd(cmd_data_t *cmd, bool wait_on_full)
64 {
65 if(unlikely(!__atomic_load_n(&aclk_sync_config.initialized, __ATOMIC_RELAXED)))
66 return false;
67
82 - struct aclk_database_cmd *t = aral_mallocz(aclk_sync_config.ar);
83 - *t = *cmd;
84 - t->prev = t->next = NULL;
85 -
86 - spinlock_lock(&aclk_sync_config.cmd_queue_lock);
87 - DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(aclk_sync_config.cmd_base, t, prev, next);
88 - spinlock_unlock(&aclk_sync_config.cmd_queue_lock);
89 -
90 - (void) uv_async_send(&aclk_sync_config.async);
91 - return true;
68 + bool added = push_cmd(&aclk_sync_config.cmd_pool, (void *)cmd, wait_on_full);
69 + if (added)
70 + (void) uv_async_send(&aclk_sync_config.async);
71 + return added;
72 }
73
74 enum {
@@ -319,39 +299,15 @@ static void async_cb(uv_async_t *handle)
299
300 #define TIMER_PERIOD_MS (1000)
301
322 -static void timer_cb(uv_timer_t *handle)
323 -{
324 - uv_stop(handle->loop);
325 - uv_update_time(handle->loop);
326 - struct aclk_sync_config_s *config = handle->data;
327 -
328 - struct aclk_database_cmd cmd = { 0 };
329 - if (aclk_online_for_alerts()) {
330 - cmd.opcode = ACLK_DATABASE_PUSH_ALERT;
331 - aclk_database_enq_cmd(&cmd);
332 - }
333 -
334 - if (config->aclk_jobs_pending > 0) {
335 - cmd.opcode = ACLK_QUERY_BATCH_EXECUTE;
336 - aclk_database_enq_cmd(&cmd);
337 - }
338 -}
339 -
340 -struct aclk_query_payload {
341 - uv_work_t request;
342 - void *data;
343 - struct aclk_sync_config_s *config;
344 -};
345 -
302 static void after_aclk_run_query_job(uv_work_t *req, int status __maybe_unused)
303 {
348 - struct aclk_query_payload *payload = req->data;
349 - struct aclk_sync_config_s *config = payload->config;
304 + worker_data_t *worker_data = req->data;
305 + struct aclk_sync_config_s *config = worker_data->config;
306 config->aclk_queries_running--;
351 - freez(payload);
307 + return_worker(&config->worker_pool, worker_data);
308 }
309
354 -static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query)
310 +static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t *query)
311 {
312 if (query->type == UNKNOWN || query->type >= ACLK_QUERY_TYPE_COUNT) {
313 error_report("Unknown query in query queue. %u", query->type);
@@ -448,9 +404,9 @@ static void aclk_run_query_job(uv_work_t *req)
404 {
405 register_libuv_worker_jobs();
406
451 - struct aclk_query_payload *payload = req->data;
452 - struct aclk_sync_config_s *config = payload->config;
453 - aclk_query_t query = (aclk_query_t) payload->data;
407 + worker_data_t *worker_data = req->data;
408 + struct aclk_sync_config_s *config = worker_data->config;
409 + aclk_query_t *query = (aclk_query_t *) worker_data->payload;
410
411 aclk_run_query(config, query);
412 worker_is_idle();
@@ -458,19 +414,19 @@ static void aclk_run_query_job(uv_work_t *req)
414
415 static void after_aclk_execute_batch(uv_work_t *req, int status __maybe_unused)
416 {
461 - struct aclk_query_payload *payload = req->data;
462 - struct aclk_sync_config_s *config = payload->config;
417 + worker_data_t *worker_data = req->data;
418 + struct aclk_sync_config_s *config = worker_data->config;
419 config->aclk_batch_job_is_running = false;
464 - freez(payload);
420 + return_worker(&config->worker_pool, worker_data);
421 }
422
423 static void aclk_execute_batch(uv_work_t *req)
424 {
425 register_libuv_worker_jobs();
426
471 - struct aclk_query_payload *payload = req->data;
472 - struct aclk_sync_config_s *config = payload->config;
473 - struct judy_list_t *aclk_query_batch = payload->data;
427 + worker_data_t *worker_data = req->data;
428 + struct aclk_sync_config_s *config = worker_data->config;
429 + struct judy_list_t *aclk_query_batch = worker_data->payload;
430
431 if (!aclk_query_batch)
432 return;
@@ -485,7 +441,7 @@ static void aclk_execute_batch(uv_work_t *req)
441 if (!*Pvalue)
442 continue;
443
488 - aclk_query_t query = *Pvalue;
444 + aclk_query_t *query = *Pvalue;
445 aclk_run_query(config, query);
446 }
447
@@ -500,11 +456,11 @@ static void aclk_execute_batch(uv_work_t *req)
456 worker_is_idle();
457 }
458
503 -struct worker_data {
504 - uv_work_t request;
505 - void *payload;
506 - struct aclk_sync_config_s *config;
507 -};
459 +//struct worker_data {
460 +// uv_work_t request;
461 +// void *payload;
462 +// struct aclk_sync_config_s *config;
463 +//};
464
465 struct notify_timer_cb_data {
466 void *payload;
@@ -513,19 +469,20 @@ struct notify_timer_cb_data {
469
470 static void after_do_unregister_node(uv_work_t *req, int status __maybe_unused)
471 {
516 - struct worker_data *data = req->data;
517 - freez(data);
472 + worker_data_t *worker_data = req->data;
473 + struct aclk_sync_config_s *config = worker_data->config;
474 + return_worker(&config->worker_pool, worker_data);
475 }
476
477 static void do_unregister_node(uv_work_t *req)
478 {
479 register_libuv_worker_jobs();
480
524 - struct worker_data *data = req->data;
481 + worker_data_t *worker_data = req->data;
482
483 worker_is_busy(UV_EVENT_UNREGISTER_NODE);
484
528 - sql_unregister_node(data->payload);
485 + sql_unregister_node(worker_data->payload);
486
487 worker_is_idle();
488 }
@@ -557,7 +514,7 @@ static void after_start_alert_push(uv_work_t *req, int status __maybe_unused)
514 struct aclk_sync_config_s *config = data->config;
515
516 config->alert_push_running = false;
560 - freez(data);
517 + return_worker(&config->worker_pool, data);
518 }
519
520 // Worker thread to scan hosts for pending metadata to store
@@ -583,16 +540,22 @@ static void start_alert_push(uv_work_t *req __maybe_unused)
540 // config->aclk_queries_running is only accessed from the vent loop
541 // On failure: free the payload
542
586 -int schedule_query_in_worker(uv_loop_t *loop, struct aclk_sync_config_s *config, aclk_query_t query) {
587 - struct aclk_query_payload *payload = mallocz(sizeof(*payload));
588 - payload->request.data = payload;
589 - payload->config = config;
590 - payload->data = query;
543 +int schedule_query_in_worker(uv_loop_t *loop, struct aclk_sync_config_s *config, aclk_query_t *query) {
544 +
545 + worker_data_t *worker_data = get_worker(&config->worker_pool);
546 + if (!worker_data) {
547 + netdata_log_error("Failed to get worker data");
548 + return -1;
549 + }
550 + worker_data->request.data = worker_data;
551 + worker_data->payload = query;
552 + worker_data->config = config;
553 +
554 config->aclk_queries_running++;
592 - int rc = uv_queue_work(loop, &payload->request, aclk_run_query_job, after_aclk_run_query_job);
555 + int rc = uv_queue_work(loop, &worker_data->request, aclk_run_query_job, after_aclk_run_query_job);
556 if (rc) {
557 config->aclk_queries_running--;
595 - freez(payload);
558 + return_worker(&config->worker_pool, worker_data);
559 }
560 return rc;
561 }
@@ -602,7 +565,7 @@ static void free_query_list(Pvoid_t JudyL)
565 bool first = true;
566 Pvoid_t *Pvalue;
567 Word_t Index = 0;
605 - aclk_query_t query;
568 + aclk_query_t *query;
569 while ((Pvalue = JudyLFirstThenNext(JudyL, &Index, &first))) {
570 if (!*Pvalue)
571 continue;
@@ -611,7 +574,31 @@ static void free_query_list(Pvoid_t JudyL)
574 }
575 }
576
577 +static void timer_cb(uv_timer_t *handle)
578 +{
579 + uv_stop(handle->loop);
580 + uv_update_time(handle->loop);
581 + struct aclk_sync_config_s *config = handle->data;
582 +
583 + if (aclk_online_for_alerts()) {
584 + worker_data_t *worker_data;
585 + if (!config->alert_push_running && (worker_data = get_worker(&config->worker_pool))) {
586 + worker_data->request.data = worker_data;
587 + worker_data->config = config;
588 + config->alert_push_running = true;
589 + if (uv_queue_work(handle->loop, &worker_data->request, start_alert_push, after_start_alert_push)) {
590 + config->alert_push_running = false;
591 + return_worker(&config->worker_pool, worker_data);
592 + }
593 + }
594 + }
595 +
596 + if (config->aclk_jobs_pending > 0)
597 + config->run_query_batch = true;
598 +}
599 +
600 #define MAX_SHUTDOWN_TIMEOUT_SECONDS (5)
601 +#define CMD_POOL_SIZE (2048)
602
603 #define ACLK_SYNC_SHOULD_BE_RUNNING \
604 (!shutdown_requested || config->aclk_queries_running || config->alert_push_running || \
@@ -621,23 +608,22 @@ static void *aclk_synchronization_event_loop(void *arg)
608 {
609 struct aclk_sync_config_s *config = arg;
610 uv_thread_set_name_np("ACLKSYNC");
624 - config->ar = aral_by_size_acquire(sizeof(struct aclk_database_cmd));
611 + init_cmd_pool(&config->cmd_pool, CMD_POOL_SIZE);
612 +
613 worker_register("ACLKSYNC");
614
615 service_register(NULL, NULL, NULL);
616
617 worker_register_job_name(ACLK_DATABASE_NOOP, "noop");
618 worker_register_job_name(ACLK_DATABASE_NODE_STATE, "node state");
631 - worker_register_job_name(ACLK_DATABASE_PUSH_ALERT, "alert push");
619 worker_register_job_name(ACLK_DATABASE_PUSH_ALERT_CONFIG, "alert conf push");
633 - worker_register_job_name(ACLK_QUERY_EXECUTE_SYNC, "aclk query execute sync");
620 worker_register_job_name(ACLK_QUERY_BATCH_EXECUTE, "aclk batch execute");
621 worker_register_job_name(ACLK_QUERY_BATCH_ADD, "aclk batch add");
622 worker_register_job_name(ACLK_MQTT_WSS_CLIENT_SET, "config mqtt client");
623 worker_register_job_name(ACLK_MQTT_WSS_CLIENT_RESET, "reset mqtt client");
624 worker_register_job_name(ACLK_DATABASE_NODE_UNREGISTER, "unregister node");
625 worker_register_job_name(ACLK_CANCEL_NODE_UPDATE_TIMER, "cancel node update timer");
640 - worker_register_job_name(ACLK_QUEUE_NODE_INFO, "queue node info");
626 + worker_register_job_name(ACLK_QUEUE_NODE_INFO, "queue node info");
627
628 uv_loop_t *loop = &config->loop;
629 fatal_assert(0 == uv_loop_init(loop));
@@ -656,9 +642,9 @@ static void *aclk_synchronization_event_loop(void *arg)
642 int query_thread_count = netdata_conf_cloud_query_threads();
643 netdata_log_info("Starting ACLK synchronization thread with %d parallel query threads", query_thread_count);
644
659 - struct worker_data *data;
645 + //struct worker_data *worker_datadata;
646 struct notify_timer_cb_data *timer_cb_data;
661 - aclk_query_t query;
647 + aclk_query_t *query;
648
649 // This holds queries that need to be executed one by one
650 struct judy_list_t *aclk_query_batch = NULL;
@@ -667,7 +653,7 @@ static void *aclk_synchronization_event_loop(void *arg)
653 size_t pending_queries = 0;
654
655 Pvoid_t *Pvalue;
670 - struct aclk_query_payload *payload;
656 + worker_data_t *worker_data;
657
658 unsigned cmd_batch_size;
659
@@ -695,11 +681,20 @@ static void *aclk_synchronization_event_loop(void *arg)
681 /* wait for commands */
682 cmd_batch_size = 0;
683 do {
684 + cmd_data_t cmd;
685 +
686 if (unlikely(++cmd_batch_size >= MAX_BATCH_SIZE))
687 break;
688
701 - struct aclk_database_cmd cmd = aclk_database_deq_cmd();
702 - opcode = cmd.opcode;
689 + if (config->run_query_batch) {
690 + opcode = ACLK_QUERY_BATCH_EXECUTE;
691 + config->run_query_batch = false;
692 + }
693 + else
694 + {
695 + cmd = aclk_database_deq_cmd();
696 + opcode = cmd.opcode;
697 + }
698
699 if(likely(opcode != ACLK_DATABASE_NOOP && opcode != ACLK_QUERY_EXECUTE))
700 worker_is_busy(opcode);
@@ -777,35 +772,24 @@ static void *aclk_synchronization_event_loop(void *arg)
772 break;
773
774 case ACLK_DATABASE_NODE_UNREGISTER:
780 - data = mallocz(sizeof(*data));
781 - data->request.data = data;
782 - data->config = config;
783 - data->payload = cmd.param[0];
784 -
785 - if (uv_queue_work(loop, &data->request, do_unregister_node, after_do_unregister_node)) {
786 - freez(data->payload);
787 - freez(data);
775 + worker_data = get_worker(&config->worker_pool);
776 + if (!worker_data) {
777 + nd_log_daemon(NDLP_ERR, "Failed to get worker for unregister node");
778 + freez(cmd.param[0]);
779 + break;
780 + }
781 + worker_data->request.data = worker_data;
782 + worker_data->config = config;
783 + worker_data->payload = cmd.param[0];
784 +
785 + if (uv_queue_work(loop, &worker_data->request, do_unregister_node, after_do_unregister_node)) {
786 + freez(cmd.param[0]);
787 + return_worker(&config->worker_pool, worker_data);
788 }
789 break;
790 case ACLK_DATABASE_PUSH_ALERT_CONFIG:
791 aclk_push_alert_config_event(cmd.param[0], cmd.param[1]);
792 break;
793 - case ACLK_DATABASE_PUSH_ALERT:
794 -
795 - if (config->alert_push_running)
796 - break;
797 -
798 - config->alert_push_running = true;
799 -
800 - data = mallocz(sizeof(*data));
801 - data->request.data = data;
802 - data->config = config;
803 -
804 - if (uv_queue_work(loop, &data->request, start_alert_push, after_start_alert_push)) {
805 - freez(data);
806 - config->alert_push_running = false;
807 - }
808 - break;
793 case ACLK_MQTT_WSS_CLIENT_SET:
794 config->client = (mqtt_wss_client)cmd.param[0];
795 break;
@@ -815,7 +799,7 @@ static void *aclk_synchronization_event_loop(void *arg)
799 completion_mark_complete(comp);
800 break;
801 case ACLK_QUERY_EXECUTE:
818 - query = (aclk_query_t)cmd.param[0];
802 + query = (aclk_query_t *) cmd.param[0];
803
804 bool too_busy = (config->aclk_queries_running >= query_thread_count);
805
@@ -847,7 +831,7 @@ static void *aclk_synchronization_event_loop(void *arg)
831 // We have nothing, leave
832 if (Pvalue == NULL)
833 break;
850 - aclk_query_t query_in_queue = *Pvalue;
834 + aclk_query_t *query_in_queue = *Pvalue;
835
836 // Schedule it and increase running
837 too_busy = schedule_query_in_worker(loop, config, query_in_queue);
@@ -881,7 +865,7 @@ static void *aclk_synchronization_event_loop(void *arg)
865
866 // Note: The following two opcodes must be in this order
867 case ACLK_QUERY_BATCH_ADD:
884 - query = (aclk_query_t)cmd.param[0];
868 + query = (aclk_query_t *)cmd.param[0];
869 if (!query)
870 break;
871
@@ -900,19 +884,23 @@ static void *aclk_synchronization_event_loop(void *arg)
884 if (!aclk_query_batch || config->aclk_batch_job_is_running)
885 break;
886
903 - payload = mallocz(sizeof(*payload));
904 - payload->request.data = payload;
905 - payload->config = config;
906 - payload->data = aclk_query_batch;
887 + worker_data = get_worker(&config->worker_pool);
888 + if (!worker_data) {
889 + nd_log_daemon(NDLP_ERR, "Failed to get worker for ACLK batch job");
890 + break;
891 + }
892 + worker_data->request.data = worker_data;
893 + worker_data->config = config;
894 + worker_data->payload = aclk_query_batch;
895
896 config->aclk_batch_job_is_running = true;
897 config->aclk_jobs_pending -= aclk_query_batch->count;
898 aclk_query_batch = NULL;
899
912 - if (uv_queue_work(loop, &payload->request, aclk_execute_batch, after_aclk_execute_batch)) {
913 - aclk_query_batch = payload->data;
900 + if (uv_queue_work(loop, &worker_data->request, aclk_execute_batch, after_aclk_execute_batch)) {
901 + aclk_query_batch = worker_data->payload;
902 config->aclk_jobs_pending += aclk_query_batch->count;
915 - freez(payload);
903 + return_worker(&config->worker_pool, worker_data);
904 config->aclk_batch_job_is_running = false;
905 }
906 break;
@@ -949,7 +937,7 @@ static void *aclk_synchronization_event_loop(void *arg)
937 freez(aclk_query_batch);
938 }
939
952 - aral_by_size_release(config->ar);
940 + release_cmd_pool(&config->cmd_pool);
941 completion_mark_complete(&config->start_stop_complete);
942
943 worker_unregister();
@@ -963,6 +951,8 @@ static void aclk_initialize_event_loop(void)
951 memset(&aclk_sync_config, 0, sizeof(aclk_sync_config));
952 completion_init(&aclk_sync_config.start_stop_complete);
953
954 + init_worker_pool(&aclk_sync_config.worker_pool);
955 +
956 aclk_sync_config.thread = nd_thread_create("ACLKSYNC", NETDATA_THREAD_OPTION_DEFAULT, aclk_synchronization_event_loop, &aclk_sync_config);
957 fatal_assert(NULL != aclk_sync_config.thread);
958
@@ -1051,11 +1041,11 @@ void aclk_synchronization_init(void)
1041
1042 static inline bool queue_aclk_sync_cmd(enum aclk_database_opcode opcode, const void *param0, const void *param1)
1043 {
1054 - struct aclk_database_cmd cmd;
1044 + cmd_data_t cmd;
1045 cmd.opcode = opcode;
1046 cmd.param[0] = (void *) param0;
1047 cmd.param[1] = (void *) param1;
1058 - return aclk_database_enq_cmd(&cmd);
1048 + return aclk_database_enq_cmd(&cmd, true);
1049 }
1050
1051 void aclk_synchronization_shutdown(void)
@@ -1084,7 +1074,7 @@ void aclk_push_alert_config(const char *node_id, const char *config_hash)
1074 queue_aclk_sync_cmd(ACLK_DATABASE_PUSH_ALERT_CONFIG, strdupz(node_id), strdupz(config_hash));
1075 }
1076
1087 -void aclk_execute_query(aclk_query_t query)
1077 +void aclk_execute_query(aclk_query_t *query)
1078 {
1079 if (unlikely(!query))
1080 return;
@@ -1092,7 +1082,7 @@ void aclk_execute_query(aclk_query_t query)
1082 (void) queue_aclk_sync_cmd(ACLK_QUERY_EXECUTE, query, NULL);
1083 }
1084
1095 -void aclk_add_job(aclk_query_t query)
1085 +void aclk_add_job(aclk_query_t *query)
1086 {
1087 if (unlikely(!query))
1088 return;
src/database/sqlite/sqlite_aclk.h
-10
@@ -18,16 +18,13 @@ static inline int uuid_parse_fix(char *in, nd_uuid_t uuid)
18 enum aclk_database_opcode {
19 ACLK_DATABASE_NOOP = 0,
20 ACLK_DATABASE_NODE_STATE,
21 - ACLK_DATABASE_PUSH_ALERT,
21 ACLK_DATABASE_PUSH_ALERT_CONFIG,
22 ACLK_DATABASE_NODE_UNREGISTER,
23 ACLK_MQTT_WSS_CLIENT_SET,
24 ACLK_MQTT_WSS_CLIENT_RESET,
25 ACLK_CANCEL_NODE_UPDATE_TIMER,
26 ACLK_QUEUE_NODE_INFO,
28 - ACLK_MQTT_WSS_CLIENT,
27 ACLK_QUERY_EXECUTE,
30 - ACLK_QUERY_EXECUTE_SYNC,
28 ACLK_QUERY_BATCH_ADD,
29 ACLK_QUERY_BATCH_EXECUTE,
30 ACLK_SYNC_SHUTDOWN,
@@ -37,12 +34,6 @@ enum aclk_database_opcode {
34 ACLK_MAX_ENUMERATIONS_DEFINED
35 };
36
40 -struct aclk_database_cmd {
41 - enum aclk_database_opcode opcode;
42 - void *param[2];
43 - struct aclk_database_cmd *prev, *next;
44 -};
45 -
37 typedef struct aclk_sync_cfg_t {
38 RRDHOST *host;
39 uv_timer_t timer;
@@ -64,7 +55,6 @@ void aclk_synchronization_shutdown(void);
55 void aclk_push_alert_config(const char *node_id, const char *config_hash);
56 void schedule_node_state_update(RRDHOST *host, uint64_t delay);
57 void unregister_node(const char *machine_guid);
67 -void cancel_node_update_timer(const RRDHOST *host, struct completion *completion);
58 void aclk_queue_node_info(RRDHOST *host, bool immediate);
59
60 #endif //NETDATA_SQLITE_ACLK_H
src/database/sqlite/sqlite_metadata.c
+65 -110
@@ -202,12 +202,10 @@ enum metadata_opcode {
202 METADATA_STORE_CLAIM_ID,
203 METADATA_STORE,
204 METADATA_LOAD_HOST_CONTEXT,
205 - METADATA_DELETE_HOST_CHART_LABELS,
205 METADATA_ADD_HOST_AE,
206 METADATA_DEL_HOST_AE,
207 METADATA_ADD_CTX_CLEANUP,
208 METADATA_EXECUTE_STORE_STATEMENT,
210 - METADATA_MAINTENANCE,
209 METADATA_SYNC_SHUTDOWN,
210 METADATA_UNITTEST,
211 // leave this last
@@ -215,13 +213,6 @@ enum metadata_opcode {
213 METADATA_MAX_ENUMERATIONS_DEFINED
214 };
215
218 -#define MAX_PARAM_LIST (2)
219 -struct metadata_cmd {
220 - enum metadata_opcode opcode;
221 - const void *param[MAX_PARAM_LIST];
222 - struct metadata_cmd *prev, *next;
223 -};
224 -
216 struct meta_config_s {
217 ND_THREAD *thread;
218 uv_loop_t loop;
@@ -232,12 +223,11 @@ struct meta_config_s {
223 bool initialized;
224 bool ctx_load_running;
225 bool metadata_running;
226 + bool store_metadata;
227 bool shutdown_requested;
236 - SPINLOCK cmd_queue_lock;
228 struct completion start_stop_complete;
238 - /* FIFO command queue */
239 - struct metadata_cmd *cmd_base;
240 - ARAL *ar;
229 + CmdPool cmd_pool;
230 + WorkerPool worker_pool;
231 } meta_config;
232
233 //
@@ -1519,51 +1509,23 @@ static void cleanup_health_log(struct meta_config_s *config)
1509 // EVENT LOOP STARTS HERE
1510 //
1511
1522 -static void metadata_free_cmd_queue(struct meta_config_s *wc)
1523 -{
1524 - spinlock_lock(&wc->cmd_queue_lock);
1525 - while(wc->cmd_base) {
1526 - struct metadata_cmd *t = wc->cmd_base;
1527 - DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wc->cmd_base, t, prev, next);
1528 - aral_freez(wc->ar, t);
1529 - }
1530 - spinlock_unlock(&wc->cmd_queue_lock);
1531 -}
1512
1533 -static bool metadata_enq_cmd(struct metadata_cmd *cmd)
1513 +static bool metadata_enq_cmd(cmd_data_t *cmd, bool wait_on_full)
1514 {
1515 if(unlikely(!__atomic_load_n(&meta_config.initialized, __ATOMIC_RELAXED)))
1516 return false;
1517
1538 - struct metadata_cmd *t = aral_mallocz(meta_config.ar);
1539 - *t = *cmd;
1540 - t->prev = t->next = NULL;
1541 -
1542 - spinlock_lock(&meta_config.cmd_queue_lock);
1543 - DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(meta_config.cmd_base, t, prev, next);
1544 - spinlock_unlock(&meta_config.cmd_queue_lock);
1545 -
1546 - (void) uv_async_send(&meta_config.async);
1547 - return true;
1518 + bool added = push_cmd(&meta_config.cmd_pool, (void *)cmd, wait_on_full);
1519 + if (added)
1520 + (void) uv_async_send(&meta_config.async);
1521 + return added;
1522 }
1523
1550 -static struct metadata_cmd metadata_deq_cmd(struct meta_config_s *wc)
1524 +static cmd_data_t metadata_deq_cmd()
1525 {
1552 - struct metadata_cmd ret, *to_free = NULL;
1553 -
1554 - spinlock_lock(&wc->cmd_queue_lock);
1555 - if(wc->cmd_base) {
1556 - struct metadata_cmd *t = wc->cmd_base;
1557 - DOUBLE_LINKED_LIST_REMOVE_ITEM_UNSAFE(wc->cmd_base, t, prev, next);
1558 - ret = *t;
1559 - to_free = t;
1560 - }
1561 - else
1562 - ret.opcode = METADATA_DATABASE_NOOP;
1563 - spinlock_unlock(&wc->cmd_queue_lock);
1564 -
1565 - aral_freez(wc->ar, to_free);
1566 -
1526 + cmd_data_t ret;
1527 + ret.opcode = METADATA_DATABASE_NOOP;
1528 + (void) pop_cmd(&meta_config.cmd_pool, (cmd_data_t *) &ret);
1529 return ret;
1530 }
1531
@@ -1576,19 +1538,14 @@ static void async_cb(uv_async_t *handle)
1538 #define TIMER_INITIAL_PERIOD_MS (1000)
1539 #define TIMER_REPEAT_PERIOD_MS (1000)
1540
1579 -static void timer_cb(uv_timer_t* handle)
1541 +static void timer_cb(uv_timer_t *handle)
1542 {
1543 uv_stop(handle->loop);
1544 uv_update_time(handle->loop);
1545
1584 - struct meta_config_s *wc = handle->data;
1585 -
1586 - if (wc->metadata_check_after < now_realtime_sec()) {
1587 - struct metadata_cmd cmd;
1588 - memset(&cmd, 0, sizeof(cmd));
1589 - cmd.opcode = METADATA_STORE;
1590 - (void) metadata_enq_cmd(&cmd);
1591 - }
1546 + struct meta_config_s *config = handle->data;
1547 + if (config->metadata_check_after < now_realtime_sec())
1548 + config->store_metadata = true;
1549 }
1550
1551 void vacuum_database(sqlite3 *database, const char *db_alias, int threshold, int vacuum_pc)
@@ -1761,16 +1718,6 @@ void run_metadata_cleanup(struct meta_config_s *config)
1718 (void) sqlite3_wal_checkpoint(db_meta, NULL);
1719 }
1720
1764 -struct work_payload {
1765 - uv_work_t request;
1766 - struct meta_config_s *config;
1767 - void *pending_alert_list;
1768 - void *pending_ctx_cleanup_list;
1769 - void *pending_uuid_deletion;
1770 - void *pending_sql_statement;
1771 - BUFFER *work_buffer;
1772 -};
1773 -
1721 struct host_context_load_thread {
1722 ND_THREAD *thread;
1723 RRDHOST *host;
@@ -1843,10 +1790,10 @@ static void *restore_host_context(void *arg)
1790 // Callback after scan of hosts is done
1791 static void after_ctx_hosts_load(uv_work_t *req, int status __maybe_unused)
1792 {
1846 - struct work_payload *data = req->data;
1847 - struct meta_config_s *config = data->config;
1793 + worker_data_t *worker = req->data;
1794 + struct meta_config_s *config = worker->config;
1795 config->ctx_load_running = false;
1849 - freez(data);
1796 + return_worker(&config->worker_pool, worker);
1797 }
1798
1799 static bool cleanup_finished_threads(struct host_context_load_thread *hclt, size_t max_thread_slots, bool wait, size_t *free_slot)
@@ -1900,7 +1847,7 @@ static void ctx_hosts_load(uv_work_t *req)
1847 {
1848 register_libuv_worker_jobs();
1849
1903 - struct work_payload *data = req->data;
1850 + worker_data_t *data = req->data;
1851 struct meta_config_s *config = data->config;
1852
1853 worker_is_busy(UV_EVENT_HOST_CONTEXT_LOAD);
@@ -1989,8 +1936,8 @@ static void ctx_hosts_load(uv_work_t *req)
1936 // Callback after scan of hosts is done
1937 static void after_metadata_hosts(uv_work_t *req, int status __maybe_unused)
1938 {
1992 - struct work_payload *data = req->data;
1993 - struct meta_config_s *config = data->config;
1939 + worker_data_t *worker = req->data;
1940 + struct meta_config_s *config = worker->config;
1941
1942 bool first = true;
1943 Word_t Index = 0;
@@ -2006,7 +1953,7 @@ static void after_metadata_hosts(uv_work_t *req, int status __maybe_unused)
1953 }
1954
1955 config->metadata_running = false;
2009 - freez(data);
1956 + return_worker(&config->worker_pool, worker);
1957 }
1958
1959 #ifdef ENABLE_DBENGINE
@@ -2447,7 +2394,7 @@ static void start_metadata_hosts(uv_work_t *req)
2394 {
2395 register_libuv_worker_jobs();
2396
2450 - struct work_payload *data = req->data;
2397 + worker_data_t *data = req->data;
2398 struct meta_config_s *config = data->config;
2399
2400 BUFFER *work_buffer = data->work_buffer;
@@ -2480,6 +2427,7 @@ static void start_metadata_hosts(uv_work_t *req)
2427
2428 #define MAX_SHUTDOWN_TIMEOUT_SECONDS (10)
2429 #define SHUTDOWN_SLEEP_INTERVAL_MS (100)
2430 +#define CMD_POOL_SIZE (32768)
2431
2432 static void *metadata_event_loop(void *arg)
2433 {
@@ -2488,7 +2436,7 @@ static void *metadata_event_loop(void *arg)
2436 service_register(NULL, NULL, NULL);
2437 worker_register(EVENT_LOOP_NAME);
2438
2491 - config->ar = aral_by_size_acquire(sizeof(struct metadata_cmd));
2439 + init_cmd_pool(&config->cmd_pool, CMD_POOL_SIZE);
2440
2441 worker_register_job_name(METADATA_DATABASE_NOOP, "noop");
2442 worker_register_job_name(METADATA_DEL_DIMENSION, "delete dimension");
@@ -2513,7 +2461,7 @@ static void *metadata_event_loop(void *arg)
2461 config->metadata_check_after = now_realtime_sec() + METADATA_HOST_CHECK_FIRST_CHECK;
2462
2463 BUFFER *work_buffer = buffer_create(1024, &netdata_buffers_statistics.buffers_sqlite);
2516 - struct work_payload *data;
2464 + worker_data_t *worker;
2465 Pvoid_t *Pvalue;
2466 struct judy_list_t *pending_alert_list = NULL;
2467 struct judy_list_t *pending_ctx_cleanup_list = NULL;
@@ -2540,8 +2488,15 @@ static void *metadata_event_loop(void *arg)
2488 if (unlikely(++cmd_batch_size >= METADATA_MAX_BATCH_SIZE))
2489 break;
2490
2543 - struct metadata_cmd cmd = metadata_deq_cmd(config);
2544 - opcode = cmd.opcode;
2491 + cmd_data_t cmd;
2492 + if (config->store_metadata && !config->metadata_running) {
2493 + config->store_metadata = false;
2494 + opcode = METADATA_STORE;
2495 + }
2496 + else {
2497 + cmd = metadata_deq_cmd();
2498 + opcode = cmd.opcode;
2499 + }
2500
2501 if (likely(opcode != METADATA_DATABASE_NOOP))
2502 worker_is_busy(opcode);
@@ -2588,42 +2543,42 @@ static void *metadata_event_loop(void *arg)
2543 if (config->metadata_running || unittest_running)
2544 break;
2545
2591 - data = mallocz(sizeof(*data));
2592 - data->request.data = data;
2593 - data->config = config;
2594 - data->pending_alert_list = pending_alert_list;
2595 - data->pending_ctx_cleanup_list = pending_ctx_cleanup_list;
2596 - data->pending_uuid_deletion = pending_uuid_deletion;
2597 - data->pending_sql_statement = pending_sql_statement;
2546 + worker = get_worker(&config->worker_pool);
2547 + worker->request.data = worker;
2548 + worker->config = config;
2549 + worker->pending_alert_list = pending_alert_list;
2550 + worker->pending_ctx_cleanup_list = pending_ctx_cleanup_list;
2551 + worker->pending_uuid_deletion = pending_uuid_deletion;
2552 + worker->pending_sql_statement = pending_sql_statement;
2553
2599 - data->work_buffer = work_buffer;
2554 + worker->work_buffer = work_buffer;
2555 pending_alert_list = NULL;
2556 pending_ctx_cleanup_list = NULL;
2557 pending_uuid_deletion = NULL;
2558 pending_sql_statement = NULL;
2559 config->metadata_running = true;
2605 - if (uv_queue_work(loop, &data->request, start_metadata_hosts, after_metadata_hosts)) {
2606 - pending_alert_list = data->pending_alert_list;
2607 - pending_ctx_cleanup_list = data->pending_ctx_cleanup_list;
2608 - pending_uuid_deletion = data->pending_uuid_deletion;
2609 - pending_sql_statement = data->pending_sql_statement;
2610 - freez(data);
2560 + if (uv_queue_work(loop, &worker->request, start_metadata_hosts, after_metadata_hosts)) {
2561 + pending_alert_list = worker->pending_alert_list;
2562 + pending_ctx_cleanup_list = worker->pending_ctx_cleanup_list;
2563 + pending_uuid_deletion = worker->pending_uuid_deletion;
2564 + pending_sql_statement = worker->pending_sql_statement;
2565 config->metadata_running = false;
2566 + return_worker(&config->worker_pool, worker);
2567 }
2568 break;
2569 case METADATA_LOAD_HOST_CONTEXT:
2570 if (config->ctx_load_running || unittest_running)
2571 break;
2572
2573 + worker = get_worker(&config->worker_pool);
2574 config->ctx_load_running = true;
2619 - data = callocz(1, sizeof(*data));
2620 - data->request.data = data;
2621 - data->config = config;
2622 - if (uv_queue_work(loop, &data->request, ctx_hosts_load, after_ctx_hosts_load)) {
2623 - freez(data);
2575 + worker->request.data = worker;
2576 + worker->config = config;
2577 + if (uv_queue_work(loop, &worker->request, ctx_hosts_load, after_ctx_hosts_load)) {
2578 config->ctx_load_running = false;
2579 // Fallback reset context so hosts will load on demand
2580 reset_host_context_load_flag();
2581 + return_worker(&config->worker_pool, worker);
2582 }
2583 break;
2584 case METADATA_ADD_HOST_AE:
@@ -2710,8 +2665,7 @@ static void *metadata_event_loop(void *arg)
2665 }
2666
2667 buffer_free(work_buffer);
2713 - metadata_free_cmd_queue(config);
2714 - aral_by_size_release(config->ar);
2668 + release_cmd_pool(&config->cmd_pool);
2669 worker_unregister();
2670
2671 completion_mark_complete(&config->start_stop_complete);
@@ -2721,14 +2675,14 @@ static void *metadata_event_loop(void *arg)
2675
2676 void metadata_sync_shutdown(void)
2677 {
2724 - struct metadata_cmd cmd;
2678 + cmd_data_t cmd;
2679 memset(&cmd, 0, sizeof(cmd));
2680 cmd.opcode = METADATA_SYNC_SHUTDOWN;
2681
2682 // if we can't sent command return
2683 // This should not happen but if we wait we may not get a completion
2684 // and shutdown will timeout
2731 - if (!metadata_enq_cmd(&cmd)) {
2685 + if (!metadata_enq_cmd(&cmd, true)) {
2686 nd_log_daemon(NDLP_WARNING, "METADATA: Failed to send a shutdown command");
2687 return;
2688 }
@@ -2752,6 +2706,7 @@ void metadata_sync_init(void)
2706 memset(&meta_config, 0, sizeof(meta_config));
2707 completion_init(&meta_config.start_stop_complete);
2708
2709 + init_worker_pool(&meta_config.worker_pool);
2710 meta_config.thread = nd_thread_create("METASYNC", NETDATA_THREAD_OPTION_DEFAULT, metadata_event_loop, &meta_config);
2711 fatal_assert(NULL != meta_config.thread);
2712
@@ -2764,13 +2719,13 @@ void metadata_sync_init(void)
2719
2720 // Helpers
2721
2767 -static inline bool queue_metadata_cmd(enum metadata_opcode opcode, const void *param0, const void *param1)
2722 +static inline bool queue_metadata_cmd(enum metadata_opcode opcode, void *param0, void *param1)
2723 {
2769 - struct metadata_cmd cmd;
2724 + cmd_data_t cmd;
2725 cmd.opcode = opcode;
2726 cmd.param[0] = param0;
2727 cmd.param[1] = param1;
2773 - return metadata_enq_cmd(&cmd);
2728 + return metadata_enq_cmd(&cmd, true);
2729 }
2730
2731 // Public
@@ -2953,15 +2908,15 @@ void get_agent_event_time_median_init(void) {
2908 static void *unittest_queue_metadata(void *arg) {
2909 struct thread_unittest *tu = arg;
2910
2956 - struct metadata_cmd cmd;
2911 + cmd_data_t cmd;
2912 cmd.opcode = METADATA_UNITTEST;
2913 cmd.param[0] = tu;
2914 cmd.param[1] = NULL;
2960 - metadata_enq_cmd(&cmd);
2915 + metadata_enq_cmd(&cmd, true);
2916
2917 do {
2918 __atomic_fetch_add(&tu->added, 1, __ATOMIC_SEQ_CST);
2964 - metadata_enq_cmd(&cmd);
2919 + metadata_enq_cmd(&cmd, true);
2920 sleep_usec(10000);
2921 } while (!__atomic_load_n(&tu->join, __ATOMIC_RELAXED));
2922 return arg;