@cryptotaxi247 / netdata-1 / commits / ed52c959d

Revert "Query queue only for queries" (#13452)

Revert "Query queue only for queries (#13431)" This reverts commit 221fd512873613a10b3d95b25a8a4d542b2c4801.

Stelios Fragkakis committed Jul 28, 2022 at 23:59 UTC ed52c959de2e0232170087a3e743562522c1ebfa
14 files changed +253 -77
aclk/aclk.c
+23 -15
@@ -742,8 +742,8 @@ void aclk_host_state_update(RRDHOST *host, int cmd)
742 }
743 if (ret < 0) {
744 // node_id not found
745 - size_t payload_len;
746 -
745 + aclk_query_t create_query;
746 + create_query = aclk_query_new(REGISTER_NODE);
747 rrdhost_aclk_state_lock(localhost);
748 node_instance_creation_t node_instance_creation = {
749 .claim_id = localhost->aclk_state.claimed_id,
@@ -751,14 +751,16 @@ void aclk_host_state_update(RRDHOST *host, int cmd)
751 .hostname = host->hostname,
752 .machine_guid = host->machine_guid
753 };
754 - char *payload = generate_node_instance_creation(&payload_len, &node_instance_creation);
754 + create_query->data.bin_payload.payload = generate_node_instance_creation(&create_query->data.bin_payload.size, &node_instance_creation);
755 rrdhost_aclk_state_unlock(localhost);
756 -
756 + create_query->data.bin_payload.topic = ACLK_TOPICID_CREATE_NODE;
757 + create_query->data.bin_payload.msg_name = "CreateNodeInstance";
758 info("Registering host=%s, hops=%u",host->machine_guid, host->system_info->hops);
758 - aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_CREATE_NODE, "CreateNodeInstance");
759 + aclk_queue_query(create_query);
760 return;
761 }
762
763 + aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
764 node_instance_connection_t node_state_update = {
765 .hops = host->system_info->hops,
766 .live = cmd,
@@ -779,14 +781,15 @@ void aclk_host_state_update(RRDHOST *host, int cmd)
781
782 rrdhost_aclk_state_lock(localhost);
783 node_state_update.claim_id = localhost->aclk_state.claimed_id;
782 - size_t payload_len;
783 - char *payload = generate_node_instance_connection(&payload_len, &node_state_update);
784 + query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
785 rrdhost_aclk_state_unlock(localhost);
786
787 info("Queuing status update for node=%s, live=%d, hops=%u",(char*)node_state_update.node_id, cmd,
788 host->system_info->hops);
789 freez((void*)node_state_update.node_id);
789 - aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_NODE_CONN, "UpdateNodeInstanceConnection");
790 + query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
791 + query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
792 + aclk_queue_query(query);
793 }
794
795 void aclk_send_node_instances()
@@ -799,6 +802,7 @@ void aclk_send_node_instances()
802 }
803 while (!uuid_is_null(list->host_id)) {
804 if (!uuid_is_null(list->node_id)) {
805 + aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
806 node_instance_connection_t node_state_update = {
807 .live = list->live,
808 .hops = list->hops,
@@ -823,30 +827,34 @@ void aclk_send_node_instances()
827
828 rrdhost_aclk_state_lock(localhost);
829 node_state_update.claim_id = localhost->aclk_state.claimed_id;
826 - size_t payload_len;
827 - char *payload = generate_node_instance_connection(&payload_len, &node_state_update);
830 + query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
831 rrdhost_aclk_state_unlock(localhost);
832 info("Queuing status update for node=%s, live=%d, hops=%d",(char*)node_state_update.node_id,
833 list->live,
834 list->hops);
835 freez((void*)node_state_update.node_id);
833 - aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_NODE_CONN, "UpdateNodeInstanceConnection");
836 + query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
837 + query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
838 + aclk_queue_query(query);
839 } else {
840 + aclk_query_t create_query;
841 + create_query = aclk_query_new(REGISTER_NODE);
842 node_instance_creation_t node_instance_creation = {
843 .hops = list->hops,
844 .hostname = list->hostname,
845 };
846 node_instance_creation.machine_guid = mallocz(UUID_STR_LEN);
847 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";
850 rrdhost_aclk_state_lock(localhost);
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);
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);
853 rrdhost_aclk_state_unlock(localhost);
854 info("Queuing registration for host=%s, hops=%d",(char*)node_instance_creation.machine_guid,
855 list->hops);
856 freez(node_instance_creation.machine_guid);
849 - aclk_send_bin_message_subtopic_pid(mqttwss_client, payload, payload_len, ACLK_TOPICID_CREATE_NODE, "CreateNodeInstance");
857 + aclk_queue_query(create_query);
858 }
859 freez(list->hostname);
860
aclk/aclk.h
-11
@@ -34,17 +34,6 @@ 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 -
37 char *ng_aclk_state(void);
38 char *ng_aclk_state_json(void);
39
aclk/aclk_alarm_api.c
+22 -4
@@ -10,20 +10,38 @@
10
11 void aclk_send_alarm_log_health(struct alarm_log_health *log_health)
12 {
13 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_HEALTH, "AlarmLogHealth", generate_alarm_log_health, log_health);
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);
18 }
19
20 void aclk_send_alarm_log_entry(struct alarm_log_entry *log_entry)
21 {
18 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_LOG, "AlarmLogEntry", generate_alarm_log_entry, log_entry);
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);
29 }
30
31 void aclk_send_provide_alarm_cfg(struct provide_alarm_configuration *cfg)
32 {
23 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_CONFIG, "ProvideAlarmConfiguration", generate_provide_alarm_configuration, cfg);
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);
38 }
39
40 void aclk_send_alarm_snapshot(alarm_snapshot_proto_ptr_t snapshot)
41 {
28 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_ALARM_SNAPSHOT, "AlarmSnapshot", generate_alarm_snapshot_bin, snapshot);
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);
47 }
aclk/aclk_charts_api.c
+39 -10
@@ -3,46 +3,75 @@
3
4 #include "aclk_query_queue.h"
5
6 -#include "aclk.h"
7 -
6 #define CHART_DIM_UPDATE_NAME "ChartsAndDimensionsUpdated"
7
8 void aclk_chart_inst_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
9 {
12 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_DIMS, CHART_DIM_UPDATE_NAME, generate_charts_updated, payloads, payload_sizes, new_positions);
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);
14 }
15
16 void aclk_chart_dim_update(char **payloads, size_t *payload_sizes, struct aclk_message_position *new_positions)
17 {
17 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_DIMS, CHART_DIM_UPDATE_NAME, generate_chart_dimensions_updated, payloads, payload_sizes, new_positions);
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);
23 }
24
25 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)
26 {
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);
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);
32 }
33
34 void aclk_chart_config_updated(struct chart_config_updated *config_list, int list_size)
35 {
27 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_CONFIGS_UPDATED, "ChartConfigsUpdated", generate_chart_configs_updated, config_list, list_size);
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);
41 }
42
43 void aclk_chart_reset(chart_reset_t reset)
44 {
32 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CHART_RESET, "ResetChartMessages", generate_reset_chart_messages, reset);
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);
50 }
51
52 void aclk_retention_updated(struct retention_updated *data)
53 {
37 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_RETENTION_UPDATED, "RetentionUpdated", generate_retention_updated, data);
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);
59 }
60
61 void aclk_update_node_info(struct update_node_info *info)
62 {
42 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_NODE_INFO, "UpdateNodeInfo", generate_update_node_info_message, info);
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);
68 }
69
70 void aclk_update_node_collectors(struct update_node_collectors *collectors)
71 {
47 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_NODE_COLLECTORS, "UpdateNodeCollectors", generate_update_node_collectors_message, collectors);
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);
77 }
aclk/aclk_contexts_api.c
+10 -4
@@ -4,14 +4,20 @@
4
5 #include "aclk_contexts_api.h"
6
7 -#include "aclk.h"
8 -
7 void aclk_send_contexts_snapshot(contexts_snapshot_t data)
8 {
11 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CTXS_SNAPSHOT, "ContextsSnapshot", contexts_snapshot_2bin, data);
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);
14 }
15
16 void aclk_send_contexts_updated(contexts_updated_t data)
17 {
16 - GENERATE_AND_SEND_PAYLOAD(ACLK_TOPICID_CTXS_UPDATED, "ContextsUpdated", contexts_updated_2bin, data);
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);
23 }
aclk/aclk_query.c
+55 -13
@@ -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->http_api_v2.query);
87 + buffer_strcat(log_buffer, query->data.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->http_api_v2.query, NODE_ID_QUERY, strlen(NODE_ID_QUERY))) {
106 - char *node_uuid = query->http_api_v2.query + strlen(NODE_ID_QUERY);
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);
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->http_api_v2.query, '?');
130 + char *mysep = strchr(query->data.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->http_api_v2.query, NETDATA_WEB_REQUEST_URL_SIZE + 1);
135 + url_decode_r(w->decoded_query_string, query->data.http_api_v2.query, NETDATA_WEB_REQUEST_URL_SIZE + 1);
136
137 - mysep = strrchr(query->http_api_v2.query, '/');
137 + mysep = strrchr(query->data.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->http_api_v2.payload, WEB_HDR_ACCEPT_ENC))) {
154 + if ((start = strstr((char *)query->data.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,16 +256,57 @@ 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 +
288 static void aclk_query_process_msg(struct aclk_query_thread *query_thr, aclk_query_t query)
289 {
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);
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 + }
304
305 if (aclk_stats_enabled) {
306 ACLK_STATS_LOCK;
307 aclk_metrics_per_sample.queries_dispatched++;
308 aclk_queries_per_thread[query_thr->idx]++;
309 + aclk_metrics_per_sample.queries_per_type[query->type]++;
310 ACLK_STATS_UNLOCK;
311 }
312
@@ -285,10 +326,11 @@ int aclk_query_process_msgs(struct aclk_query_thread *query_thr)
326 return 0;
327 }
328
288 -static void worker_aclk_register(void)
289 -{
329 +static void worker_aclk_register(void) {
330 worker_register("ACLKQUERY");
291 - worker_register_job_name(0, "http query");
331 + for (int i = 1; i < ACLK_QUERY_TYPE_COUNT; i++) {
332 + worker_register_job_name(i, aclk_query_get_name(i));
333 + }
334 }
335
336 /**
aclk/aclk_query.h
+2
@@ -31,4 +31,6 @@ 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 +
36 #endif //NETDATA_AGENT_CLOUD_LINK_H
aclk/aclk_query_queue.c
+30 -5
@@ -95,16 +95,41 @@ void aclk_queue_flush(void)
95 };
96 }
97
98 -aclk_query_t aclk_query_new()
98 +aclk_query_t aclk_query_new(aclk_query_type_t type)
99 {
100 - return callocz(1, sizeof(struct aclk_query));
100 + aclk_query_t query = callocz(1, sizeof(struct aclk_query));
101 + query->type = type;
102 + return query;
103 }
104
105 void aclk_query_free(aclk_query_t query)
106 {
105 - freez(query->http_api_v2.payload);
106 - if (query->http_api_v2.query != query->dedup_id)
107 - freez(query->http_api_v2.query);
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 + }
133
134 freez(query->dedup_id);
135 freez(query->callback_topic);
aclk/aclk_query_queue.h
+33 -3
@@ -9,6 +9,24 @@
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 +
30 struct aclk_query_http_api_v2 {
31 char *payload;
32 char *query;
@@ -23,6 +41,8 @@ struct aclk_bin_payload {
41
42 typedef struct aclk_query *aclk_query_t;
43 struct aclk_query {
44 + aclk_query_type_t type;
45 +
46 // dedup_id is used to deduplicate queries in the list
47 // if type and dedup_id is the same message is deduplicated
48 // set dedup_id to NULL to never deduplicate the message
@@ -39,11 +59,13 @@ struct aclk_query {
59
60 // TODO maybe remove?
61 int version;
42 -
43 - struct aclk_query_http_api_v2 http_api_v2;
62 + union {
63 + struct aclk_query_http_api_v2 http_api_v2;
64 + struct aclk_bin_payload bin_payload;
65 + } data;
66 };
67
46 -aclk_query_t aclk_query_new();
68 +aclk_query_t aclk_query_new(aclk_query_type_t type);
69 void aclk_query_free(aclk_query_t query);
70
71 int aclk_queue_query(aclk_query_t query);
@@ -53,4 +75,12 @@ void aclk_queue_flush(void);
75 void aclk_queue_lock(void);
76 void aclk_queue_unlock(void);
77
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 +
86 #endif /* NETDATA_ACLK_QUERY_QUEUE_H */
aclk/aclk_rx_msgs.c
+9 -8
@@ -5,7 +5,6 @@
5 #include "aclk_stats.h"
6 #include "aclk_query_queue.h"
7 #include "aclk.h"
8 -#include "aclk_tx_msgs.h"
8
9 #include "schema-wrappers/proto_2_json.h"
10
@@ -132,14 +131,14 @@ static int aclk_handle_cloud_http_request_v2(struct aclk_request *cloud_to_agent
131 return 1;
132 }
133
135 - query = aclk_query_new();
134 + query = aclk_query_new(HTTP_API_V2);
135
137 - if (unlikely(aclk_extract_v2_data(raw_payload, &query->http_api_v2.payload))) {
136 + if (unlikely(aclk_extract_v2_data(raw_payload, &query->data.http_api_v2.payload))) {
137 error("Error extracting payload expected after the JSON dictionary.");
138 goto error;
139 }
140
142 - if (unlikely(aclk_v2_payload_get_query(query->http_api_v2.payload, &query->dedup_id))) {
141 + if (unlikely(aclk_v2_payload_get_query(query->data.http_api_v2.payload, &query->dedup_id))) {
142 error("Could not extract payload from query");
143 goto error;
144 }
@@ -159,7 +158,7 @@ static int aclk_handle_cloud_http_request_v2(struct aclk_request *cloud_to_agent
158 query->timeout = cloud_to_agent->timeout;
159 // for clarity and code readability as when we process the request
160 // it would be strange to get URL from `dedup_id`
162 - query->http_api_v2.query = query->dedup_id;
161 + query->data.http_api_v2.query = query->dedup_id;
162 query->msg_id = cloud_to_agent->msg_id;
163 aclk_queue_query(query);
164 return 0;
@@ -266,6 +265,7 @@ int create_node_instance_result(const char *msg, size_t msg_len)
265 }
266 update_node_id(&host_id, &node_id);
267
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,14 +298,15 @@ 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;
301 rrdhost_aclk_state_lock(localhost);
302 node_state_update.claim_id = localhost->aclk_state.claimed_id;
304 - char *payload = generate_node_instance_connection(&payload_len, &node_state_update);
303 + query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
304 rrdhost_aclk_state_unlock(localhost);
305
307 - aclk_send_bin_msg(payload, payload_len, ACLK_TOPICID_NODE_CONN, "UpdateNodeInstanceConnection");
306 + query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
307 + query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
308
309 + aclk_queue_query(query);
310 freez(res.node_id);
311 freez(res.machine_guid);
312 return 0;
aclk/aclk_stats.c
+23
@@ -118,6 +118,28 @@ 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 +
143 static char *cloud_req_http_type_names[ACLK_STATS_CLOUD_HTTP_REQ_TYPE_CNT] = {
144 "other",
145 "info",
@@ -329,6 +351,7 @@ void *aclk_stats_main_thread(void *ptr)
351 #endif
352
353 aclk_stats_cloud_req(&per_sample);
354 + aclk_stats_cloud_req_type(&per_sample);
355 aclk_stats_cloud_req_http_type(&per_sample);
356
357 aclk_stats_query_threads(aclk_queries_per_thread_sample);
aclk/aclk_stats.h
+3
@@ -51,6 +51,9 @@ 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 +
57 // HTTP-specific request types.
58 volatile uint32_t cloud_req_http_by_type[ACLK_STATS_CLOUD_HTTP_REQ_TYPE_CNT];
59
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(size_t *len, contexts_snapshot_t ctxs_snapshot)
60 +char *contexts_snapshot_2bin(contexts_snapshot_t ctxs_snapshot, size_t *len)
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(size_t *len, contexts_updated_t ctxs_updated)
112 +char *contexts_updated_2bin(contexts_updated_t ctxs_updated, size_t *len)
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(size_t *len, contexts_snapshot_t ctxs_snapshot);
39 +char *contexts_snapshot_2bin(contexts_snapshot_t ctxs_snapshot, size_t *len);
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(size_t *len, contexts_updated_t ctxs_updated);
46 +char *contexts_updated_2bin(contexts_updated_t ctxs_updated, size_t *len);
47
48
49 #ifdef __cplusplus