Query queue only for queries (#13431)
simplify and clean up
Timotej S committed
Jul 28, 2022 at 10:29 UTC
221fd512873613a10b3d95b25a8a4d542b2c4801
14 files changed
+77
-253
aclk/aclk.c
+15
-23
@@ -742,8 +742,8 @@ void aclk_host_state_update(RRDHOST *host, int cmd)
742
}
743
if (ret < 0) {
744
// node_id not found
745
- aclk_query_t create_query;
746
- create_query = aclk_query_new(REGISTER_NODE);
745
+ size_t payload_len;
746
+
747
rrdhost_aclk_state_lock(localhost);
748
node_instance_creation_t node_instance_creation = {
749
.claim_id = localhost->aclk_state.claimed_id,
@@ -751,16 +751,14 @@ void aclk_host_state_update(RRDHOST *host, int cmd)
751
.hostname = host->hostname,
752
.machine_guid = host->machine_guid
753
};
754
- create_query->data.bin_payload.payload = generate_node_instance_creation(&create_query->data.bin_payload.size, &node_instance_creation);
754
+ char *payload = generate_node_instance_creation(&payload_len, &node_instance_creation);
755
rrdhost_aclk_state_unlock(localhost);
756
- create_query->data.bin_payload.topic = ACLK_TOPICID_CREATE_NODE;
757
- create_query->data.bin_payload.msg_name = "CreateNodeInstance";
756
+
757
info("Registering host=%s, hops=%u",host->machine_guid, host->system_info->hops);
759
- aclk_queue_query(create_query);
758
+ aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_CREATE_NODE, "CreateNodeInstance");
759
return;
760
}
761
763
- aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
762
node_instance_connection_t node_state_update = {
763
.hops = host->system_info->hops,
764
.live = cmd,
@@ -781,15 +779,14 @@ void aclk_host_state_update(RRDHOST *host, int cmd)
779
780
rrdhost_aclk_state_lock(localhost);
781
node_state_update.claim_id = localhost->aclk_state.claimed_id;
784
- query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
782
+ size_t payload_len;
783
+ char *payload = generate_node_instance_connection(&payload_len, &node_state_update);
784
rrdhost_aclk_state_unlock(localhost);
785
786
info("Queuing status update for node=%s, live=%d, hops=%u",(char*)node_state_update.node_id, cmd,
787
host->system_info->hops);
788
freez((void*)node_state_update.node_id);
790
- query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
791
- query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
792
- aclk_queue_query(query);
789
+ aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_NODE_CONN, "UpdateNodeInstanceConnection");
790
}
791
792
void aclk_send_node_instances()
@@ -802,7 +799,6 @@ void aclk_send_node_instances()
799
}
800
while (!uuid_is_null(list->host_id)) {
801
if (!uuid_is_null(list->node_id)) {
805
- aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
802
node_instance_connection_t node_state_update = {
803
.live = list->live,
804
.hops = list->hops,
@@ -827,34 +823,30 @@ void aclk_send_node_instances()
823
824
rrdhost_aclk_state_lock(localhost);
825
node_state_update.claim_id = localhost->aclk_state.claimed_id;
830
- query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
826
+ size_t payload_len;
827
+ char *payload = generate_node_instance_connection(&payload_len, &node_state_update);
828
rrdhost_aclk_state_unlock(localhost);
829
info("Queuing status update for node=%s, live=%d, hops=%d",(char*)node_state_update.node_id,
830
list->live,
831
list->hops);
832
freez((void*)node_state_update.node_id);
836
- query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
837
- query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
838
- aclk_queue_query(query);
833
+ aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_NODE_CONN, "UpdateNodeInstanceConnection");
834
} else {
840
- aclk_query_t create_query;
841
- create_query = aclk_query_new(REGISTER_NODE);
835
node_instance_creation_t node_instance_creation = {
836
.hops = list->hops,
837
.hostname = list->hostname,
838
};
839
node_instance_creation.machine_guid = mallocz(UUID_STR_LEN);
840
uuid_unparse_lower(list->host_id, (char*)node_instance_creation.machine_guid);
848
- create_query->data.bin_payload.topic = ACLK_TOPICID_CREATE_NODE;
849
- create_query->data.bin_payload.msg_name = "CreateNodeInstance";
841
rrdhost_aclk_state_lock(localhost);
851
- node_instance_creation.claim_id = localhost->aclk_state.claimed_id,
852
- create_query->data.bin_payload.payload = generate_node_instance_creation(&create_query->data.bin_payload.size, &node_instance_creation);
842
+ node_instance_creation.claim_id = localhost->aclk_state.claimed_id;
843
+ size_t payload_len;
844
+ char *payload = generate_node_instance_creation(&payload_len, &node_instance_creation);
845
rrdhost_aclk_state_unlock(localhost);
846
info("Queuing registration for host=%s, hops=%d",(char*)node_instance_creation.machine_guid,
847
list->hops);
848
freez(node_instance_creation.machine_guid);
857
- aclk_queue_query(create_query);
849
+ aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_CREATE_NODE, "CreateNodeInstance");
850
}
851
freez(list->hostname);
852
aclk/aclk.h
+11
@@ -34,6 +34,17 @@ void aclk_send_node_instances(void);
34
35
void aclk_send_bin_msg(char *msg, size_t msg_len, enum aclk_topics subtopic, const char *msgname);
36
37
+#define GENERATE_AND_SEND_PAYLOAD(topic, msg_name, generator_fnc, generator_data...) \
38
+ size_t payload_len; \
39
+ char *payload = generator_fnc(&payload_len, generator_data); \
40
+ if (unlikely(payload == NULL)) { \
41
+ error("Failed to generate payload (%s)", __FUNCTION__); \
42
+ return; \
43
+ } \
44
+ aclk_send_bin_msg(payload, payload_len, topic, msg_name); \
45
+ if (!use_mqtt_5) \
46
+ freez(payload);
47
+
48
char *ng_aclk_state(void);
49
char *ng_aclk_state_json(void);
50
aclk/aclk_alarm_api.c
+4
-22
@@ -10,38 +10,20 @@
10
11
void aclk_send_alarm_log_health(struct alarm_log_health *log_health)
12
{
13
- aclk_query_t query = aclk_query_new(ALARM_LOG_HEALTH);
14
- query->data.bin_payload.payload = generate_alarm_log_health(&query->data.bin_payload.size, log_health);
15
- query->data.bin_payload.topic = ACLK_TOPICID_ALARM_HEALTH;
16
- query->data.bin_payload.msg_name = "AlarmLogHealth";
17
- QUEUE_IF_PAYLOAD_PRESENT(query);
13
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_HEALTH, "AlarmLogHealth", generate_alarm_log_health, log_health);
14
}
15
16
void aclk_send_alarm_log_entry(struct alarm_log_entry *log_entry)
17
{
22
- size_t payload_size;
23
- char *payload = generate_alarm_log_entry(&payload_size, log_entry);
24
-
25
- aclk_send_bin_msg(payload, payload_size, ACLK_TOPICID_ALARM_LOG, "AlarmLogEntry");
26
-
27
- if (!use_mqtt_5)
28
- freez(payload);
18
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_LOG, "AlarmLogEntry", generate_alarm_log_entry, log_entry);
19
}
20
21
void aclk_send_provide_alarm_cfg(struct provide_alarm_configuration *cfg)
22
{
33
- aclk_query_t query = aclk_query_new(ALARM_PROVIDE_CFG);
34
- query->data.bin_payload.payload = generate_provide_alarm_configuration(&query->data.bin_payload.size, cfg);
35
- query->data.bin_payload.topic = ACLK_TOPICID_ALARM_CONFIG;
36
- query->data.bin_payload.msg_name = "ProvideAlarmConfiguration";
37
- QUEUE_IF_PAYLOAD_PRESENT(query);
23
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_CONFIG, "ProvideAlarmConfiguration", generate_provide_alarm_configuration, cfg);
24
}
25
26
void aclk_send_alarm_snapshot(alarm_snapshot_proto_ptr_t snapshot)
27
{
42
- aclk_query_t query = aclk_query_new(ALARM_SNAPSHOT);
43
- query->data.bin_payload.payload = generate_alarm_snapshot_bin(&query->data.bin_payload.size, snapshot);
44
- query->data.bin_payload.topic = ACLK_TOPICID_ALARM_SNAPSHOT;
45
- query->data.bin_payload.msg_name = "AlarmSnapshot";
46
- QUEUE_IF_PAYLOAD_PRESENT(query);
28
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_SNAPSHOT, "AlarmSnapshot", generate_alarm_snapshot_bin, snapshot);
29
}
aclk/aclk_charts_api.c
+10
-39
@@ -3,75 +3,46 @@
3
4
#include "aclk_query_queue.h"
5
6
+#include "aclk.h"
7
+
8
#define CHART_DIM_UPDATE_NAME "ChartsAndDimensionsUpdated"
9
10
void aclk_chart_inst_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
11
{
10
- aclk_query_t query = aclk_query_new(CHART_DIMS_UPDATE);
11
- query->data.bin_payload.payload = generate_charts_updated(&query->data.bin_payload.size, payloads, payload_sizes, new_positions);
12
- query->data.bin_payload.msg_name = CHART_DIM_UPDATE_NAME;
13
- QUEUE_IF_PAYLOAD_PRESENT(query);
12
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_DIMS, CHART_DIM_UPDATE_NAME, generate_charts_updated, payloads, payload_sizes, new_positions);
13
}
14
15
void aclk_chart_dim_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
16
{
18
- aclk_query_t query = aclk_query_new(CHART_DIMS_UPDATE);
19
- query->data.bin_payload.topic = ACLK_TOPICID_CHART_DIMS;
20
- query->data.bin_payload.payload = generate_chart_dimensions_updated(&query->data.bin_payload.size, payloads, payload_sizes, new_positions);
21
- query->data.bin_payload.msg_name = CHART_DIM_UPDATE_NAME;
22
- QUEUE_IF_PAYLOAD_PRESENT(query);
17
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_DIMS, CHART_DIM_UPDATE_NAME, generate_chart_dimensions_updated, payloads, payload_sizes, new_positions);
18
}
19
20
void aclk_chart_inst_and_dim_update(char **payloads, size_t *payload_sizes, int *is_dim, struct aclk_message_position *new_positions, uint64_t batch_id)
21
{
27
- aclk_query_t query = aclk_query_new(CHART_DIMS_UPDATE);
28
- query->data.bin_payload.topic = ACLK_TOPICID_CHART_DIMS;
29
- query->data.bin_payload.payload = generate_charts_and_dimensions_updated(&query->data.bin_payload.size, payloads, payload_sizes, is_dim, new_positions, batch_id);
30
- query->data.bin_payload.msg_name = CHART_DIM_UPDATE_NAME;
31
- QUEUE_IF_PAYLOAD_PRESENT(query);
22
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_DIMS, CHART_DIM_UPDATE_NAME, generate_charts_and_dimensions_updated, payloads, payload_sizes, is_dim, new_positions, batch_id);
23
}
24
25
void aclk_chart_config_updated(struct chart_config_updated *config_list, int list_size)
26
{
36
- aclk_query_t query = aclk_query_new(CHART_CONFIG_UPDATED);
37
- query->data.bin_payload.topic = ACLK_TOPICID_CHART_CONFIGS_UPDATED;
38
- query->data.bin_payload.payload = generate_chart_configs_updated(&query->data.bin_payload.size, config_list, list_size);
39
- query->data.bin_payload.msg_name = "ChartConfigsUpdated";
40
- QUEUE_IF_PAYLOAD_PRESENT(query);
27
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_CONFIGS_UPDATED, "ChartConfigsUpdated", generate_chart_configs_updated, config_list, list_size);
28
}
29
30
void aclk_chart_reset(chart_reset_t reset)
31
{
45
- aclk_query_t query = aclk_query_new(CHART_RESET);
46
- query->data.bin_payload.topic = ACLK_TOPICID_CHART_RESET;
47
- query->data.bin_payload.payload = generate_reset_chart_messages(&query->data.bin_payload.size, reset);
48
- query->data.bin_payload.msg_name = "ResetChartMessages";
49
- QUEUE_IF_PAYLOAD_PRESENT(query);
32
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_RESET, "ResetChartMessages", generate_reset_chart_messages, reset);
33
}
34
35
void aclk_retention_updated(struct retention_updated *data)
36
{
54
- aclk_query_t query = aclk_query_new(RETENTION_UPDATED);
55
- query->data.bin_payload.topic = ACLK_TOPICID_RETENTION_UPDATED;
56
- query->data.bin_payload.payload = generate_retention_updated(&query->data.bin_payload.size, data);
57
- query->data.bin_payload.msg_name = "RetentionUpdated";
58
- QUEUE_IF_PAYLOAD_PRESENT(query);
37
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_RETENTION_UPDATED, "RetentionUpdated", generate_retention_updated, data);
38
}
39
40
void aclk_update_node_info(struct update_node_info *info)
41
{
63
- aclk_query_t query = aclk_query_new(UPDATE_NODE_INFO);
64
- query->data.bin_payload.topic = ACLK_TOPICID_NODE_INFO;
65
- query->data.bin_payload.payload = generate_update_node_info_message(&query->data.bin_payload.size, info);
66
- query->data.bin_payload.msg_name = "UpdateNodeInfo";
67
- QUEUE_IF_PAYLOAD_PRESENT(query);
42
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_NODE_INFO, "UpdateNodeInfo", generate_update_node_info_message, info);
43
}
44
45
void aclk_update_node_collectors(struct update_node_collectors *collectors)
46
{
72
- aclk_query_t query = aclk_query_new(UPDATE_NODE_COLLECTORS);
73
- query->data.bin_payload.topic = ACLK_TOPICID_NODE_COLLECTORS;
74
- query->data.bin_payload.payload = generate_update_node_collectors_message(&query->data.bin_payload.size, collectors);
75
- query->data.bin_payload.msg_name = "UpdateNodeCollectors";
76
- QUEUE_IF_PAYLOAD_PRESENT(query);
47
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_NODE_COLLECTORS, "UpdateNodeCollectors", generate_update_node_collectors_message, collectors);
48
}
aclk/aclk_contexts_api.c
+4
-10
@@ -4,20 +4,14 @@
4
5
#include "aclk_contexts_api.h"
6
7
+#include "aclk.h"
8
+
9
void aclk_send_contexts_snapshot(contexts_snapshot_t data)
10
{
9
- aclk_query_t query = aclk_query_new(PROTO_BIN_MESSAGE);
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";
13
- QUEUE_IF_PAYLOAD_PRESENT(query);
11
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CTXS_SNAPSHOT, "ContextsSnapshot", contexts_snapshot_2bin, data);
12
}
13
14
void aclk_send_contexts_updated(contexts_updated_t data)
15
{
18
- aclk_query_t query = aclk_query_new(PROTO_BIN_MESSAGE);
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";
22
- QUEUE_IF_PAYLOAD_PRESENT(query);
16
+ GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CTXS_UPDATED, "ContextsUpdated", contexts_updated_2bin, data);
17
}
aclk/aclk_query.c
+13
-55
@@ -84,7 +84,7 @@ static int http_api_v2(struct aclk_query_thread *query_thr, aclk_query_t query)
84
w->cookie2[0] = 0; // Simulate web_client_create_on_fd()
85
w->acl = 0x1f;
86
87
- buffer_strcat(log_buffer, query->data.http_api_v2.query);
87
+ buffer_strcat(log_buffer, query->http_api_v2.query);
88
size_t size = 0;
89
size_t sent = 0;
90
w->tv_in = query->created_tv;
@@ -102,8 +102,8 @@ static int http_api_v2(struct aclk_query_thread *query_thr, aclk_query_t query)
102
}
103
104
RRDHOST *temp_host = NULL;
105
- if (!strncmp(query->data.http_api_v2.query, NODE_ID_QUERY, strlen(NODE_ID_QUERY))) {
106
- char *node_uuid = query->data.http_api_v2.query + strlen(NODE_ID_QUERY);
105
+ if (!strncmp(query->http_api_v2.query, NODE_ID_QUERY, strlen(NODE_ID_QUERY))) {
106
+ char *node_uuid = query->http_api_v2.query + strlen(NODE_ID_QUERY);
107
char nodeid[UUID_STR_LEN];
108
if (strlen(node_uuid) < (UUID_STR_LEN - 1)) {
109
error_report(CLOUD_EMSG_MALFORMED_NODE_ID);
@@ -127,14 +127,14 @@ static int http_api_v2(struct aclk_query_thread *query_thr, aclk_query_t query)
127
}
128
}
129
130
- char *mysep = strchr(query->data.http_api_v2.query, '?');
130
+ char *mysep = strchr(query->http_api_v2.query, '?');
131
if (mysep) {
132
url_decode_r(w->decoded_query_string, mysep, NETDATA_WEB_REQUEST_URL_SIZE + 1);
133
*mysep = '\0';
134
} else
135
- url_decode_r(w->decoded_query_string, query->data.http_api_v2.query, NETDATA_WEB_REQUEST_URL_SIZE + 1);
135
+ url_decode_r(w->decoded_query_string, query->http_api_v2.query, NETDATA_WEB_REQUEST_URL_SIZE + 1);
136
137
- mysep = strrchr(query->data.http_api_v2.query, '/');
137
+ mysep = strrchr(query->http_api_v2.query, '/');
138
139
if (aclk_stats_enabled) {
140
ACLK_STATS_LOCK;
@@ -151,7 +151,7 @@ static int http_api_v2(struct aclk_query_thread *query_thr, aclk_query_t query)
151
152
#ifdef NETDATA_WITH_ZLIB
153
// check if gzip encoding can and should be used
154
- if ((start = strstr((char *)query->data.http_api_v2.payload, WEB_HDR_ACCEPT_ENC))) {
154
+ if ((start = strstr((char *)query->http_api_v2.payload, WEB_HDR_ACCEPT_ENC))) {
155
start += strlen(WEB_HDR_ACCEPT_ENC);
156
end = strstr(start, "\x0D\x0A");
157
start = strstr(start, "gzip");
@@ -256,57 +256,16 @@ cleanup:
256
return retval;
257
}
258
259
-static int send_bin_msg(struct aclk_query_thread *query_thr, aclk_query_t query)
260
-{
261
- // this will be simplified when legacy support is removed
262
- aclk_send_bin_message_subtopic_pid(query_thr->client, query->data.bin_payload.payload, query->data.bin_payload.size, query->data.bin_payload.topic, query->data.bin_payload.msg_name);
263
- return 0;
264
-}
265
-
266
-const char *aclk_query_get_name(aclk_query_type_t qt)
267
-{
268
- switch (qt) {
269
- case HTTP_API_V2: return "http_api_request_v2";
270
- case REGISTER_NODE: return "register_node";
271
- case NODE_STATE_UPDATE: return "node_state_update";
272
- case CHART_DIMS_UPDATE: return "chart_and_dim_update";
273
- case CHART_CONFIG_UPDATED: return "chart_config_updated";
274
- case CHART_RESET: return "reset_chart_messages";
275
- case RETENTION_UPDATED: return "update_retention_info";
276
- case UPDATE_NODE_INFO: return "update_node_info";
277
- case ALARM_LOG_HEALTH: return "alarm_log_health";
278
- case ALARM_PROVIDE_CFG: return "provide_alarm_config";
279
- case ALARM_SNAPSHOT: return "alarm_snapshot";
280
- case UPDATE_NODE_COLLECTORS: return "update_node_collectors";
281
- case PROTO_BIN_MESSAGE: return "generic_binary_proto_message";
282
- default:
283
- error_report("Unknown query type used %d", (int) qt);
284
- return "unknown";
285
- }
286
-}
287
-
259
static void aclk_query_process_msg(struct aclk_query_thread *query_thr, aclk_query_t query)
260
{
290
- if (query->type == UNKNOWN || query->type >= ACLK_QUERY_TYPE_COUNT) {
291
- error_report("Unknown query in query queue. %u", query->type);
292
- aclk_query_free(query);
293
- return;
294
- }
295
-
296
- worker_is_busy(query->type);
297
- if (query->type == HTTP_API_V2) {
298
- debug(D_ACLK, "Processing Queued Message of type: \"http_api_request_v2\"");
299
- http_api_v2(query_thr, query);
300
- } else {
301
- debug(D_ACLK, "Processing Queued Message of type: \"%s\"", query->data.bin_payload.msg_name);
302
- send_bin_msg(query_thr, query);
303
- }
261
+ worker_is_busy(0);
262
+ debug(D_ACLK, "Processing Queued Message of type: \"http_api_request_v2\"");
263
+ http_api_v2(query_thr, query);
264
265
if (aclk_stats_enabled) {
266
ACLK_STATS_LOCK;
267
aclk_metrics_per_sample.queries_dispatched++;
268
aclk_queries_per_thread[query_thr->idx]++;
309
- aclk_metrics_per_sample.queries_per_type[query->type]++;
269
ACLK_STATS_UNLOCK;
270
}
271
@@ -326,11 +285,10 @@ int aclk_query_process_msgs(struct aclk_query_thread *query_thr)
285
return 0;
286
}
287
329
-static void worker_aclk_register(void) {
288
+static void worker_aclk_register(void)
289
+{
290
worker_register("ACLKQUERY");
331
- for (int i = 1; i < ACLK_QUERY_TYPE_COUNT; i++) {
332
- worker_register_job_name(i, aclk_query_get_name(i));
333
- }
291
+ worker_register_job_name(0, "http query");
292
}
293
294
/**
aclk/aclk_query.h
-2
@@ -31,6 +31,4 @@ struct aclk_query_threads {
31
void aclk_query_threads_start(struct aclk_query_threads *query_threads, mqtt_wss_client client);
32
void aclk_query_threads_cleanup(struct aclk_query_threads *query_threads);
33
34
-const char *aclk_query_get_name(aclk_query_type_t qt);
35
-
34
#endif //NETDATA_AGENT_CLOUD_LINK_H
aclk/aclk_query_queue.c
+5
-30
@@ -95,41 +95,16 @@ void aclk_queue_flush(void)
95
};
96
}
97
98
-aclk_query_t aclk_query_new(aclk_query_type_t type)
98
+aclk_query_t aclk_query_new()
99
{
100
- aclk_query_t query = callocz(1, sizeof(struct aclk_query));
101
- query->type = type;
102
- return query;
100
+ return callocz(1, sizeof(struct aclk_query));
101
}
102
103
void aclk_query_free(aclk_query_t query)
104
{
107
- switch (query->type) {
108
- case HTTP_API_V2:
109
- freez(query->data.http_api_v2.payload);
110
- if (query->data.http_api_v2.query != query->dedup_id)
111
- freez(query->data.http_api_v2.query);
112
- break;
113
-
114
- case NODE_STATE_UPDATE:
115
- case REGISTER_NODE:
116
- case CHART_DIMS_UPDATE:
117
- case CHART_CONFIG_UPDATED:
118
- case CHART_RESET:
119
- case RETENTION_UPDATED:
120
- case UPDATE_NODE_INFO:
121
- case ALARM_LOG_HEALTH:
122
- case ALARM_PROVIDE_CFG:
123
- case ALARM_SNAPSHOT:
124
- case UPDATE_NODE_COLLECTORS:
125
- case PROTO_BIN_MESSAGE:
126
- if (!use_mqtt_5)
127
- freez(query->data.bin_payload.payload);
128
- break;
129
-
130
- default:
131
- break;
132
- }
105
+ freez(query->http_api_v2.payload);
106
+ if (query->http_api_v2.query != query->dedup_id)
107
+ freez(query->http_api_v2.query);
108
109
freez(query->dedup_id);
110
freez(query->callback_topic);
aclk/aclk_query_queue.h
+3
-33
@@ -9,24 +9,6 @@
9
10
#include "aclk_util.h"
11
12
-typedef enum {
13
- UNKNOWN = 0,
14
- HTTP_API_V2,
15
- REGISTER_NODE,
16
- NODE_STATE_UPDATE,
17
- CHART_DIMS_UPDATE,
18
- CHART_CONFIG_UPDATED,
19
- CHART_RESET,
20
- RETENTION_UPDATED,
21
- UPDATE_NODE_INFO,
22
- ALARM_LOG_HEALTH,
23
- ALARM_PROVIDE_CFG,
24
- ALARM_SNAPSHOT,
25
- UPDATE_NODE_COLLECTORS,
26
- PROTO_BIN_MESSAGE,
27
- ACLK_QUERY_TYPE_COUNT // always keep this as last
28
-} aclk_query_type_t;
29
-
12
struct aclk_query_http_api_v2 {
13
char *payload;
14
char *query;
@@ -41,8 +23,6 @@ struct aclk_bin_payload {
23
24
typedef struct aclk_query *aclk_query_t;
25
struct aclk_query {
44
- aclk_query_type_t type;
45
-
26
// dedup_id is used to deduplicate queries in the list
27
// if type and dedup_id is the same message is deduplicated
28
// set dedup_id to NULL to never deduplicate the message
@@ -59,13 +39,11 @@ struct aclk_query {
39
40
// TODO maybe remove?
41
int version;
62
- union {
63
- struct aclk_query_http_api_v2 http_api_v2;
64
- struct aclk_bin_payload bin_payload;
65
- } data;
42
+
43
+ struct aclk_query_http_api_v2 http_api_v2;
44
};
45
68
-aclk_query_t aclk_query_new(aclk_query_type_t type);
46
+aclk_query_t aclk_query_new();
47
void aclk_query_free(aclk_query_t query);
48
49
int aclk_queue_query(aclk_query_t query);
@@ -75,12 +53,4 @@ void aclk_queue_flush(void);
53
void aclk_queue_lock(void);
54
void aclk_queue_unlock(void);
55
78
-#define QUEUE_IF_PAYLOAD_PRESENT(query) \
79
- if (likely(query->data.bin_payload.payload)) { \
80
- aclk_queue_query(query); \
81
- } else { \
82
- error("Failed to generate payload (%s)", __FUNCTION__); \
83
- aclk_query_free(query); \
84
- }
85
-
56
#endif /* NETDATA_ACLK_QUERY_QUEUE_H */
aclk/aclk_rx_msgs.c
+8
-9
@@ -5,6 +5,7 @@
5
#include "aclk_stats.h"
6
#include "aclk_query_queue.h"
7
#include "aclk.h"
8
+#include "aclk_tx_msgs.h"
9
10
#include "schema-wrappers/proto_2_json.h"
11
@@ -131,14 +132,14 @@ static int aclk_handle_cloud_http_request_v2(struct aclk_request *cloud_to_agent
132
return 1;
133
}
134
134
- query = aclk_query_new(HTTP_API_V2);
135
+ query = aclk_query_new();
136
136
- if (unlikely(aclk_extract_v2_data(raw_payload, &query->data.http_api_v2.payload))) {
137
+ if (unlikely(aclk_extract_v2_data(raw_payload, &query->http_api_v2.payload))) {
138
error("Error extracting payload expected after the JSON dictionary.");
139
goto error;
140
}
141
141
- if (unlikely(aclk_v2_payload_get_query(query->data.http_api_v2.payload, &query->dedup_id))) {
142
+ if (unlikely(aclk_v2_payload_get_query(query->http_api_v2.payload, &query->dedup_id))) {
143
error("Could not extract payload from query");
144
goto error;
145
}
@@ -158,7 +159,7 @@ static int aclk_handle_cloud_http_request_v2(struct aclk_request *cloud_to_agent
159
query->timeout = cloud_to_agent->timeout;
160
// for clarity and code readability as when we process the request
161
// it would be strange to get URL from `dedup_id`
161
- query->data.http_api_v2.query = query->dedup_id;
162
+ query->http_api_v2.query = query->dedup_id;
163
query->msg_id = cloud_to_agent->msg_id;
164
aclk_queue_query(query);
165
return 0;
@@ -265,7 +266,6 @@ int create_node_instance_result(const char *msg, size_t msg_len)
266
}
267
update_node_id(&host_id, &node_id);
268
268
- aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
269
node_instance_connection_t node_state_update = {
270
.hops = 1,
271
.live = 0,
@@ -298,15 +298,14 @@ int create_node_instance_result(const char *msg, size_t msg_len)
298
};
299
node_state_update.capabilities = caps;
300
301
+ size_t payload_len;
302
rrdhost_aclk_state_lock(localhost);
303
node_state_update.claim_id = localhost->aclk_state.claimed_id;
303
- query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
304
+ char *payload = generate_node_instance_connection(&payload_len, &node_state_update);
305
rrdhost_aclk_state_unlock(localhost);
306
306
- query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
307
- query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
307
+ aclk_send_bin_msg(payload, payload_len, ACLK_TOPICID_NODE_CONN, "UpdateNodeInstanceConnection");
308
309
- aclk_queue_query(query);
309
freez(res.node_id);
310
freez(res.machine_guid);
311
return 0;
aclk/aclk_stats.c
-23
@@ -118,28 +118,6 @@ static void aclk_stats_cloud_req(struct aclk_metrics_per_sample *per_sample)
118
rrdset_done(st);
119
}
120
121
-static void aclk_stats_cloud_req_type(struct aclk_metrics_per_sample *per_sample)
122
-{
123
- static RRDSET *st = NULL;
124
- static RRDDIM *dims[ACLK_QUERY_TYPE_COUNT];
125
-
126
- if (unlikely(!st)) {
127
- st = rrdset_create_localhost(
128
- "netdata", "aclk_processed_query_type", NULL, "aclk", NULL, "Query thread commands processed by their type", "cmd/s",
129
- "netdata", "stats", 200006, localhost->rrd_update_every, RRDSET_TYPE_STACKED);
130
-
131
- for (int i = 0; i < ACLK_QUERY_TYPE_COUNT; i++)
132
- dims[i] = rrddim_add(st, aclk_query_get_name(i), NULL, 1, localhost->rrd_update_every, RRD_ALGORITHM_ABSOLUTE);
133
-
134
- } else
135
- rrdset_next(st);
136
-
137
- for (int i = 0; i < ACLK_QUERY_TYPE_COUNT; i++)
138
- rrddim_set_by_pointer(st, dims[i], per_sample->queries_per_type[i]);
139
-
140
- rrdset_done(st);
141
-}
142
-
121
static char *cloud_req_http_type_names[ACLK_STATS_CLOUD_HTTP_REQ_TYPE_CNT] = {
122
"other",
123
"info",
@@ -351,7 +329,6 @@ void *aclk_stats_main_thread(void *ptr)
329
#endif
330
331
aclk_stats_cloud_req(&per_sample);
354
- aclk_stats_cloud_req_type(&per_sample);
332
aclk_stats_cloud_req_http_type(&per_sample);
333
334
aclk_stats_query_threads(aclk_queries_per_thread_sample);
aclk/aclk_stats.h
-3
@@ -51,9 +51,6 @@ extern struct aclk_metrics_per_sample {
51
volatile uint32_t cloud_req_recvd;
52
volatile uint32_t cloud_req_err;
53
54
- // query types.
55
- volatile uint32_t queries_per_type[ACLK_QUERY_TYPE_COUNT];
56
-
54
// HTTP-specific request types.
55
volatile uint32_t cloud_req_http_by_type[ACLK_STATS_CLOUD_HTTP_REQ_TYPE_CNT];
56
aclk/schema-wrappers/context.cc
+2
-2
@@ -57,7 +57,7 @@ void contexts_snapshot_add_ctx_update(contexts_snapshot_t ctxs_snapshot, struct
57
fill_ctx_updated(ctx, ctx_update);
58
}
59
60
-char *contexts_snapshot_2bin(contexts_snapshot_t ctxs_snapshot, size_t *len)
60
+char *contexts_snapshot_2bin(size_t *len, contexts_snapshot_t ctxs_snapshot)
61
{
62
ContextsSnapshot *ctxs_snap = (ContextsSnapshot *)ctxs_snapshot;
63
*len = PROTO_COMPAT_MSG_SIZE_PTR(ctxs_snap);
@@ -109,7 +109,7 @@ void contexts_updated_add_ctx_update(contexts_updated_t ctxs_updated, struct con
109
fill_ctx_updated(ctx, ctx_update);
110
}
111
112
-char *contexts_updated_2bin(contexts_updated_t ctxs_updated, size_t *len)
112
+char *contexts_updated_2bin(size_t *len, contexts_updated_t ctxs_updated)
113
{
114
ContextsUpdated *ctxs_update = (ContextsUpdated *)ctxs_updated;
115
*len = PROTO_COMPAT_MSG_SIZE_PTR(ctxs_update);
aclk/schema-wrappers/context.h
+2
-2
@@ -36,14 +36,14 @@ contexts_snapshot_t contexts_snapshot_new(const char *claim_id, const char *node
36
void contexts_snapshot_delete(contexts_snapshot_t ctxs_snapshot);
37
void contexts_snapshot_set_version(contexts_snapshot_t ctxs_snapshot, uint64_t version);
38
void contexts_snapshot_add_ctx_update(contexts_snapshot_t ctxs_snapshot, struct context_updated *ctx_update);
39
-char *contexts_snapshot_2bin(contexts_snapshot_t ctxs_snapshot, size_t *len);
39
+char *contexts_snapshot_2bin(size_t *len, contexts_snapshot_t ctxs_snapshot);
40
41
// ContextS Updated related
42
contexts_updated_t contexts_updated_new(const char *claim_id, const char *node_id, uint64_t version_hash, uint64_t created_at);
43
void contexts_updated_delete(contexts_updated_t ctxs_updated);
44
void contexts_updated_update_version_hash(contexts_updated_t ctxs_updated, uint64_t version_hash);
45
void contexts_updated_add_ctx_update(contexts_updated_t ctxs_updated, struct context_updated *ctx_update);
46
-char *contexts_updated_2bin(contexts_updated_t ctxs_updated, size_t *len);
46
+char *contexts_updated_2bin(size_t *len, contexts_updated_t ctxs_updated);
47
48
49
#ifdef __cplusplus