UpdateNodeCollectors message (#13330)
* add new aclk-schemas. remove services related * add updatenodecollectors message * build with --disable-cloud
Emmanuel Vasilakis committed
Jul 7, 2022 at 21:50 UTC
19d9a0030db7d8ff6e43f2ff76cea7f9761c6cd7
14 files changed
+115
-10
aclk/aclk-schemas
+1
-1
@@ -1 +1 @@
1
-Subproject commit d8342ee6d932c152a78c9fe886281fe28170a6c4
1
+Subproject commit fa46ccca237a9bdb613b3b1f2809a25b7b45c7c4
aclk/aclk_charts_api.c
+9
@@ -66,3 +66,12 @@ void aclk_update_node_info(struct update_node_info *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
+{
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_charts_api.h
+2
@@ -17,4 +17,6 @@ void aclk_retention_updated(struct retention_updated *data);
17
18
void aclk_update_node_info(struct update_node_info *info);
19
20
+void aclk_update_node_collectors(struct update_node_collectors *collectors);
21
+
22
#endif /* ACLK_CHARTS_H */
aclk/aclk_query.c
+1
@@ -277,6 +277,7 @@ const char *aclk_query_get_name(aclk_query_type_t qt)
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
default:
282
error_report("Unknown query type used %d", (int) qt);
283
return "unknown";
aclk/aclk_query_queue.c
+1
@@ -121,6 +121,7 @@ void aclk_query_free(aclk_query_t query)
121
case ALARM_LOG_HEALTH:
122
case ALARM_PROVIDE_CFG:
123
case ALARM_SNAPSHOT:
124
+ case UPDATE_NODE_COLLECTORS:
125
if (!use_mqtt_5)
126
freez(query->data.bin_payload.payload);
127
break;
aclk/aclk_query_queue.h
+1
@@ -22,6 +22,7 @@ typedef enum {
22
ALARM_LOG_HEALTH,
23
ALARM_PROVIDE_CFG,
24
ALARM_SNAPSHOT,
25
+ UPDATE_NODE_COLLECTORS,
26
ACLK_QUERY_TYPE_COUNT // always keep this as last
27
} aclk_query_type_t;
28
aclk/aclk_util.c
+2
@@ -123,6 +123,7 @@ struct topic_name {
123
{ .id = ACLK_TOPICID_ALARM_HEALTH, .name = "alarm-health" },
124
{ .id = ACLK_TOPICID_ALARM_CONFIG, .name = "alarm-config" },
125
{ .id = ACLK_TOPICID_ALARM_SNAPSHOT, .name = "alarm-snapshot" },
126
+ { .id = ACLK_TOPICID_NODE_COLLECTORS, .name = "node-instance-collectors" },
127
{ .id = ACLK_TOPICID_UNKNOWN, .name = NULL }
128
};
129
@@ -145,6 +146,7 @@ enum aclk_topics compulsory_topics[] = {
146
ACLK_TOPICID_ALARM_HEALTH,
147
ACLK_TOPICID_ALARM_CONFIG,
148
ACLK_TOPICID_ALARM_SNAPSHOT,
149
+ ACLK_TOPICID_NODE_COLLECTORS,
150
ACLK_TOPICID_UNKNOWN
151
};
152
aclk/aclk_util.h
+2
-1
@@ -87,7 +87,8 @@ enum aclk_topics {
87
ACLK_TOPICID_ALARM_LOG = 14,
88
ACLK_TOPICID_ALARM_HEALTH = 15,
89
ACLK_TOPICID_ALARM_CONFIG = 16,
90
- ACLK_TOPICID_ALARM_SNAPSHOT = 17
90
+ ACLK_TOPICID_ALARM_SNAPSHOT = 17,
91
+ ACLK_TOPICID_NODE_COLLECTORS = 18
92
};
93
94
const char *aclk_get_topic(enum aclk_topics topic);
aclk/schema-wrappers/node_info.cc
+24
-3
@@ -55,9 +55,6 @@ static int generate_node_info(nodeinstance::info::v1::NodeInfo *info, struct acl
55
if (data->custom_info)
56
info->set_custom_info(data->custom_info);
57
58
- for (size_t i = 0; i < data->service_count; i++)
59
- info->add_services(data->services[i]);
60
-
58
if (data->machine_guid)
59
info->set_machine_guid(data->machine_guid);
60
@@ -113,3 +110,27 @@ char *generate_update_node_info_message(size_t *len, struct update_node_info *in
110
111
return bin;
112
}
113
+
114
+char *generate_update_node_collectors_message(size_t *len, struct update_node_collectors *upd_node_collectors)
115
+{
116
+ nodeinstance::info::v1::UpdateNodeCollectors msg;
117
+
118
+ msg.set_node_id(upd_node_collectors->node_id);
119
+ msg.set_claim_id(upd_node_collectors->claim_id);
120
+
121
+ void *colls;
122
+ dfe_start_read(upd_node_collectors->node_collectors, colls) {
123
+ struct collector_info *c =(struct collector_info *)colls;
124
+ nodeinstance::info::v1::CollectorInfo *col = msg.add_collectors();
125
+ col->set_plugin(c->plugin);
126
+ col->set_module(c->module);
127
+ }
128
+ dfe_done(colls);
129
+
130
+ *len = PROTO_COMPAT_MSG_SIZE(msg);
131
+ char *bin = (char*)malloc(*len);
132
+ if (bin)
133
+ msg.SerializeToArray(bin, *len);
134
+
135
+ return bin;
136
+}
aclk/schema-wrappers/node_info.h
+13
-3
@@ -50,9 +50,6 @@ struct aclk_node_info {
50
51
char *custom_info;
52
53
- char **services;
54
- size_t service_count;
55
-
53
char *machine_guid;
54
55
DICTIONARY *host_labels_ptr;
@@ -74,8 +71,21 @@ struct update_node_info {
71
struct capability *node_instance_capabilities;
72
};
73
74
+struct collector_info {
75
+ char *module;
76
+ char *plugin;
77
+};
78
+
79
+struct update_node_collectors {
80
+ char *claim_id;
81
+ char *node_id;
82
+ DICTIONARY *node_collectors;
83
+};
84
+
85
char *generate_update_node_info_message(size_t *len, struct update_node_info *info);
86
87
+char *generate_update_node_collectors_message(size_t *len, struct update_node_collectors *collectors);
88
+
89
#ifdef __cplusplus
90
}
91
#endif
database/sqlite/sqlite_aclk.c
+10
@@ -423,6 +423,7 @@ void aclk_database_worker(void *arg)
423
worker_register_job_name(ACLK_DATABASE_CLEANUP, "cleanup");
424
worker_register_job_name(ACLK_DATABASE_DELETE_HOST, "node delete");
425
worker_register_job_name(ACLK_DATABASE_NODE_INFO, "node info");
426
+ worker_register_job_name(ACLK_DATABASE_NODE_COLLECTORS, "node collectors");
427
worker_register_job_name(ACLK_DATABASE_PUSH_ALERT, "alert push");
428
worker_register_job_name(ACLK_DATABASE_PUSH_ALERT_CONFIG, "alert conf push");
429
worker_register_job_name(ACLK_DATABASE_PUSH_ALERT_SNAPSHOT, "alert snapshot");
@@ -583,6 +584,10 @@ void aclk_database_worker(void *arg)
584
debug(D_ACLK_SYNC,"Sending node info for %s", wc->uuid_str);
585
sql_build_node_info(wc, cmd);
586
break;
587
+ case ACLK_DATABASE_NODE_COLLECTORS:
588
+ debug(D_ACLK_SYNC,"Sending node collectors info for %s", wc->uuid_str);
589
+ sql_build_node_collectors(wc);
590
+ break;
591
#ifdef ENABLE_ACLK
592
case ACLK_DATABASE_DIM_DELETION:
593
debug(D_ACLK_SYNC,"Sending dimension deletion information %s", wc->uuid_str);
@@ -634,6 +639,11 @@ void aclk_database_worker(void *arg)
639
cmd.completion = NULL;
640
wc->node_info_send = aclk_database_enq_cmd_noblock(wc, &cmd);
641
}
642
+ if (wc->node_collectors_send && wc->node_collectors_send + 30 < now_realtime_sec()) {
643
+ cmd.opcode = ACLK_DATABASE_NODE_COLLECTORS;
644
+ cmd.completion = NULL;
645
+ wc->node_collectors_send = aclk_database_enq_cmd_noblock(wc, &cmd);
646
+ }
647
if (localhost == wc->host)
648
(void) sqlite3_wal_checkpoint(db_meta, NULL);
649
break;
database/sqlite/sqlite_aclk.h
+2
@@ -132,6 +132,7 @@ enum aclk_database_opcode {
132
ACLK_DATABASE_PUSH_ALERT_CONFIG,
133
ACLK_DATABASE_PUSH_ALERT_SNAPSHOT,
134
ACLK_DATABASE_QUEUE_REMOVED_ALERTS,
135
+ ACLK_DATABASE_NODE_COLLECTORS,
136
ACLK_DATABASE_TIMER,
137
138
// leave this last
@@ -195,6 +196,7 @@ struct aclk_database_worker_config {
196
int alert_updates;
197
time_t batch_created;
198
int node_info_send;
199
+ time_t node_collectors_send;
200
int chart_pending;
201
int chart_reset_count;
202
int retention_running;
database/sqlite/sqlite_aclk_node.c
+46
-2
@@ -7,6 +7,50 @@
7
#include "../../aclk/aclk_charts_api.h"
8
#endif
9
10
+#ifdef ENABLE_ACLK
11
+DICTIONARY *collectors_from_charts(RRDHOST *host, DICTIONARY *dict) {
12
+ RRDSET *st;
13
+ char name[500];
14
+
15
+ rrdhost_rdlock(host);
16
+ rrdset_foreach_read(st, host) {
17
+ if (rrdset_is_available_for_viewers(st)) {
18
+ struct collector_info col = {
19
+ .plugin = st->plugin_name ? st->plugin_name : "",
20
+ .module = st->module_name ? st->module_name : ""
21
+ };
22
+ snprintfz(name, 499, "%s:%s", col.plugin, col.module);
23
+ dictionary_set(dict, name, &col, sizeof(struct collector_info));
24
+ }
25
+ }
26
+ rrdhost_unlock(host);
27
+
28
+ return dict;
29
+}
30
+#endif
31
+
32
+void sql_build_node_collectors(struct aclk_database_worker_config *wc)
33
+{
34
+#ifdef ENABLE_ACLK
35
+ struct update_node_collectors upd_node_collectors;
36
+ DICTIONARY *dict = dictionary_create(DICTIONARY_FLAG_SINGLE_THREADED);
37
+
38
+ upd_node_collectors.node_id = wc->node_id;
39
+ upd_node_collectors.claim_id = is_agent_claimed();
40
+
41
+ upd_node_collectors.node_collectors = collectors_from_charts(wc->host, dict);
42
+ aclk_update_node_collectors(&upd_node_collectors);
43
+
44
+ dictionary_destroy(dict);
45
+ freez(upd_node_collectors.claim_id);
46
+
47
+ log_access("ACLK RES [%s (%s)]: NODE COLLECTORS SENT", wc->node_id, wc->host->hostname);
48
+#else
49
+ UNUSED(wc);
50
+#endif
51
+ return;
52
+}
53
+
54
void sql_build_node_info(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
55
{
56
UNUSED(cmd);
@@ -61,8 +105,6 @@ void sql_build_node_info(struct aclk_database_worker_config *wc, struct aclk_dat
105
node_info.data.virtualization_type = host->system_info->virtualization ? host->system_info->virtualization : "unknown";
106
node_info.data.container_type = host->system_info->container ? host->system_info->container : "unknown";
107
node_info.data.custom_info = config_get(CONFIG_SECTION_WEB, "custom dashboard_info.js", "");
64
- node_info.data.services = NULL; // char **
65
- node_info.data.service_count = 0;
108
node_info.data.machine_guid = wc->host_guid;
109
110
struct capability node_caps[] = {
@@ -83,6 +125,8 @@ void sql_build_node_info(struct aclk_database_worker_config *wc, struct aclk_dat
125
rrd_unlock();
126
freez(node_info.claim_id);
127
freez(host_version);
128
+
129
+ wc->node_collectors_send = now_realtime_sec();
130
#else
131
UNUSED(wc);
132
#endif
database/sqlite/sqlite_aclk_node.h
+1
@@ -4,4 +4,5 @@
4
#define NETDATA_SQLITE_ACLK_NODE_H
5
6
void sql_build_node_info(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
7
+void sql_build_node_collectors(struct aclk_database_worker_config *wc);
8
#endif //NETDATA_SQLITE_ACLK_NODE_H