@cryptotaxi247 / netdata-1 / commits / 0d2c327ae

Add a checkpoint message to alerts stream (#14847)

* pull aclk schemas * resolve capas * handle checkpoints and removed from health * build with disable-cloud * codacy 1 * misc changes * one more char in hash * free buffer * change topic * misc fixes * skip removed alert variables * change hash functions * use create and destroy for compatibility with older openssl

Emmanuel Vasilakis committed Apr 21, 2023 at 12:24 UTC 0d2c327ae58af5de35c9cb35a20cce9f59e92b9c
21 files changed +344 -365
aclk/aclk-schemas
+1 -1
@@ -1 +1 @@
1 -Subproject commit 3252118bd547640251356629f0df05eaf952ac39
1 +Subproject commit d3a5c636b6dacf364834f2ba99ce0170c71ef861
aclk/aclk.c
-12
@@ -49,8 +49,6 @@ float last_backoff_value = 0;
49
50 time_t aclk_block_until = 0;
51
52 -int aclk_alert_reloaded = 0; //1 on health log exchange, and again on health_reload
53 -
52 #ifdef ENABLE_ACLK
53 mqtt_wss_client mqttwss_client;
54
@@ -928,14 +926,10 @@ static void fill_alert_status_for_host(BUFFER *wb, RRDHOST *host)
926 }
927 buffer_sprintf(wb,
928 "\n\t\tUpdates: %d"
931 - "\n\t\tBatch ID: %"PRIu64
932 - "\n\t\tLast Acked Seq ID: %"PRIu64
929 "\n\t\tPending Min Seq ID: %"PRIu64
930 "\n\t\tPending Max Seq ID: %"PRIu64
931 "\n\t\tLast Submitted Seq ID: %"PRIu64,
932 status.alert_updates,
937 - status.alerts_batch_id,
938 - status.last_acked_sequence_id,
933 status.pending_min_sequence_id,
934 status.pending_max_sequence_id,
935 status.last_submitted_sequence_id
@@ -1043,12 +1037,6 @@ static void fill_alert_status_for_host_json(json_object *obj, RRDHOST *host)
1037 json_object *tmp = json_object_new_int(status.alert_updates);
1038 json_object_object_add(obj, "updates", tmp);
1039
1046 - tmp = json_object_new_int(status.alerts_batch_id);
1047 - json_object_object_add(obj, "batch-id", tmp);
1048 -
1049 - tmp = json_object_new_int(status.last_acked_sequence_id);
1050 - json_object_object_add(obj, "last-acked-seq-id", tmp);
1051 -
1040 tmp = json_object_new_int(status.pending_min_sequence_id);
1041 json_object_object_add(obj, "pending-min-seq-id", tmp);
1042
aclk/aclk.h
-2
@@ -26,8 +26,6 @@ extern time_t aclk_block_until;
26
27 extern int disconnect_req;
28
29 -extern int aclk_alert_reloaded;
30 -
29 #ifdef ENABLE_ACLK
30 void *aclk_main(void *ptr);
31
aclk/aclk_alarm_api.c
+5 -5
@@ -8,12 +8,12 @@
8
9 #include "aclk.h"
10
11 -void aclk_send_alarm_log_health(struct alarm_log_health *log_health)
11 +void aclk_send_provide_alarm_checkpoint(struct alarm_checkpoint *checkpoint)
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";
13 + aclk_query_t query = aclk_query_new(ALARM_PROVIDE_CHECKPOINT);
14 + query->data.bin_payload.payload = generate_alarm_checkpoint(&query->data.bin_payload.size, checkpoint);
15 + query->data.bin_payload.topic = ACLK_TOPICID_ALARM_CHECKPOINT;
16 + query->data.bin_payload.msg_name = "AlarmCheckpoint";
17 QUEUE_IF_PAYLOAD_PRESENT(query);
18 }
19
aclk/aclk_alarm_api.h
+1 -1
@@ -6,7 +6,7 @@
6 #include "../daemon/common.h"
7 #include "schema-wrappers/schema_wrappers.h"
8
9 -void aclk_send_alarm_log_health(struct alarm_log_health *log_health);
9 +void aclk_send_provide_alarm_checkpoint(struct alarm_checkpoint *checkpoint);
10 void aclk_send_alarm_log_entry(struct alarm_log_entry *log_entry);
11 void aclk_send_provide_alarm_cfg(struct provide_alarm_configuration *cfg);
12 void aclk_send_alarm_snapshot(alarm_snapshot_proto_ptr_t snapshot);
aclk/aclk_capas.c
+4
@@ -14,6 +14,7 @@ const struct capability *aclk_get_agent_capas()
14 { .name = "ctx", .version = 1, .enabled = 1 },
15 { .name = "funcs", .version = 1, .enabled = 1 },
16 { .name = "http_api_v2", .version = 1, .enabled = 1 },
17 + { .name = "health", .version = 1, .enabled = 0 },
18 { .name = NULL, .version = 0, .enabled = 0 }
19 };
20 agent_capabilities[2].version = ml_capable() ? 1 : 0;
@@ -22,6 +23,8 @@ const struct capability *aclk_get_agent_capas()
23 agent_capabilities[3].version = enable_metric_correlations ? metric_correlations_version : 0;
24 agent_capabilities[3].enabled = enable_metric_correlations;
25
26 + agent_capabilities[7].enabled = localhost->health.health_enabled;
27 +
28 return agent_capabilities;
29 }
30
@@ -36,6 +39,7 @@ struct capability *aclk_get_node_instance_capas(RRDHOST *host)
39 { .name = "ctx", .version = 1, .enabled = 1 },
40 { .name = "funcs", .version = 0, .enabled = 0 },
41 { .name = "http_api_v2", .version = 1, .enabled = 1 },
42 + { .name = "health", .version = 1, .enabled = host->health.health_enabled },
43 { .name = NULL, .version = 0, .enabled = 0 }
44 };
45
aclk/aclk_query.c
+13 -13
@@ -185,19 +185,19 @@ static int send_bin_msg(struct aclk_query_thread *query_thr, aclk_query_t query)
185 const char *aclk_query_get_name(aclk_query_type_t qt, int unknown_ok)
186 {
187 switch (qt) {
188 - case HTTP_API_V2: return "http_api_request_v2";
189 - case REGISTER_NODE: return "register_node";
190 - case NODE_STATE_UPDATE: return "node_state_update";
191 - case CHART_DIMS_UPDATE: return "chart_and_dim_update";
192 - case CHART_CONFIG_UPDATED: return "chart_config_updated";
193 - case CHART_RESET: return "reset_chart_messages";
194 - case RETENTION_UPDATED: return "update_retention_info";
195 - case UPDATE_NODE_INFO: return "update_node_info";
196 - case ALARM_LOG_HEALTH: return "alarm_log_health";
197 - case ALARM_PROVIDE_CFG: return "provide_alarm_config";
198 - case ALARM_SNAPSHOT: return "alarm_snapshot";
199 - case UPDATE_NODE_COLLECTORS: return "update_node_collectors";
200 - case PROTO_BIN_MESSAGE: return "generic_binary_proto_message";
188 + case HTTP_API_V2: return "http_api_request_v2";
189 + case REGISTER_NODE: return "register_node";
190 + case NODE_STATE_UPDATE: return "node_state_update";
191 + case CHART_DIMS_UPDATE: return "chart_and_dim_update";
192 + case CHART_CONFIG_UPDATED: return "chart_config_updated";
193 + case CHART_RESET: return "reset_chart_messages";
194 + case RETENTION_UPDATED: return "update_retention_info";
195 + case UPDATE_NODE_INFO: return "update_node_info";
196 + case ALARM_PROVIDE_CHECKPOINT: return "alarm_checkpoint";
197 + case ALARM_PROVIDE_CFG: return "provide_alarm_config";
198 + case ALARM_SNAPSHOT: return "alarm_snapshot";
199 + case UPDATE_NODE_COLLECTORS: return "update_node_collectors";
200 + case PROTO_BIN_MESSAGE: return "generic_binary_proto_message";
201 default:
202 if (!unknown_ok)
203 error_report("Unknown query type used %d", (int) qt);
aclk/aclk_query_queue.h
+1 -1
@@ -19,7 +19,7 @@ typedef enum {
19 CHART_RESET,
20 RETENTION_UPDATED,
21 UPDATE_NODE_INFO,
22 - ALARM_LOG_HEALTH,
22 + ALARM_PROVIDE_CHECKPOINT,
23 ALARM_PROVIDE_CFG,
24 ALARM_SNAPSHOT,
25 UPDATE_NODE_COLLECTORS,
aclk/aclk_rx_msgs.c
+14 -12
@@ -339,25 +339,27 @@ int update_chart_configs(const char *msg, size_t msg_len)
339 int start_alarm_streaming(const char *msg, size_t msg_len)
340 {
341 struct start_alarm_streaming res = parse_start_alarm_streaming(msg, msg_len);
342 - if (!res.node_id || !res.batch_id) {
342 + if (!res.node_id) {
343 error("Error parsing StartAlarmStreaming");
344 - freez(res.node_id);
344 return 1;
345 }
347 - aclk_start_alert_streaming(res.node_id, res.batch_id, res.start_seq_id);
346 + aclk_start_alert_streaming(res.node_id, res.resets);
347 freez(res.node_id);
348 return 0;
349 }
350
352 -int send_alarm_log_health(const char *msg, size_t msg_len)
351 +int send_alarm_checkpoint(const char *msg, size_t msg_len)
352 {
354 - char *node_id = parse_send_alarm_log_health(msg, msg_len);
355 - if (!node_id) {
356 - error("Error parsing SendAlarmLogHealth");
353 + struct send_alarm_checkpoint sac = parse_send_alarm_checkpoint(msg, msg_len);
354 + if (!sac.node_id || !sac.claim_id) {
355 + error("Error parsing SendAlarmCheckpoint");
356 + freez(sac.node_id);
357 + freez(sac.claim_id);
358 return 1;
359 }
359 - aclk_send_alarm_health_log(node_id);
360 - freez(node_id);
360 + aclk_send_alarm_checkpoint(sac.node_id, sac.claim_id);
361 + freez(sac.node_id);
362 + freez(sac.claim_id);
363 return 0;
364 }
365
@@ -377,12 +379,12 @@ int send_alarm_configuration(const char *msg, size_t msg_len)
379 int send_alarm_snapshot(const char *msg, size_t msg_len)
380 {
381 struct send_alarm_snapshot *sas = parse_send_alarm_snapshot(msg, msg_len);
380 - if (!sas->node_id || !sas->claim_id) {
382 + if (!sas->node_id || !sas->claim_id || !sas->snapshot_uuid) {
383 error("Error parsing SendAlarmSnapshot");
384 destroy_send_alarm_snapshot(sas);
385 return 1;
386 }
385 - aclk_process_send_alarm_snapshot(sas->node_id, sas->claim_id, sas->snapshot_id, sas->sequence_id);
387 + aclk_process_send_alarm_snapshot(sas->node_id, sas->claim_id, sas->snapshot_uuid);
388 destroy_send_alarm_snapshot(sas);
389 return 0;
390 }
@@ -458,7 +460,7 @@ new_cloud_rx_msg_t rx_msgs[] = {
460 { .name = "ChartsAndDimensionsAck", .name_hash = 0, .fnc = charts_and_dimensions_ack },
461 { .name = "UpdateChartConfigs", .name_hash = 0, .fnc = update_chart_configs },
462 { .name = "StartAlarmStreaming", .name_hash = 0, .fnc = start_alarm_streaming },
461 - { .name = "SendAlarmLogHealth", .name_hash = 0, .fnc = send_alarm_log_health },
463 + { .name = "SendAlarmCheckpoint", .name_hash = 0, .fnc = send_alarm_checkpoint },
464 { .name = "SendAlarmConfiguration", .name_hash = 0, .fnc = send_alarm_configuration },
465 { .name = "SendAlarmSnapshot", .name_hash = 0, .fnc = send_alarm_snapshot },
466 { .name = "DisconnectReq", .name_hash = 0, .fnc = handle_disconnect_req },
aclk/aclk_util.c
+4 -4
@@ -120,10 +120,10 @@ struct topic_name {
120 { .id = ACLK_TOPICID_CHART_RESET, .name = "reset-charts" },
121 { .id = ACLK_TOPICID_RETENTION_UPDATED, .name = "chart-retention-updated" },
122 { .id = ACLK_TOPICID_NODE_INFO, .name = "node-instance-info" },
123 - { .id = ACLK_TOPICID_ALARM_LOG, .name = "alarm-log" },
124 - { .id = ACLK_TOPICID_ALARM_HEALTH, .name = "alarm-health" },
123 + { .id = ACLK_TOPICID_ALARM_LOG, .name = "alarm-log-v2" },
124 + { .id = ACLK_TOPICID_ALARM_CHECKPOINT, .name = "alarm-checkpoint" },
125 { .id = ACLK_TOPICID_ALARM_CONFIG, .name = "alarm-config" },
126 - { .id = ACLK_TOPICID_ALARM_SNAPSHOT, .name = "alarm-snapshot" },
126 + { .id = ACLK_TOPICID_ALARM_SNAPSHOT, .name = "alarm-snapshot-v2" },
127 { .id = ACLK_TOPICID_NODE_COLLECTORS, .name = "node-instance-collectors" },
128 { .id = ACLK_TOPICID_CTXS_SNAPSHOT, .name = "contexts-snapshot" },
129 { .id = ACLK_TOPICID_CTXS_UPDATED, .name = "contexts-updated" },
@@ -146,7 +146,7 @@ enum aclk_topics compulsory_topics[] = {
146 ACLK_TOPICID_RETENTION_UPDATED,
147 ACLK_TOPICID_NODE_INFO,
148 ACLK_TOPICID_ALARM_LOG,
149 - ACLK_TOPICID_ALARM_HEALTH,
149 + ACLK_TOPICID_ALARM_CHECKPOINT,
150 ACLK_TOPICID_ALARM_CONFIG,
151 ACLK_TOPICID_ALARM_SNAPSHOT,
152 ACLK_TOPICID_NODE_COLLECTORS,
aclk/aclk_util.h
+1 -1
@@ -85,7 +85,7 @@ enum aclk_topics {
85 ACLK_TOPICID_RETENTION_UPDATED = 12,
86 ACLK_TOPICID_NODE_INFO = 13,
87 ACLK_TOPICID_ALARM_LOG = 14,
88 - ACLK_TOPICID_ALARM_HEALTH = 15,
88 + ACLK_TOPICID_ALARM_CHECKPOINT = 15,
89 ACLK_TOPICID_ALARM_CONFIG = 16,
90 ACLK_TOPICID_ALARM_SNAPSHOT = 17,
91 ACLK_TOPICID_NODE_COLLECTORS = 18,
aclk/schema-wrappers/alarm_stream.cc
+32 -48
@@ -21,57 +21,24 @@ struct start_alarm_streaming parse_start_alarm_streaming(const char *data, size_
21 return ret;
22
23 ret.node_id = strdupz(msg.node_id().c_str());
24 - ret.batch_id = msg.batch_id();
25 - ret.start_seq_id = msg.start_sequnce_id();
24 + ret.resets = msg.resets();
25
26 return ret;
27 }
28
30 -char *parse_send_alarm_log_health(const char *data, size_t len)
29 +struct send_alarm_checkpoint parse_send_alarm_checkpoint(const char *data, size_t len)
30 {
32 - SendAlarmLogHealth msg;
33 - if (!msg.ParseFromArray(data, len))
34 - return NULL;
35 - return strdupz(msg.node_id().c_str());
36 -}
37 -
38 -char *generate_alarm_log_health(size_t *len, struct alarm_log_health *data)
39 -{
40 - AlarmLogHealth msg;
41 - LogEntries *entries;
42 -
43 - msg.set_claim_id(data->claim_id);
44 - msg.set_node_id(data->node_id);
45 - msg.set_enabled(data->enabled);
46 -
47 - switch (data->status) {
48 - case alarm_log_status_aclk::ALARM_LOG_STATUS_IDLE:
49 - msg.set_status(alarms::v1::ALARM_LOG_STATUS_IDLE);
50 - break;
51 - case alarm_log_status_aclk::ALARM_LOG_STATUS_RUNNING:
52 - msg.set_status(alarms::v1::ALARM_LOG_STATUS_RUNNING);
53 - break;
54 - case alarm_log_status_aclk::ALARM_LOG_STATUS_UNSPECIFIED:
55 - msg.set_status(alarms::v1::ALARM_LOG_STATUS_UNSPECIFIED);
56 - break;
57 - default:
58 - error("Unknown status of AlarmLogHealth LogEntry");
59 - return NULL;
60 - }
61 -
62 - entries = msg.mutable_log_entries();
63 - entries->set_first_sequence_id(data->log_entries.first_seq_id);
64 - entries->set_last_sequence_id(data->log_entries.last_seq_id);
31 + struct send_alarm_checkpoint ret;
32 + memset(&ret, 0, sizeof(ret));
33
66 - set_google_timestamp_from_timeval(data->log_entries.first_when, entries->mutable_first_when());
67 - set_google_timestamp_from_timeval(data->log_entries.last_when, entries->mutable_last_when());
34 + SendAlarmCheckpoint msg;
35 + if (!msg.ParseFromArray(data, len))
36 + return ret;
37
69 - *len = PROTO_COMPAT_MSG_SIZE(msg);
70 - char *bin = (char*)mallocz(*len);
71 - if (!msg.SerializeToArray(bin, *len))
72 - return NULL;
38 + ret.node_id = strdupz(msg.node_id().c_str());
39 + ret.claim_id = strdupz(msg.claim_id().c_str());
40
74 - return bin;
41 + return ret;
42 }
43
44 static alarms::v1::AlarmStatus aclk_alarm_status_to_proto(enum aclk_alarm_status status)
@@ -131,8 +98,6 @@ static void fill_alarm_log_entry(struct alarm_log_entry *data, AlarmLogEntry *pr
98 if (data->family)
99 proto->set_family(data->family);
100
134 - proto->set_batch_id(data->batch_id);
135 - proto->set_sequence_id(data->sequence_id);
101 proto->set_when(data->when);
102
103 proto->set_config_hash(data->config_hash);
@@ -187,6 +152,24 @@ char *generate_alarm_log_entry(size_t *len, struct alarm_log_entry *data)
152 return bin;
153 }
154
155 +char *generate_alarm_checkpoint(size_t *len, struct alarm_checkpoint *data)
156 +{
157 + AlarmCheckpoint msg;
158 +
159 + msg.set_claim_id(data->claim_id);
160 + msg.set_node_id(data->node_id);
161 + msg.set_checksum(data->checksum);
162 +
163 + *len = PROTO_COMPAT_MSG_SIZE(msg);
164 + char *bin = (char*)mallocz(*len);
165 + if (!msg.SerializeToArray(bin, *len)) {
166 + freez(bin);
167 + return NULL;
168 + }
169 +
170 + return bin;
171 +}
172 +
173 struct send_alarm_snapshot *parse_send_alarm_snapshot(const char *data, size_t len)
174 {
175 SendAlarmSnapshot msg;
@@ -198,8 +181,8 @@ struct send_alarm_snapshot *parse_send_alarm_snapshot(const char *data, size_t l
181 ret->claim_id = strdupz(msg.claim_id().c_str());
182 if (msg.node_id().c_str())
183 ret->node_id = strdupz(msg.node_id().c_str());
201 - ret->snapshot_id = msg.snapshot_id();
202 - ret->sequence_id = msg.sequence_id();
184 + if (msg.snapshot_uuid().c_str())
185 + ret->snapshot_uuid = strdupz(msg.snapshot_uuid().c_str());
186
187 return ret;
188 }
@@ -208,6 +191,7 @@ void destroy_send_alarm_snapshot(struct send_alarm_snapshot *ptr)
191 {
192 freez(ptr->claim_id);
193 freez(ptr->node_id);
194 + freez(ptr->snapshot_uuid);
195 freez(ptr);
196 }
197
@@ -218,7 +202,7 @@ alarm_snapshot_proto_ptr_t generate_alarm_snapshot_proto(struct alarm_snapshot *
202
203 msg->set_node_id(data->node_id);
204 msg->set_claim_id(data->claim_id);
221 - msg->set_snapshot_id(data->snapshot_id);
205 + msg->set_snapshot_uuid(data->snapshot_uuid);
206 msg->set_chunks(data->chunks);
207 msg->set_chunk(data->chunk);
208
aclk/schema-wrappers/alarm_stream.h
+17 -30
@@ -11,38 +11,12 @@
11 extern "C" {
12 #endif
13
14 -enum alarm_log_status_aclk {
15 - ALARM_LOG_STATUS_UNSPECIFIED = 0,
16 - ALARM_LOG_STATUS_RUNNING = 1,
17 - ALARM_LOG_STATUS_IDLE = 2
18 -};
19 -
20 -struct alarm_log_entries {
21 - int64_t first_seq_id;
22 - struct timeval first_when;
23 -
24 - int64_t last_seq_id;
25 - struct timeval last_when;
26 -};
27 -
28 -struct alarm_log_health {
29 - char *claim_id;
30 - char *node_id;
31 - int enabled;
32 - enum alarm_log_status_aclk status;
33 - struct alarm_log_entries log_entries;
34 -};
35 -
14 struct start_alarm_streaming {
15 char *node_id;
38 - uint64_t batch_id;
39 - uint64_t start_seq_id;
16 + bool resets;
17 };
18
19 struct start_alarm_streaming parse_start_alarm_streaming(const char *data, size_t len);
43 -char *parse_send_alarm_log_health(const char *data, size_t len);
44 -
45 -char *generate_alarm_log_health(size_t *len, struct alarm_log_health *data);
20
21 enum aclk_alarm_status {
22 ALARM_STATUS_NULL = 0,
@@ -101,17 +75,27 @@ struct alarm_log_entry {
75 char *chart_context;
76 };
77
78 +struct send_alarm_checkpoint {
79 + char *node_id;
80 + char *claim_id;
81 +};
82 +
83 +struct alarm_checkpoint {
84 + char *node_id;
85 + char *claim_id;
86 + char *checksum;
87 +};
88 +
89 struct send_alarm_snapshot {
90 char *node_id;
91 char *claim_id;
107 - uint64_t snapshot_id;
108 - uint64_t sequence_id;
92 + char *snapshot_uuid;
93 };
94
95 struct alarm_snapshot {
96 char *node_id;
97 char *claim_id;
114 - uint64_t snapshot_id;
98 + char *snapshot_uuid;
99 uint32_t chunks;
100 uint32_t chunk;
101 };
@@ -125,6 +109,9 @@ char *generate_alarm_log_entry(size_t *len, struct alarm_log_entry *data);
109 struct send_alarm_snapshot *parse_send_alarm_snapshot(const char *data, size_t len);
110 void destroy_send_alarm_snapshot(struct send_alarm_snapshot *ptr);
111
112 +struct send_alarm_checkpoint parse_send_alarm_checkpoint(const char *data, size_t len);
113 +char *generate_alarm_checkpoint(size_t *len, struct alarm_checkpoint *data);
114 +
115 alarm_snapshot_proto_ptr_t generate_alarm_snapshot_proto(struct alarm_snapshot *data);
116 void add_alarm_log_entry2snapshot(alarm_snapshot_proto_ptr_t snapshot, struct alarm_log_entry *data);
117 char *generate_alarm_snapshot_bin(size_t *len, alarm_snapshot_proto_ptr_t snapshot);
aclk/schema-wrappers/proto_2_json.cc
+4 -4
@@ -29,8 +29,8 @@ static google::protobuf::Message *msg_name_to_protomsg(const char *msgname)
29 return new nodeinstance::create::v1::CreateNodeInstance;
30 if (!strcmp(msgname, "UpdateNodeInfo"))
31 return new nodeinstance::info::v1::UpdateNodeInfo;
32 - if (!strcmp(msgname, "AlarmLogHealth"))
33 - return new alarms::v1::AlarmLogHealth;
32 + if (!strcmp(msgname, "AlarmCheckpoint"))
33 + return new alarms::v1::AlarmCheckpoint;
34 if (!strcmp(msgname, "ProvideAlarmConfiguration"))
35 return new alarms::v1::ProvideAlarmConfiguration;
36 if (!strcmp(msgname, "AlarmSnapshot"))
@@ -51,8 +51,8 @@ static google::protobuf::Message *msg_name_to_protomsg(const char *msgname)
51 return new agent::v1::SendNodeInstances;
52 if (!strcmp(msgname, "StartAlarmStreaming"))
53 return new alarms::v1::StartAlarmStreaming;
54 - if (!strcmp(msgname, "SendAlarmLogHealth"))
55 - return new alarms::v1::SendAlarmLogHealth;
54 + if (!strcmp(msgname, "SendAlarmCheckpoint"))
55 + return new alarms::v1::SendAlarmCheckpoint;
56 if (!strcmp(msgname, "SendAlarmConfiguration"))
57 return new alarms::v1::SendAlarmConfiguration;
58 if (!strcmp(msgname, "SendAlarmSnapshot"))
database/sqlite/sqlite_aclk.c
+1 -12
@@ -367,12 +367,12 @@ static void aclk_synchronization(void *arg __maybe_unused)
367 service_register(SERVICE_THREAD_TYPE_EVENT_LOOP, NULL, NULL, NULL, true);
368
369 worker_register_job_name(ACLK_DATABASE_NOOP, "noop");
370 - worker_register_job_name(ACLK_DATABASE_ALARM_HEALTH_LOG, "alert log");
370 worker_register_job_name(ACLK_DATABASE_CLEANUP, "cleanup");
371 worker_register_job_name(ACLK_DATABASE_DELETE_HOST, "node delete");
372 worker_register_job_name(ACLK_DATABASE_NODE_STATE, "node state");
373 worker_register_job_name(ACLK_DATABASE_PUSH_ALERT, "alert push");
374 worker_register_job_name(ACLK_DATABASE_PUSH_ALERT_CONFIG, "alert conf push");
375 + worker_register_job_name(ACLK_DATABASE_PUSH_ALERT_CHECKPOINT,"alert checkpoint");
376 worker_register_job_name(ACLK_DATABASE_PUSH_ALERT_SNAPSHOT, "alert snapshot");
377 worker_register_job_name(ACLK_DATABASE_QUEUE_REMOVED_ALERTS, "alerts check");
378 worker_register_job_name(ACLK_DATABASE_TIMER, "timer");
@@ -437,9 +437,6 @@ static void aclk_synchronization(void *arg __maybe_unused)
437 case ACLK_DATABASE_PUSH_ALERT:
438 aclk_push_alert_events_for_all_hosts();
439 break;
440 - case ACLK_DATABASE_ALARM_HEALTH_LOG:
441 - aclk_push_alarm_health_log(cmd.param[0]);
442 - break;
440 case ACLK_DATABASE_PUSH_ALERT_SNAPSHOT:;
441 aclk_push_alert_snapshot_event(cmd.param[0]);
442 break;
@@ -605,14 +602,6 @@ void aclk_push_node_alert_snapshot(const char *node_id)
602 }
603
604
608 -void aclk_push_node_health_log(const char *node_id)
609 -{
610 - if (unlikely(!aclk_sync_config.initialized))
611 - return;
612 -
613 - queue_aclk_sync_cmd(ACLK_DATABASE_ALARM_HEALTH_LOG, strdupz(node_id), NULL);
614 -}
615 -
605 void aclk_push_node_removed_alerts(const char *node_id)
606 {
607 if (unlikely(!aclk_sync_config.initialized))
database/sqlite/sqlite_aclk.h
+4 -5
@@ -41,13 +41,13 @@ static inline int claimed()
41 enum aclk_database_opcode {
42 ACLK_DATABASE_NOOP = 0,
43
44 - ACLK_DATABASE_ALARM_HEALTH_LOG,
44 ACLK_DATABASE_CLEANUP,
45 ACLK_DATABASE_DELETE_HOST,
46 ACLK_DATABASE_NODE_STATE,
47 ACLK_DATABASE_PUSH_ALERT,
48 ACLK_DATABASE_PUSH_ALERT_CONFIG,
49 ACLK_DATABASE_PUSH_ALERT_SNAPSHOT,
50 + ACLK_DATABASE_PUSH_ALERT_CHECKPOINT,
51 ACLK_DATABASE_QUEUE_REMOVED_ALERTS,
52 ACLK_DATABASE_TIMER,
53
@@ -72,14 +72,13 @@ struct aclk_database_cmdqueue {
72 struct aclk_sync_host_config {
73 RRDHOST *host;
74 int alert_updates;
75 + int alert_checkpoint_req;
76 + int alert_queue_removed;
77 time_t node_info_send_time;
78 time_t node_collectors_send;
79 char uuid_str[UUID_STR_LEN];
80 char node_id[UUID_STR_LEN];
79 - uint64_t alerts_batch_id; // batch id for alerts to use
80 - uint64_t alerts_start_seq_id; // cloud has asked to start streaming from
81 - uint64_t alerts_snapshot_id; // will contain the snapshot_id value if snapshot was requested
82 - uint64_t alerts_ack_sequence_id; // last sequence_id ack'ed from cloud via sendsnapshot message
81 + char *alerts_snapshot_uuid; // will contain the snapshot_uuid value if snapshot was requested
82 };
83
84 extern sqlite3 *db_meta;
database/sqlite/sqlite_aclk_alert.c
+178 -176
@@ -97,9 +97,6 @@ int should_send_to_cloud(RRDHOST *host, ALARM_ENTRY *ae)
97 if (unlikely(uuid_is_null(ae->config_hash_id)))
98 return 0;
99
100 - if (is_event_from_alert_variable_config(ae->unique_id, uuid_str))
101 - return 0;
102 -
100 char sql[ACLK_SYNC_QUERY_SIZE];
101 uuid_t config_hash_id;
102 RRDCALC_STATUS status;
@@ -193,10 +190,13 @@ int sql_queue_alarm_to_aclk(RRDHOST *host, ALARM_ENTRY *ae, int skip_filter)
190 }
191 }
192
196 - sqlite3_stmt *res_alert = NULL;
197 - char uuid_str[GUID_LEN + 1];
193 + char uuid_str[UUID_STR_LEN];
194 uuid_unparse_lower_fix(&host->host_uuid, uuid_str);
195
196 + if (is_event_from_alert_variable_config(ae->unique_id, uuid_str))
197 + return 0;
198 +
199 + sqlite3_stmt *res_alert = NULL;
200 char sql[ACLK_SYNC_QUERY_SIZE];
201
202 snprintfz(sql, ACLK_SYNC_QUERY_SIZE - 1, SQL_QUEUE_ALERT_TO_CLOUD, uuid_str);
@@ -278,26 +278,6 @@ void aclk_push_alert_event(struct aclk_sync_host_config *wc)
278
279 BUFFER *sql = buffer_create(1024, &netdata_buffers_statistics.buffers_sqlite);
280
281 - if (wc->alerts_start_seq_id != 0) {
282 - buffer_sprintf(
283 - sql,
284 - "UPDATE aclk_alert_%s SET date_submitted = NULL, date_cloud_ack = NULL WHERE sequence_id >= %"PRIu64
285 - "; UPDATE aclk_alert_%s SET date_cloud_ack = unixepoch() WHERE sequence_id < %"PRIu64
286 - " and date_cloud_ack is null "
287 - "; UPDATE aclk_alert_%s SET date_submitted = unixepoch() WHERE sequence_id < %"PRIu64
288 - " and date_submitted is null",
289 - wc->uuid_str,
290 - wc->alerts_start_seq_id,
291 - wc->uuid_str,
292 - wc->alerts_start_seq_id,
293 - wc->uuid_str,
294 - wc->alerts_start_seq_id);
295 - if (unlikely(db_execute(buffer_tostring(sql))))
296 - error_report("Failed to reset ACLK alert entries");
297 - buffer_reset(sql);
298 - wc->alerts_start_seq_id = 0;
299 - }
300 -
281 int limit = ACLK_MAX_ALERT_UPDATES;
282
283 sqlite3_stmt *res = NULL;
@@ -359,8 +339,8 @@ void aclk_push_alert_event(struct aclk_sync_host_config *wc)
339 alarm_log.name = strdupz((char *)sqlite3_column_text(res, 11));
340 alarm_log.family = sqlite3_column_bytes(res, 13) > 0 ? strdupz((char *)sqlite3_column_text(res, 13)) : NULL;
341
362 - alarm_log.batch_id = wc->alerts_batch_id;
363 - alarm_log.sequence_id = (uint64_t) sqlite3_column_int64(res, 0);
342 + //alarm_log.batch_id = wc->alerts_batch_id;
343 + //alarm_log.sequence_id = (uint64_t) sqlite3_column_int64(res, 0);
344 alarm_log.when = (time_t) sqlite3_column_int64(res, 5);
345
346 uuid_unparse_lower(*((uuid_t *) sqlite3_column_blob(res, 3)), uuid_str);
@@ -445,12 +425,11 @@ void aclk_push_alert_event(struct aclk_sync_host_config *wc)
425 } else {
426 if (log_first_sequence_id)
427 log_access(
448 - "ACLK RES [%s (%s)]: ALERTS SENT from %" PRIu64 " to %" PRIu64 " batch=%" PRIu64,
428 + "ACLK RES [%s (%s)]: ALERTS SENT from %" PRIu64 " to %" PRIu64 "",
429 wc->node_id,
430 wc->host ? rrdhost_hostname(wc->host) : "N/A",
431 log_first_sequence_id,
452 - log_last_sequence_id,
453 - wc->alerts_batch_id);
432 + log_last_sequence_id);
433 log_first_sequence_id = 0;
434 log_last_sequence_id = 0;
435 }
@@ -494,121 +473,17 @@ void sql_queue_existing_alerts_to_aclk(RRDHOST *host)
473 "where new_status <> 0 and new_status <> -2 and config_hash_id is not null and updated_by_id = 0 " \
474 "order by unique_id asc on conflict (alert_unique_id) do nothing;", uuid_str, uuid_str, uuid_str);
475
476 + netdata_rwlock_rdlock(&host->health_log.alarm_log_rwlock);
477 +
478 if (unlikely(db_execute(buffer_tostring(sql))))
479 error_report("Failed to queue existing ACLK alert events for host %s", rrdhost_hostname(host));
480
481 + netdata_rwlock_unlock(&host->health_log.alarm_log_rwlock);
482 +
483 buffer_free(sql);
484 rrdhost_flag_set(host, RRDHOST_FLAG_ACLK_STREAM_ALERTS);
485 }
486
504 -void aclk_send_alarm_health_log(char *node_id)
505 -{
506 - if (unlikely(!node_id))
507 - return;
508 -
509 - struct aclk_sync_host_config *wc = NULL;
510 - RRDHOST *host = find_host_by_node_id(node_id);
511 -
512 - if (unlikely(!host))
513 - return;
514 -
515 - wc = (struct aclk_sync_host_config *)host->aclk_sync_host_config;
516 - if (unlikely(!wc)) {
517 - log_access("ACLK REQ [%s (N/A)]: HEALTH LOG REQUEST RECEIVED FOR INVALID NODE", node_id);
518 - return;
519 - }
520 -
521 - log_access("ACLK REQ [%s (%s)]: HEALTH LOG REQUEST RECEIVED", node_id, rrdhost_hostname(host));
522 -
523 - aclk_push_node_health_log(node_id);
524 -}
525 -
526 -void aclk_push_alarm_health_log(char *node_id __maybe_unused)
527 -{
528 -#ifdef ENABLE_ACLK
529 - int rc;
530 -
531 - char *claim_id = get_agent_claimid();
532 - if (unlikely(!claim_id)) {
533 - freez(node_id);
534 - return;
535 - }
536 -
537 - RRDHOST *host = find_host_by_node_id(node_id);
538 -
539 - if (unlikely(!host)) {
540 - log_access("AC [%s (N/A)]: Node id not found", node_id);
541 - freez(claim_id);
542 - freez(node_id);
543 - return;
544 - }
545 -
546 - struct aclk_sync_host_config *wc = host->aclk_sync_host_config;
547 -
548 - int64_t first_sequence = 0;
549 - int64_t last_sequence = 0;
550 - struct timeval first_timestamp;
551 - struct timeval last_timestamp;
552 -
553 - char sql[ACLK_SYNC_QUERY_SIZE];
554 -
555 - sqlite3_stmt *res = NULL;
556 -
557 - //TODO: make this better: include info from health log too
558 - snprintfz(sql, ACLK_SYNC_QUERY_SIZE - 1, "SELECT MIN(sequence_id), MIN(date_created), " \
559 - "MAX(sequence_id), MAX(date_created) FROM aclk_alert_%s;", wc->uuid_str);
560 -
561 - rc = sqlite3_prepare_v2(db_meta, sql, -1, &res, 0);
562 - if (rc != SQLITE_OK) {
563 - error_report("Failed to prepare statement to get health log statistics from the database");
564 - freez(claim_id);
565 - return;
566 - }
567 -
568 - first_timestamp.tv_sec = 0;
569 - first_timestamp.tv_usec = 0;
570 - last_timestamp.tv_sec = 0;
571 - last_timestamp.tv_usec = 0;
572 -
573 - while (sqlite3_step_monitored(res) == SQLITE_ROW) {
574 - first_sequence = sqlite3_column_bytes(res, 0) > 0 ? (int64_t) sqlite3_column_int64(res, 0) : 0;
575 - if (sqlite3_column_bytes(res, 1) > 0) {
576 - first_timestamp.tv_sec = sqlite3_column_int64(res, 1);
577 - }
578 -
579 - last_sequence = sqlite3_column_bytes(res, 2) > 0 ? (int64_t) sqlite3_column_int64(res, 2) : 0;
580 - if (sqlite3_column_bytes(res, 3) > 0) {
581 - last_timestamp.tv_sec = sqlite3_column_int64(res, 3);
582 - }
583 - }
584 -
585 - struct alarm_log_entries log_entries;
586 - log_entries.first_seq_id = first_sequence;
587 - log_entries.first_when = first_timestamp;
588 - log_entries.last_seq_id = last_sequence;
589 - log_entries.last_when = last_timestamp;
590 -
591 - struct alarm_log_health alarm_log;
592 - alarm_log.claim_id = claim_id;
593 - alarm_log.node_id = wc->node_id;
594 - alarm_log.log_entries = log_entries;
595 - alarm_log.status = wc->alert_updates == 0 ? 2 : 1;
596 - alarm_log.enabled = (int)host->health.health_enabled;
597 -
598 - aclk_send_alarm_log_health(&alarm_log);
599 - log_access("ACLK RES [%s (%s)]: HEALTH LOG SENT from %ld to %ld", wc->node_id, rrdhost_hostname(host), first_sequence, last_sequence);
600 -
601 - rc = sqlite3_finalize(res);
602 - if (unlikely(rc != SQLITE_OK))
603 - error_report("Failed to reset statement to get health log statistics from the database, rc = %d", rc);
604 -
605 - freez(claim_id);
606 - freez(node_id);
607 -
608 - aclk_alert_reloaded = 1;
609 -#endif
610 -}
611 -
487 void aclk_send_alarm_configuration(char *config_hash)
488 {
489 if (unlikely(!config_hash))
@@ -755,7 +630,7 @@ bind_fail:
630
631
632 // Start streaming alerts
758 -void aclk_start_alert_streaming(char *node_id, uint64_t batch_id, uint64_t start_seq_id)
633 +void aclk_start_alert_streaming(char *node_id, bool resets)
634 {
635 if (unlikely(!node_id))
636 return;
@@ -779,20 +654,22 @@ void aclk_start_alert_streaming(char *node_id, uint64_t batch_id, uint64_t start
654 return;
655 }
656
782 - // TODO: CHECK
783 - if (unlikely(batch_id == 1) && unlikely(start_seq_id == 1))
657 + if (resets) {
658 + log_access("ACLK REQ [%s (%s)]: STREAM ALERTS ENABLED (RESET REQUESTED)", node_id, wc->host ? rrdhost_hostname(wc->host) : "N/A");
659 sql_queue_existing_alerts_to_aclk(host);
660 + } else
661 + log_access("ACLK REQ [%s (%s)]: STREAM ALERTS ENABLED", node_id, wc->host ? rrdhost_hostname(wc->host) : "N/A");
662
786 - log_access("ACLK REQ [%s (%s)]: ALERTS STREAM from %"PRIu64" batch=%"PRIu64, node_id, wc->host ? rrdhost_hostname(wc->host) : "N/A", start_seq_id, batch_id);
787 - wc->alerts_batch_id = batch_id;
788 - wc->alerts_start_seq_id = start_seq_id;
663 wc->alert_updates = 1;
664 + wc->alert_queue_removed = SEND_REMOVED_AFTER_HEALTH_LOOPS;
665 }
666
667 #define SQL_QUEUE_REMOVE_ALERTS "INSERT INTO aclk_alert_%s (alert_unique_id, date_created, filtered_alert_unique_id) " \
668 "SELECT unique_id alert_unique_id, UNIXEPOCH(), unique_id alert_unique_id FROM health_log_%s " \
669 "WHERE new_status = -2 AND updated_by_id = 0 AND unique_id NOT IN " \
795 - "(SELECT alert_unique_id FROM aclk_alert_%s) ORDER BY unique_id ASC " \
670 + "(SELECT alert_unique_id FROM aclk_alert_%s) " \
671 + "AND config_hash_id NOT IN (select hash_id from alert_hash where warn is null and crit is null) " \
672 + "ORDER BY unique_id ASC " \
673 "ON CONFLICT (alert_unique_id) DO NOTHING;"
674
675 void sql_process_queue_removed_alerts_to_aclk(char *node_id)
@@ -814,7 +691,9 @@ void sql_process_queue_removed_alerts_to_aclk(char *node_id)
691 }
692 else
693 log_access("ACLK STA [%s (%s)]: QUEUED REMOVED ALERTS", wc->node_id, rrdhost_hostname(wc->host));
694 +
695 rrdhost_flag_set(wc->host, RRDHOST_FLAG_ACLK_STREAM_ALERTS);
696 + wc->alert_queue_removed = 0;
697 }
698
699 void sql_queue_removed_alerts_to_aclk(RRDHOST *host)
@@ -831,7 +710,7 @@ void sql_queue_removed_alerts_to_aclk(RRDHOST *host)
710 aclk_push_node_removed_alerts(node_id);
711 }
712
834 -void aclk_process_send_alarm_snapshot(char *node_id, char *claim_id __maybe_unused, uint64_t snapshot_id, uint64_t sequence_id)
713 +void aclk_process_send_alarm_snapshot(char *node_id, char *claim_id __maybe_unused, char *snapshot_uuid)
714 {
715 uuid_t node_uuid;
716 if (unlikely(!node_id || uuid_parse(node_id, node_uuid)))
@@ -851,19 +730,17 @@ void aclk_process_send_alarm_snapshot(char *node_id, char *claim_id __maybe_unus
730 }
731
732 log_access(
854 - "IN [%s (%s)]: Request to send alerts snapshot, snapshot_id %" PRIu64 " and ack_sequence_id %" PRIu64,
855 - wc->node_id,
733 + "IN [%s (%s)]: Request to send alerts snapshot, snapshot_uuid %s",
734 + node_id,
735 wc->host ? rrdhost_hostname(wc->host) : "N/A",
857 - snapshot_id,
858 - sequence_id);
859 - if (wc->alerts_snapshot_id == snapshot_id)
860 - return;
861 - __sync_synchronize();
862 - wc->alerts_snapshot_id = snapshot_id;
863 - wc->alerts_ack_sequence_id = sequence_id;
864 - __sync_synchronize();
736 + snapshot_uuid);
737 + if (wc->alerts_snapshot_uuid && !strcmp(wc->alerts_snapshot_uuid,snapshot_uuid))
738 + return;
739 + __sync_synchronize();
740 + wc->alerts_snapshot_uuid = strdupz(snapshot_uuid);
741 + __sync_synchronize();
742
866 - aclk_push_node_alert_snapshot(node_id);
743 + aclk_push_node_alert_snapshot(node_id);
744 }
745
746 #ifdef ENABLE_ACLK
@@ -934,8 +811,6 @@ static int have_recent_alarm(RRDHOST *host, uint32_t alarm_id, uint32_t mark)
811 #endif
812
813 #define ALARM_EVENTS_PER_CHUNK 10
937 -#define SQL_ALERT_CLOUD_ACK "UPDATE aclk_alert_%s SET date_cloud_ack=unixepoch() WHERE sequence_id <= %"PRIu64
938 -
814 void aclk_push_alert_snapshot_event(char *node_id __maybe_unused)
815 {
816 #ifdef ENABLE_ACLK
@@ -959,21 +834,14 @@ void aclk_push_alert_snapshot_event(char *node_id __maybe_unused)
834 return;
835 }
836
962 - if (unlikely(!wc->alerts_snapshot_id))
837 + if (unlikely(!wc->alerts_snapshot_uuid))
838 return;
839
840 char *claim_id = get_agent_claimid();
841 if (unlikely(!claim_id))
842 return;
843
969 - log_access("ACLK REQ [%s (%s)]: Sending alerts snapshot, snapshot_id %" PRIu64, wc->node_id, rrdhost_hostname(wc->host), wc->alerts_snapshot_id);
970 -
971 - if (wc->alerts_ack_sequence_id) {
972 - char sql[512];
973 - snprintfz(sql, 511, SQL_ALERT_CLOUD_ACK, wc->uuid_str, wc->alerts_ack_sequence_id);
974 - if (unlikely(db_execute(sql)))
975 - error_report("Failed to set ACLK alert entries cloud ACK status for host %s", rrdhost_hostname(host));
976 - }
844 + log_access("ACLK REQ [%s (%s)]: Sending alerts snapshot, snapshot_uuid %s", wc->node_id, rrdhost_hostname(wc->host), wc->alerts_snapshot_uuid);
845
846 uint32_t cnt = 0;
847 char uuid_str[UUID_STR_LEN];
@@ -1009,7 +877,7 @@ void aclk_push_alert_snapshot_event(char *node_id __maybe_unused)
877 struct alarm_snapshot alarm_snap;
878 alarm_snap.node_id = wc->node_id;
879 alarm_snap.claim_id = claim_id;
1012 - alarm_snap.snapshot_id = wc->alerts_snapshot_id;
880 + alarm_snap.snapshot_uuid = wc->alerts_snapshot_uuid;
881 alarm_snap.chunks = chunks;
882 alarm_snap.chunk = chunk;
883
@@ -1051,7 +919,7 @@ void aclk_push_alert_snapshot_event(char *node_id __maybe_unused)
919 struct alarm_snapshot alarm_snap;
920 alarm_snap.node_id = wc->node_id;
921 alarm_snap.claim_id = claim_id;
1054 - alarm_snap.snapshot_id = wc->alerts_snapshot_id;
922 + alarm_snap.snapshot_uuid = wc->alerts_snapshot_uuid;
923 alarm_snap.chunks = chunks;
924 alarm_snap.chunk = chunk;
925
@@ -1065,7 +933,7 @@ void aclk_push_alert_snapshot_event(char *node_id __maybe_unused)
933 }
934
935 netdata_rwlock_unlock(&host->health_log.alarm_log_rwlock);
1068 - wc->alerts_snapshot_id = 0;
936 + wc->alerts_snapshot_uuid = NULL;
937
938 freez(claim_id);
939 #endif
@@ -1093,7 +961,6 @@ void sql_aclk_alert_clean_dead_entries(RRDHOST *host)
961 }
962
963 #define SQL_GET_MIN_MAX_ALERT_SEQ "SELECT MIN(sequence_id), MAX(sequence_id), " \
1096 - "(SELECT MAX(sequence_id) FROM aclk_alert_%s WHERE date_cloud_ack IS NOT NULL), " \
964 "(SELECT MAX(sequence_id) FROM aclk_alert_%s WHERE date_submitted IS NOT NULL) " \
965 "FROM aclk_alert_%s WHERE date_submitted IS NULL;"
966
@@ -1106,12 +973,11 @@ int get_proto_alert_status(RRDHOST *host, struct proto_alert_status *proto_alert
973 return 1;
974
975 proto_alert_status->alert_updates = wc->alert_updates;
1109 - proto_alert_status->alerts_batch_id = wc->alerts_batch_id;
976
977 char sql[ACLK_SYNC_QUERY_SIZE];
978 sqlite3_stmt *res = NULL;
979
1114 - snprintfz(sql, ACLK_SYNC_QUERY_SIZE - 1, SQL_GET_MIN_MAX_ALERT_SEQ, wc->uuid_str, wc->uuid_str, wc->uuid_str);
980 + snprintfz(sql, ACLK_SYNC_QUERY_SIZE - 1, SQL_GET_MIN_MAX_ALERT_SEQ, wc->uuid_str, wc->uuid_str);
981
982 rc = sqlite3_prepare_v2(db_meta, sql, -1, &res, 0);
983 if (rc != SQLITE_OK) {
@@ -1122,8 +988,7 @@ int get_proto_alert_status(RRDHOST *host, struct proto_alert_status *proto_alert
988 while (sqlite3_step_monitored(res) == SQLITE_ROW) {
989 proto_alert_status->pending_min_sequence_id = sqlite3_column_bytes(res, 0) > 0 ? (uint64_t) sqlite3_column_int64(res, 0) : 0;
990 proto_alert_status->pending_max_sequence_id = sqlite3_column_bytes(res, 1) > 0 ? (uint64_t) sqlite3_column_int64(res, 1) : 0;
1125 - proto_alert_status->last_acked_sequence_id = sqlite3_column_bytes(res, 2) > 0 ? (uint64_t) sqlite3_column_int64(res, 2) : 0;
1126 - proto_alert_status->last_submitted_sequence_id = sqlite3_column_bytes(res, 3) > 0 ? (uint64_t) sqlite3_column_int64(res, 3) : 0;
991 + proto_alert_status->last_submitted_sequence_id = sqlite3_column_bytes(res, 2) > 0 ? (uint64_t) sqlite3_column_int64(res, 2) : 0;
992 }
993
994 rc = sqlite3_finalize(res);
@@ -1132,3 +997,140 @@ int get_proto_alert_status(RRDHOST *host, struct proto_alert_status *proto_alert
997
998 return 0;
999 }
1000 +
1001 +void aclk_send_alarm_checkpoint(char *node_id, char *claim_id __maybe_unused)
1002 +{
1003 + if (unlikely(!node_id))
1004 + return;
1005 +
1006 + struct aclk_sync_host_config *wc = NULL;
1007 + RRDHOST *host = find_host_by_node_id(node_id);
1008 +
1009 + if (unlikely(!host))
1010 + return;
1011 +
1012 + wc = (struct aclk_sync_host_config *)host->aclk_sync_host_config;
1013 + if (unlikely(!wc)) {
1014 + log_access("ACLK REQ [%s (N/A)]: ALERTS CHECKPOINT REQUEST RECEIVED FOR INVALID NODE", node_id);
1015 + return;
1016 + }
1017 +
1018 + log_access("ACLK REQ [%s (%s)]: ALERTS CHECKPOINT REQUEST RECEIVED", node_id, rrdhost_hostname(host));
1019 +
1020 + wc->alert_checkpoint_req = SEND_CHECKPOINT_AFTER_HEALTH_LOOPS;
1021 +}
1022 +
1023 +typedef struct active_alerts {
1024 + char *name;
1025 + char *chart;
1026 + RRDCALC_STATUS status;
1027 +} active_alerts_t;
1028 +
1029 +static inline int compare_active_alerts(const void * a, const void * b) {
1030 + active_alerts_t *active_alerts_a = (active_alerts_t *)a;
1031 + active_alerts_t *active_alerts_b = (active_alerts_t *)b;
1032 +
1033 + if( !(strcmp(active_alerts_a->name, active_alerts_b->name)) )
1034 + {
1035 + return strcmp(active_alerts_a->chart, active_alerts_b->chart);
1036 + }
1037 + else
1038 + return strcmp(active_alerts_a->name, active_alerts_b->name);
1039 +}
1040 +
1041 +void aclk_push_alarm_checkpoint(RRDHOST *host __maybe_unused)
1042 +{
1043 +#ifdef ENABLE_ACLK
1044 + struct aclk_sync_host_config *wc = host->aclk_sync_host_config;
1045 + if (unlikely(!wc)) {
1046 + log_access("ACLK REQ [%s (N/A)]: ALERTS CHECKPOINT REQUEST RECEIVED FOR INVALID NODE", rrdhost_hostname(host));
1047 + return;
1048 + }
1049 +
1050 + //TODO: make sure all pending events are sent.
1051 + if (rrdhost_flag_check(host, RRDHOST_FLAG_ACLK_STREAM_ALERTS)) {
1052 + //postpone checkpoint send
1053 + wc->alert_checkpoint_req++;
1054 + log_access("ACLK REQ [%s (N/A)]: ALERTS CHECKPOINT POSTPONED", rrdhost_hostname(host));
1055 + return;
1056 + }
1057 +
1058 + //TODO: lock rc here, or make sure it's called when health decides
1059 + //count them
1060 + RRDCALC *rc;
1061 + uint32_t cnt = 0;
1062 + size_t len = 0;
1063 + active_alerts_t *active_alerts = NULL;
1064 +
1065 + foreach_rrdcalc_in_rrdhost_read(host, rc) {
1066 + if(unlikely(!rc->rrdset || !rc->rrdset->last_collected_time.tv_sec))
1067 + continue;
1068 +
1069 + if (rc->status == RRDCALC_STATUS_WARNING ||
1070 + rc->status == RRDCALC_STATUS_CRITICAL) {
1071 +
1072 + cnt++;
1073 + }
1074 + }
1075 + foreach_rrdcalc_in_rrdhost_done(rc);
1076 +
1077 + if (cnt) {
1078 + active_alerts = callocz(cnt, sizeof(active_alerts_t));
1079 + cnt = 0;
1080 + foreach_rrdcalc_in_rrdhost_read(host, rc) {
1081 + if(unlikely(!rc->rrdset || !rc->rrdset->last_collected_time.tv_sec))
1082 + continue;
1083 +
1084 + if (rc->status == RRDCALC_STATUS_WARNING ||
1085 + rc->status == RRDCALC_STATUS_CRITICAL) {
1086 +
1087 + active_alerts[cnt].name = (char *)rrdcalc_name(rc);
1088 + len += string_strlen(rc->name);
1089 + active_alerts[cnt].chart = (char *)rrdcalc_chart_name(rc);
1090 + len += string_strlen(rc->chart);
1091 + active_alerts[cnt].status = rc->status;
1092 + len++;
1093 + cnt++;
1094 + }
1095 + }
1096 + foreach_rrdcalc_in_rrdhost_done(rc);
1097 + }
1098 +
1099 + BUFFER *alarms_to_hash;
1100 + if (cnt) {
1101 + qsort (active_alerts, cnt, sizeof(active_alerts_t), compare_active_alerts);
1102 +
1103 + alarms_to_hash = buffer_create(len, NULL);
1104 + for (uint32_t i=0;i<cnt;i++) {
1105 + buffer_strcat(alarms_to_hash, active_alerts[i].name);
1106 + buffer_strcat(alarms_to_hash, active_alerts[i].chart);
1107 + if (active_alerts[i].status == RRDCALC_STATUS_WARNING)
1108 + buffer_strcat(alarms_to_hash, "W");
1109 + else if (active_alerts[i].status == RRDCALC_STATUS_CRITICAL)
1110 + buffer_strcat(alarms_to_hash, "C");
1111 + }
1112 + } else {
1113 + alarms_to_hash = buffer_create(1, NULL);
1114 + buffer_strcat(alarms_to_hash, "");
1115 + len = 0;
1116 + }
1117 +
1118 + char hash[SHA256_DIGEST_LENGTH + 1];
1119 + if (hash256_string((const unsigned char *)buffer_tostring(alarms_to_hash), len, hash)) {
1120 + hash[SHA256_DIGEST_LENGTH] = 0;
1121 +
1122 + struct alarm_checkpoint alarm_checkpoint;
1123 + char *claim_id = get_agent_claimid();
1124 + alarm_checkpoint.claim_id = claim_id;
1125 + alarm_checkpoint.node_id = wc->node_id;
1126 + alarm_checkpoint.checksum = (char *)hash;
1127 +
1128 + aclk_send_provide_alarm_checkpoint(&alarm_checkpoint);
1129 + log_access("ACLK RES [%s (%s)]: ALERTS CHECKPOINT SENT", wc->node_id, rrdhost_hostname(host));
1130 + } else {
1131 + log_access("ACLK RES [%s (%s)]: FAILED TO CREATE ALERTS CHECKPOINT HASH", wc->node_id, rrdhost_hostname(host));
1132 + }
1133 + wc->alert_checkpoint_req = 0;
1134 + buffer_free(alarms_to_hash);
1135 +#endif
1136 +}
database/sqlite/sqlite_aclk_alert.h
+7 -6
@@ -5,10 +5,11 @@
5
6 extern sqlite3 *db_meta;
7
8 +#define SEND_REMOVED_AFTER_HEALTH_LOOPS 3
9 +#define SEND_CHECKPOINT_AFTER_HEALTH_LOOPS 4
10 +
11 struct proto_alert_status {
12 int alert_updates;
10 - uint64_t alerts_batch_id;
11 - uint64_t last_acked_sequence_id;
13 uint64_t pending_min_sequence_id;
14 uint64_t pending_max_sequence_id;
15 uint64_t last_submitted_sequence_id;
@@ -16,16 +17,16 @@ struct proto_alert_status {
17
18 int aclk_add_alert_event(struct aclk_sync_host_config *wc, struct aclk_database_cmd cmd);
19 void aclk_push_alert_event(struct aclk_sync_host_config *wc);
19 -void aclk_send_alarm_health_log(char *node_id);
20 -void aclk_push_alarm_health_log(char *node_id);
20 void aclk_send_alarm_configuration (char *config_hash);
21 int aclk_push_alert_config_event(char *node_id, char *config_hash);
23 -void aclk_start_alert_streaming(char *node_id, uint64_t batch_id, uint64_t start_seq_id);
22 +void aclk_start_alert_streaming(char *node_id, bool resets);
23 void sql_queue_removed_alerts_to_aclk(RRDHOST *host);
24 void sql_process_queue_removed_alerts_to_aclk(char *node_id);
25 +void aclk_send_alarm_checkpoint(char *node_id, char *claim_id);
26 +void aclk_push_alarm_checkpoint(RRDHOST *host);
27
28 void aclk_push_alert_snapshot_event(char *node_id);
28 -void aclk_process_send_alarm_snapshot(char *node_id, char *claim_id, uint64_t snapshot_id, uint64_t sequence_id);
29 +void aclk_process_send_alarm_snapshot(char *node_id, char *claim_id, char *snapshot_uuid);
30 int get_proto_alert_status(RRDHOST *host, struct proto_alert_status *proto_alert_status);
31 int sql_queue_alarm_to_aclk(RRDHOST *host, ALARM_ENTRY *ae, int skip_filter);
32 void aclk_push_alert_events_for_all_hosts(void);
health/health.c
+31 -32
@@ -347,6 +347,15 @@ static void health_reload_host(RRDHOST *host) {
347 rrdcalctemplate_link_matching_templates_to_rrdset(st);
348 }
349 rrdset_foreach_done(st);
350 +
351 +#ifdef ENABLE_ACLK
352 + if (netdata_cloud_setting) {
353 + struct aclk_sync_host_config *wc = (struct aclk_sync_host_config *)host->aclk_sync_host_config;
354 + if (likely(wc)) {
355 + wc->alert_queue_removed = SEND_REMOVED_AFTER_HEALTH_LOOPS;
356 + }
357 + }
358 +#endif
359 }
360
361 /**
@@ -362,12 +371,6 @@ void health_reload(void) {
371 health_reload_host(host);
372 }
373 dfe_done(host);
365 -
366 -#ifdef ENABLE_ACLK
367 - if (netdata_cloud_setting) {
368 - aclk_alert_reloaded = 1;
369 - }
370 -#endif
374 }
375
376 // ----------------------------------------------------------------------------
@@ -977,9 +980,7 @@ void *health_main(void *ptr) {
980 rrdcalc_delete_alerts_not_matching_host_labels_from_all_hosts();
981
982 unsigned int loop = 0;
980 -#ifdef ENABLE_ACLK
981 - unsigned int marked_aclk_reload_loop = 0;
982 -#endif
983 +
984 while(service_running(SERVICE_HEALTH)) {
985 loop++;
986 debug(D_HEALTH, "Health monitoring iteration no %u started", loop);
@@ -1008,11 +1009,6 @@ void *health_main(void *ptr) {
1009 }
1010 }
1011
1011 -#ifdef ENABLE_ACLK
1012 - if (aclk_alert_reloaded && !marked_aclk_reload_loop)
1013 - marked_aclk_reload_loop = loop;
1014 -#endif
1015 -
1012 worker_is_busy(WORKER_HEALTH_JOB_RRD_LOCK);
1013 dfe_start_reentrant(rrdhost_root_index, host) {
1014
@@ -1117,7 +1113,7 @@ void *health_main(void *ptr) {
1113 rc->value = NAN;
1114
1115 #ifdef ENABLE_ACLK
1120 - if (netdata_cloud_setting && likely(!aclk_alert_reloaded))
1116 + if (netdata_cloud_setting)
1117 sql_queue_alarm_to_aclk(host, ae, 1);
1118 #endif
1119 }
@@ -1488,6 +1484,26 @@ void *health_main(void *ptr) {
1484 }
1485 break;
1486 }
1487 +#ifdef ENABLE_ACLK
1488 + if (netdata_cloud_setting) {
1489 + struct aclk_sync_host_config *wc = (struct aclk_sync_host_config *)host->aclk_sync_host_config;
1490 + if (unlikely(!wc)) {
1491 + continue;
1492 + }
1493 +
1494 + if (wc->alert_queue_removed == 1) {
1495 + sql_queue_removed_alerts_to_aclk(host);
1496 + } else if (wc->alert_queue_removed > 1) {
1497 + wc->alert_queue_removed--;
1498 + }
1499 +
1500 + if (wc->alert_checkpoint_req == 1) {
1501 + aclk_push_alarm_checkpoint(host);
1502 + } else if (wc->alert_checkpoint_req > 1) {
1503 + wc->alert_checkpoint_req--;
1504 + }
1505 + }
1506 +#endif
1507 }
1508 dfe_done(host);
1509
@@ -1500,23 +1516,6 @@ void *health_main(void *ptr) {
1516 health_alarm_wait_for_execution(ae);
1517 }
1518
1503 -#ifdef ENABLE_ACLK
1504 - if (netdata_cloud_setting && unlikely(aclk_alert_reloaded) && loop > (marked_aclk_reload_loop + 2)) {
1505 - dfe_start_reentrant(rrdhost_root_index, host) {
1506 - if(unlikely(!service_running(SERVICE_HEALTH)))
1507 - break;
1508 -
1509 - if (unlikely(!host->health.health_enabled))
1510 - continue;
1511 -
1512 - sql_queue_removed_alerts_to_aclk(host);
1513 - }
1514 - dfe_done(host);
1515 - aclk_alert_reloaded = 0;
1516 - marked_aclk_reload_loop = 0;
1517 - }
1518 -#endif
1519 -
1519 if(unlikely(!service_running(SERVICE_HEALTH)))
1520 break;
1521
libnetdata/libnetdata.c
+25
@@ -2009,3 +2009,28 @@ void timing_action(TIMING_ACTION action, TIMING_STEP step) {
2009 }
2010 }
2011 }
2012 +
2013 +int hash256_string(const unsigned char *string, size_t size, char *hash) {
2014 + EVP_MD_CTX *ctx;
2015 + ctx = EVP_MD_CTX_create();
2016 +
2017 + if (!ctx)
2018 + return 0;
2019 +
2020 + if (!EVP_DigestInit(ctx, EVP_sha256())) {
2021 + EVP_MD_CTX_destroy(ctx);
2022 + return 0;
2023 + }
2024 +
2025 + if (!EVP_DigestUpdate(ctx, string, size)) {
2026 + EVP_MD_CTX_destroy(ctx);
2027 + return 0;
2028 + }
2029 +
2030 + if (!EVP_DigestFinal(ctx, (unsigned char *)hash, NULL)) {
2031 + EVP_MD_CTX_destroy(ctx);
2032 + return 0;
2033 + }
2034 +
2035 + return 1;
2036 +}
libnetdata/libnetdata.h
+1
@@ -782,6 +782,7 @@ typedef enum {
782 #endif
783 void timing_action(TIMING_ACTION action, TIMING_STEP step);
784
785 +int hash256_string(const unsigned char *string, size_t size, char *hash);
786 # ifdef __cplusplus
787 }
788 # endif