@cryptotaxi247 / netdata-1 / commits / 4ae319931

Add alert message support for ACLK new architecture (#11552)

* add alert messages * also clear date_cloud_ack * move buffer_create * remove include file * use wc->node_id

Emmanuel Vasilakis committed Sep 23, 2021 at 17:34 UTC 4ae31993115f3d796e31e931441c4bfd371dcbfa
8 files changed +571 -10
CMakeLists.txt
+2
@@ -637,6 +637,8 @@ set(RRD_PLUGIN_FILES
637 database/sqlite/sqlite_aclk_node.h
638 database/sqlite/sqlite_aclk_chart.c
639 database/sqlite/sqlite_aclk_chart.h
640 + database/sqlite/sqlite_aclk_alert.c
641 + database/sqlite/sqlite_aclk_alert.h
642 database/sqlite/sqlite3.c
643 database/sqlite/sqlite3.h
644 database/engine/rrdengine.c
Makefile.am
+2
@@ -417,6 +417,8 @@ RRD_PLUGIN_FILES = \
417 database/sqlite/sqlite_aclk_node.h \
418 database/sqlite/sqlite_aclk_chart.c \
419 database/sqlite/sqlite_aclk_chart.h \
420 + database/sqlite/sqlite_aclk_alert.c \
421 + database/sqlite/sqlite_aclk_alert.h \
422 database/sqlite/sqlite3.c \
423 database/sqlite/sqlite3.h \
424 $(NULL)
database/sqlite/sqlite_aclk.h
+1 -1
@@ -102,7 +102,7 @@ static inline char *get_str_from_uuid(uuid_t *uuid)
102 "end;"
103
104 #define TABLE_ACLK_ALERT "CREATE TABLE IF NOT EXISTS aclk_alert_%s (sequence_id INTEGER PRIMARY KEY, " \
105 - "alert_unique_id, date_created, date_submitted, " \
105 + "alert_unique_id, date_created, date_submitted, date_cloud_ack, " \
106 "unique(alert_unique_id)); " \
107 "insert into aclk_alert_%s (alert_unique_id, date_created) " \
108 "select unique_id alert_unique_id, strftime('%%s') date_created from health_log_%s where new_status <> 0 and new_status <> -2 order by unique_id asc on conflict (alert_unique_id) do nothing;"
database/sqlite/sqlite_aclk_alert.c new
+538
@@ -0,0 +1,538 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#include "sqlite_functions.h"
4 +#include "sqlite_aclk_alert.h"
5 +
6 +#include "../../aclk/aclk_alarm_api.h"
7 +
8 +// will replace call to aclk_update_alarm in health/health_log.c
9 +// and handle both cases
10 +void sql_queue_alarm_to_aclk(RRDHOST *host, ALARM_ENTRY *ae)
11 +{
12 + //check aclk architecture and handle old json alarm update to cloud
13 + //include also the valid statuses for this case
14 + /* if (!aclk_architecture)
15 + aclk_update_alarm(host, ae); */
16 +
17 + if (ae->flags & HEALTH_ENTRY_FLAG_ACLK_QUEUED)
18 + return;
19 +
20 + if (ae->new_status == RRDCALC_STATUS_REMOVED || ae->new_status == RRDCALC_STATUS_UNINITIALIZED)
21 + return;
22 +
23 + if (unlikely(!host->dbsync_worker))
24 + return;
25 +
26 + if (unlikely(uuid_is_null(ae->config_hash_id)))
27 + return;
28 +
29 + struct aclk_database_cmd cmd;
30 + memset(&cmd, 0, sizeof(cmd));
31 + cmd.opcode = ACLK_DATABASE_ADD_ALERT;
32 + cmd.data = ae;
33 + cmd.completion = NULL;
34 + aclk_database_enq_cmd((struct aclk_database_worker_config *) host->dbsync_worker, &cmd);
35 + ae->flags |= HEALTH_ENTRY_FLAG_ACLK_QUEUED;
36 + return;
37 +}
38 +
39 +// stores an alert entry to aclk_alert_ table
40 +int aclk_add_alert_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
41 +{
42 + int rc = 0;
43 +
44 + CHECK_SQLITE_CONNECTION(db_meta);
45 +
46 + sqlite3_stmt *res_alert = NULL;
47 + ALARM_ENTRY *ae = cmd.data;
48 +
49 + BUFFER *sql = buffer_create(1024);
50 +
51 + buffer_sprintf(
52 + sql,
53 + "INSERT INTO aclk_alert_%s (alert_unique_id, date_created) "
54 + "VALUES (@alert_unique_id, strftime('%%s')); ",
55 + wc->uuid_str);
56 +
57 + rc = sqlite3_prepare_v2(db_meta, buffer_tostring(sql), -1, &res_alert, 0);
58 + if (unlikely(rc != SQLITE_OK)) {
59 + error_report("Failed to prepare statement to store alert event");
60 + buffer_free(sql);
61 + return 1;
62 + }
63 +
64 + rc = sqlite3_bind_int(res_alert, 1, ae->unique_id);
65 + if (unlikely(rc != SQLITE_OK))
66 + goto bind_fail;
67 +
68 + rc = execute_insert(res_alert);
69 + if (unlikely(rc != SQLITE_DONE))
70 + error_report("Failed to store alert event %u, rc = %d", ae->unique_id, rc);
71 +
72 +bind_fail:
73 + if (unlikely(sqlite3_finalize(res_alert) != SQLITE_OK))
74 + error_report("Failed to reset statement in store alert event, rc = %d", rc);
75 +
76 + buffer_free(sql);
77 + return (rc != SQLITE_DONE);
78 +}
79 +
80 +int rrdcalc_status_to_proto_enum(RRDCALC_STATUS status)
81 +{
82 + switch(status) {
83 + case RRDCALC_STATUS_REMOVED:
84 + return ALARM_STATUS_REMOVED;
85 +
86 + case RRDCALC_STATUS_UNDEFINED:
87 + return ALARM_STATUS_NOT_A_NUMBER;
88 +
89 + case RRDCALC_STATUS_CLEAR:
90 + return ALARM_STATUS_CLEAR;
91 +
92 + case RRDCALC_STATUS_WARNING:
93 + return ALARM_STATUS_WARNING;
94 +
95 + case RRDCALC_STATUS_CRITICAL:
96 + return ALARM_STATUS_CRITICAL;
97 +
98 + default:
99 + return ALARM_STATUS_UNKNOWN;
100 + }
101 +}
102 +
103 +void aclk_push_alert_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
104 +{
105 +#ifndef ACLK_NG
106 + UNUSED(wc);
107 + UNUSED(cmd);
108 +#else
109 + int rc;
110 +
111 + if (unlikely(!wc->alert_updates)) {
112 + debug(D_ACLK_SYNC,"Ignoring alert push event, updates have been turned off for node %s", wc->node_id);
113 + return;
114 + }
115 +
116 + char *claim_id = is_agent_claimed();
117 + if (unlikely(!claim_id))
118 + return;
119 +
120 + BUFFER *sql = buffer_create(1024);
121 +
122 + if (wc->alerts_start_seq_id != 0) {
123 + buffer_sprintf(
124 + sql,
125 + "UPDATE aclk_alert_%s SET date_submitted = NULL, date_cloud_ack = NULL WHERE sequence_id >= %" PRIu64
126 + "; UPDATE aclk_alert_%s SET date_cloud_ack = strftime('%%s','now') WHERE sequence_id < %" PRIu64
127 + " and date_cloud_ack is null",
128 + wc->uuid_str,
129 + wc->alerts_start_seq_id,
130 + wc->uuid_str,
131 + wc->alerts_start_seq_id);
132 + db_execute(buffer_tostring(sql));
133 + buffer_reset(sql);
134 + wc->alerts_start_seq_id = 0;
135 + }
136 +
137 + int limit = cmd.count > 0 ? cmd.count : 1;
138 +
139 + sqlite3_stmt *res = NULL;
140 +
141 + buffer_sprintf(sql, "select aa.sequence_id, hl.unique_id, hl.alarm_id, hl.config_hash_id, hl.updated_by_id, hl.when_key, \
142 + hl.duration, hl.non_clear_duration, hl.flags, hl.exec_run_timestamp, hl.delay_up_to_timestamp, hl.name, \
143 + hl.chart, hl.family, hl.exec, hl.recipient, hl.source, hl.units, hl.info, hl.exec_code, hl.new_status, \
144 + hl.old_status, hl.delay, hl.new_value, hl.old_value, hl.last_repeat \
145 + from health_log_%s hl, aclk_alert_%s aa \
146 + where hl.unique_id = aa.alert_unique_id and aa.date_submitted is null \
147 + order by aa.sequence_id asc limit %d;", wc->uuid_str, wc->uuid_str, limit);
148 +
149 + rc = sqlite3_prepare_v2(db_meta, buffer_tostring(sql), -1, &res, 0);
150 + if (rc != SQLITE_OK) {
151 + error_report("Failed to prepare statement when trying to send an alert update via ACLK");
152 + buffer_free(sql);
153 + freez(claim_id);
154 + return;
155 + }
156 +
157 + char uuid_str[GUID_LEN + 1];
158 + uint64_t first_sequence_id = 0;
159 + uint64_t last_sequence_id = 0;
160 +
161 + while (sqlite3_step(res) == SQLITE_ROW) {
162 + struct alarm_log_entry alarm_log;
163 + char old_value_string[100 + 1];
164 + char new_value_string[100 + 1];
165 +
166 + alarm_log.node_id = strdupz(wc->node_id);
167 + alarm_log.claim_id = claim_id;
168 +
169 + alarm_log.chart = strdupz((char *)sqlite3_column_text(res, 12));
170 + alarm_log.name = strdupz((char *)sqlite3_column_text(res, 11));
171 + alarm_log.family = sqlite3_column_bytes(res, 13) > 0 ? strdupz((char *)sqlite3_column_text(res, 13)) : NULL;
172 +
173 + alarm_log.batch_id = wc->alerts_batch_id;
174 + alarm_log.sequence_id = (uint64_t) sqlite3_column_int64(res, 0);
175 + alarm_log.when = (time_t) sqlite3_column_int64(res, 5);
176 +
177 + uuid_unparse_lower(*((uuid_t *) sqlite3_column_blob(res, 3)), uuid_str);
178 + alarm_log.config_hash = strdupz((char *)uuid_str);
179 +
180 + alarm_log.utc_offset = wc->host->utc_offset;
181 + alarm_log.timezone = strdupz((char *)wc->host->abbrev_timezone);
182 + alarm_log.exec_path = sqlite3_column_bytes(res, 14) > 0 ? strdupz((char *)sqlite3_column_text(res, 14)) : strdupz((char *)wc->host->health_default_exec);
183 + alarm_log.conf_source = strdupz((char *)sqlite3_column_text(res, 16));
184 +
185 + char *edit_command = sqlite3_column_bytes(res, 16) > 0 ? health_edit_command_from_source((char *)sqlite3_column_text(res, 16)) : strdupz("UNKNOWN=0");
186 + alarm_log.command = strdupz(edit_command);
187 +
188 + alarm_log.duration = (time_t) sqlite3_column_int64(res, 6);
189 + alarm_log.non_clear_duration = (time_t) sqlite3_column_int64(res, 7);
190 + alarm_log.status = rrdcalc_status_to_proto_enum((RRDCALC_STATUS) sqlite3_column_int(res, 20));
191 + alarm_log.old_status = rrdcalc_status_to_proto_enum((RRDCALC_STATUS) sqlite3_column_int(res, 21));
192 + alarm_log.delay = (int) sqlite3_column_int(res, 22);
193 + alarm_log.delay_up_to_timestamp = (time_t) sqlite3_column_int64(res, 10);
194 + alarm_log.last_repeat = (time_t) sqlite3_column_int64(res, 25);
195 +
196 + alarm_log.silenced = ( (sqlite3_column_int64(res, 8) & HEALTH_ENTRY_FLAG_SILENCED) || ( sqlite3_column_type(res, 15) != SQLITE_NULL && !strncmp((char *)sqlite3_column_text(res,15), "silent", 6)) ) ? 1 : 0;
197 +
198 + alarm_log.value_string = sqlite3_column_type(res, 23) == SQLITE_NULL ? strdupz((char *)"-") : strdupz((char *)format_value_and_unit(new_value_string, 100, sqlite3_column_double(res, 23), (char *) sqlite3_column_text(res, 17), -1));
199 + alarm_log.old_value_string = sqlite3_column_type(res, 24) == SQLITE_NULL ? strdupz((char *)"-") : strdupz((char *)format_value_and_unit(old_value_string, 100, sqlite3_column_double(res, 24), (char *) sqlite3_column_text(res, 17), -1));
200 +
201 + alarm_log.value = (calculated_number) sqlite3_column_double(res, 23);
202 + alarm_log.old_value = (calculated_number) sqlite3_column_double(res, 24);
203 +
204 + alarm_log.updated = (sqlite3_column_int64(res, 8) & HEALTH_ENTRY_FLAG_UPDATED) ? 1 : 0;
205 + alarm_log.rendered_info = strdupz((char *)sqlite3_column_text(res, 18));
206 +
207 + info("DEBUG: %s pushing alert seq %" PRIu64 " - %" PRIu64"", wc->uuid_str, (uint64_t) sqlite3_column_int64(res, 0), (uint64_t) sqlite3_column_int64(res, 1));
208 + aclk_send_alarm_log_entry(&alarm_log);
209 +
210 + if (first_sequence_id == 0)
211 + first_sequence_id = (uint64_t) sqlite3_column_int64(res, 0);
212 + last_sequence_id = (uint64_t) sqlite3_column_int64(res, 0);
213 +
214 + destroy_alarm_log_entry(&alarm_log);
215 + freez(edit_command);
216 + }
217 + buffer_flush(sql);
218 +
219 + buffer_sprintf(sql, "UPDATE aclk_alert_%s SET date_submitted=strftime('%%s') "
220 + "WHERE date_submitted IS NULL AND sequence_id BETWEEN %" PRIu64 " AND %" PRIu64 ";",
221 + wc->uuid_str, first_sequence_id, last_sequence_id);
222 + db_execute(buffer_tostring(sql));
223 +
224 + rc = sqlite3_finalize(res);
225 + if (unlikely(rc != SQLITE_OK))
226 + error_report("Failed to finalize statement to send alert entries from the database, rc = %d", rc);
227 +
228 + freez(claim_id);
229 + buffer_free(sql);
230 +#endif
231 +
232 + return;
233 +}
234 +
235 +void aclk_send_alarm_health_log(char *node_id)
236 +{
237 + if (unlikely(!node_id))
238 + return;
239 +
240 + struct aclk_database_worker_config *wc = NULL;
241 + struct aclk_database_cmd cmd;
242 + memset(&cmd, 0, sizeof(cmd));
243 + cmd.opcode = ACLK_DATABASE_ALARM_HEALTH_LOG;
244 +
245 + rrd_wrlock();
246 + RRDHOST *host = find_host_by_node_id(node_id);
247 + if (likely(host))
248 + wc = (struct aclk_database_worker_config *)host->dbsync_worker;
249 + rrd_unlock();
250 + if (wc)
251 + aclk_database_enq_cmd(wc, &cmd);
252 + else {
253 + if (aclk_worker_enq_cmd(node_id, &cmd))
254 + error_report("ACLK synchronization thread is not active for node id %s", node_id);
255 + }
256 + return;
257 +}
258 +
259 +void aclk_push_alarm_health_log(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
260 +{
261 + UNUSED(cmd);
262 +#ifndef ACLK_NG
263 + UNUSED(wc);
264 +#else
265 + int rc;
266 +
267 + char *claim_id = is_agent_claimed();
268 + if (unlikely(!claim_id))
269 + return;
270 +
271 + uint64_t first_sequence = 0;
272 + uint64_t last_sequence = 0;
273 + struct timeval first_timestamp;
274 + struct timeval last_timestamp;
275 +
276 + BUFFER *sql = buffer_create(1024);
277 +
278 + sqlite3_stmt *res = NULL;
279 +
280 + //TODO: make this better: include info from health log too
281 + buffer_sprintf(sql, "select aa.sequence_id, aa.date_created, \
282 + (select laa.sequence_id from aclk_alert_%s laa \
283 + order by laa.sequence_id desc limit 1), \
284 + (select laa.date_created from aclk_alert_%s laa \
285 + order by laa.sequence_id desc limit 1) \
286 + from aclk_alert_%s aa order by aa.sequence_id asc limit 1;", wc->uuid_str, wc->uuid_str, wc->uuid_str);
287 +
288 + rc = sqlite3_prepare_v2(db_meta, buffer_tostring(sql), -1, &res, 0);
289 + if (rc != SQLITE_OK) {
290 + error_report("Failed to prepare statement to get health log statistics from the database");
291 + buffer_free(sql);
292 + freez(claim_id);
293 + return;
294 + }
295 +
296 + first_timestamp.tv_sec = 0;
297 + first_timestamp.tv_usec = 0;
298 + last_timestamp.tv_sec = 0;
299 + last_timestamp.tv_usec = 0;
300 +
301 + while (sqlite3_step(res) == SQLITE_ROW) {
302 + first_sequence = sqlite3_column_bytes(res, 0) > 0 ? (uint64_t) sqlite3_column_int64(res, 0) : 0;
303 + if (sqlite3_column_bytes(res, 1) > 0) {
304 + first_timestamp.tv_sec = sqlite3_column_int64(res, 1);
305 + }
306 +
307 + last_sequence = sqlite3_column_bytes(res, 2) > 0 ? (uint64_t) sqlite3_column_int64(res, 2) : 0;
308 + if (sqlite3_column_bytes(res, 3) > 0) {
309 + last_timestamp.tv_sec = sqlite3_column_int64(res, 3);
310 + }
311 + }
312 +
313 + struct alarm_log_entries log_entries;
314 + log_entries.first_seq_id = first_sequence;
315 + log_entries.first_when = first_timestamp;
316 + log_entries.last_seq_id = last_sequence;
317 + log_entries.last_when = last_timestamp;
318 +
319 + struct alarm_log_health alarm_log;
320 + alarm_log.claim_id = claim_id;
321 + alarm_log.node_id = strdupz(wc->node_id);
322 + alarm_log.log_entries = log_entries;
323 + alarm_log.status = wc->alert_updates == 0 ? 2 : 1;
324 +
325 + wc->alert_sequence_id = last_sequence;
326 +
327 + aclk_send_alarm_log_health(&alarm_log);
328 +
329 + rc = sqlite3_finalize(res);
330 + if (unlikely(rc != SQLITE_OK))
331 + error_report("Failed to reset statement to get health log statistics from the database, rc = %d", rc);
332 +
333 + freez((char *)alarm_log.node_id);
334 + freez(claim_id);
335 + buffer_free(sql);
336 +#endif
337 +
338 + return;
339 +}
340 +
341 +void aclk_send_alarm_configuration(char *config_hash)
342 +{
343 + if (unlikely(!config_hash))
344 + return;
345 +
346 + struct aclk_database_worker_config *wc = (struct aclk_database_worker_config *) localhost->dbsync_worker;
347 +
348 + if (unlikely(!wc)) {
349 + return;
350 + }
351 +
352 + struct aclk_database_cmd cmd;
353 + memset(&cmd, 0, sizeof(cmd));
354 + cmd.opcode = ACLK_DATABASE_PUSH_ALERT_CONFIG;
355 + cmd.data_param = (void *) strdupz(config_hash);
356 + cmd.completion = NULL;
357 + aclk_database_enq_cmd(wc, &cmd);
358 +
359 + return;
360 +}
361 +
362 +int aclk_push_alert_config_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
363 +{
364 + UNUSED(wc);
365 +#ifndef ACLK_NG
366 + UNUSED(cmd);
367 +#else
368 + int rc = 0;
369 +
370 + CHECK_SQLITE_CONNECTION(db_meta);
371 +
372 + sqlite3_stmt *res = NULL;
373 +
374 + char *config_hash = (char *) cmd.data_param;
375 + BUFFER *sql = buffer_create(1024);
376 + buffer_sprintf(
377 + sql,
378 + "SELECT alarm, template, on_key, class, type, component, os, hosts, plugin, module, charts, families, lookup, every, units, green, red, calc, warn, crit, to_key, exec, delay, repeat, info, options, host_labels, p_db_lookup_dimensions, p_db_lookup_method, p_db_lookup_options, p_db_lookup_after, p_db_lookup_before, p_update_every FROM alert_hash WHERE hash_id = @hash_id;");
379 +
380 + rc = sqlite3_prepare_v2(db_meta, buffer_tostring(sql), -1, &res, 0);
381 + if (rc != SQLITE_OK) {
382 + error_report("Failed to prepare statement when trying to fetch a chart hash configuration");
383 + goto fail;
384 + }
385 +
386 + uuid_t hash_uuid;
387 + if (uuid_parse(config_hash, hash_uuid))
388 + goto fail;
389 +
390 + rc = sqlite3_bind_blob(res, 1, &hash_uuid , sizeof(hash_uuid), SQLITE_STATIC);
391 + if (unlikely(rc != SQLITE_OK))
392 + goto bind_fail;
393 +
394 + struct aclk_alarm_configuration alarm_config;
395 + struct provide_alarm_configuration p_alarm_config;
396 + p_alarm_config.cfg_hash = NULL;
397 +
398 + if (sqlite3_step(res) == SQLITE_ROW) {
399 +
400 + alarm_config.alarm = sqlite3_column_bytes(res, 0) > 0 ? strdupz((char *)sqlite3_column_text(res, 0)) : NULL;
401 + alarm_config.tmpl = sqlite3_column_bytes(res, 1) > 0 ? strdupz((char *)sqlite3_column_text(res, 1)) : NULL;
402 + alarm_config.on_chart = sqlite3_column_bytes(res, 2) > 0 ? strdupz((char *)sqlite3_column_text(res, 2)) : NULL;
403 + alarm_config.classification = sqlite3_column_bytes(res, 3) > 0 ? strdupz((char *)sqlite3_column_text(res, 3)) : NULL;
404 + alarm_config.type = sqlite3_column_bytes(res, 4) > 0 ? strdupz((char *)sqlite3_column_text(res, 4)) : NULL;
405 + alarm_config.component = sqlite3_column_bytes(res, 5) > 0 ? strdupz((char *)sqlite3_column_text(res, 5)) : NULL;
406 +
407 + alarm_config.os = sqlite3_column_bytes(res, 6) > 0 ? strdupz((char *)sqlite3_column_text(res, 6)) : NULL;
408 + alarm_config.hosts = sqlite3_column_bytes(res, 7) > 0 ? strdupz((char *)sqlite3_column_text(res, 7)) : NULL;
409 + alarm_config.plugin = sqlite3_column_bytes(res, 8) > 0 ? strdupz((char *)sqlite3_column_text(res, 8)) : NULL;
410 + alarm_config.module = sqlite3_column_bytes(res, 9) > 0 ? strdupz((char *)sqlite3_column_text(res, 9)) : NULL;
411 + alarm_config.charts = sqlite3_column_bytes(res, 10) > 0 ? strdupz((char *)sqlite3_column_text(res, 10)) : NULL;
412 + alarm_config.families = sqlite3_column_bytes(res, 11) > 0 ? strdupz((char *)sqlite3_column_text(res, 11)) : NULL;
413 + alarm_config.lookup = sqlite3_column_bytes(res, 12) > 0 ? strdupz((char *)sqlite3_column_text(res, 12)) : NULL;
414 + alarm_config.every = sqlite3_column_bytes(res, 13) > 0 ? strdupz((char *)sqlite3_column_text(res, 13)) : NULL;
415 + alarm_config.units = sqlite3_column_bytes(res, 14) > 0 ? strdupz((char *)sqlite3_column_text(res, 14)) : NULL;
416 +
417 + alarm_config.green = sqlite3_column_bytes(res, 15) > 0 ? strdupz((char *)sqlite3_column_text(res, 15)) : NULL;
418 + alarm_config.red = sqlite3_column_bytes(res, 16) > 0 ? strdupz((char *)sqlite3_column_text(res, 16)) : NULL;
419 +
420 + alarm_config.calculation_expr = sqlite3_column_bytes(res, 17) > 0 ? strdupz((char *)sqlite3_column_text(res, 17)) : NULL;
421 + alarm_config.warning_expr = sqlite3_column_bytes(res, 18) > 0 ? strdupz((char *)sqlite3_column_text(res, 18)) : NULL;
422 + alarm_config.critical_expr = sqlite3_column_bytes(res, 19) > 0 ? strdupz((char *)sqlite3_column_text(res, 19)) : NULL;
423 +
424 + alarm_config.recipient = sqlite3_column_bytes(res, 20) > 0 ? strdupz((char *)sqlite3_column_text(res, 20)) : NULL;
425 + alarm_config.exec = sqlite3_column_bytes(res, 21) > 0 ? strdupz((char *)sqlite3_column_text(res, 21)) : NULL;
426 + alarm_config.delay = sqlite3_column_bytes(res, 22) > 0 ? strdupz((char *)sqlite3_column_text(res, 22)) : NULL;
427 + alarm_config.repeat = sqlite3_column_bytes(res, 23) > 0 ? strdupz((char *)sqlite3_column_text(res, 23)) : NULL;
428 + alarm_config.info = sqlite3_column_bytes(res, 24) > 0 ? strdupz((char *)sqlite3_column_text(res, 24)) : NULL;
429 + alarm_config.options = sqlite3_column_bytes(res, 25) > 0 ? strdupz((char *)sqlite3_column_text(res, 25)) : NULL;
430 + alarm_config.host_labels = sqlite3_column_bytes(res, 26) > 0 ? strdupz((char *)sqlite3_column_text(res, 26)) : NULL;
431 +
432 + alarm_config.p_db_lookup_dimensions = NULL;
433 + alarm_config.p_db_lookup_method = NULL;
434 + alarm_config.p_db_lookup_options = NULL;
435 + alarm_config.p_db_lookup_after = 0;
436 + alarm_config.p_db_lookup_before = 0;
437 +
438 + if (sqlite3_column_bytes(res, 30) > 0) {
439 +
440 + alarm_config.p_db_lookup_dimensions = sqlite3_column_bytes(res, 27) > 0 ? strdupz((char *)sqlite3_column_text(res, 27)) : NULL;
441 + alarm_config.p_db_lookup_method = sqlite3_column_bytes(res, 28) > 0 ? strdupz((char *)sqlite3_column_text(res, 28)) : NULL;
442 +
443 + BUFFER *tmp_buf = buffer_create(1024);
444 + buffer_data_options2string(tmp_buf, sqlite3_column_int(res, 29));
445 + alarm_config.p_db_lookup_options = strdupz((char *)buffer_tostring(tmp_buf));
446 + buffer_free(tmp_buf);
447 +
448 + alarm_config.p_db_lookup_after = sqlite3_column_int(res, 30);
449 + alarm_config.p_db_lookup_before = sqlite3_column_int(res, 31);
450 + }
451 +
452 + alarm_config.p_update_every = sqlite3_column_int(res, 32);
453 +
454 + p_alarm_config.cfg_hash = strdupz((char *) config_hash);
455 + p_alarm_config.cfg = alarm_config;
456 + }
457 +
458 + if (likely(p_alarm_config.cfg_hash)) {
459 + debug(D_ACLK_SYNC, "Sending alert config for %s", config_hash);
460 + aclk_send_provide_alarm_cfg(&p_alarm_config);
461 + freez((char *) cmd.data_param);
462 + freez(p_alarm_config.cfg_hash);
463 + destroy_aclk_alarm_configuration(&alarm_config);
464 + }
465 + else
466 + info("DEBUG: Alert config for %s not found", config_hash);
467 +
468 + bind_fail:
469 + rc = sqlite3_finalize(res);
470 + if (unlikely(rc != SQLITE_OK))
471 + error_report("Failed to reset statement when pushing alarm config hash, rc = %d", rc);
472 +
473 + fail:
474 + buffer_free(sql);
475 +
476 + return rc;
477 +#endif
478 + return 0;
479 +}
480 +
481 +
482 +// Start streaming alerts
483 +void aclk_start_alert_streaming(char *node_id, uint64_t batch_id, uint64_t start_seq_id)
484 +{
485 +#ifdef ACLK_NG
486 + if (unlikely(!node_id))
487 + return;
488 +
489 + uuid_t node_uuid;
490 + if (uuid_parse(node_id, node_uuid))
491 + return;
492 +
493 + struct aclk_database_worker_config *wc = NULL;
494 + rrd_wrlock();
495 + RRDHOST *host = find_host_by_node_id(node_id);
496 + if (likely(host))
497 + wc = (struct aclk_database_worker_config *)host->dbsync_worker;
498 + rrd_unlock();
499 +
500 + if (likely(wc)) {
501 + info("START streaming alerts for %s enabled with batch_id %"PRIu64" and start_seq_id %"PRIu64, node_id, batch_id, start_seq_id);
502 + __sync_synchronize();
503 + wc->alerts_batch_id = batch_id;
504 + wc->alerts_start_seq_id = start_seq_id;
505 + wc->alert_updates = 1;
506 + __sync_synchronize();
507 + }
508 + else
509 + error("ACLK synchronization thread is not active for host %s", host->hostname);
510 +
511 +#else
512 + UNUSED(node_id);
513 + UNUSED(start_seq_id);
514 + UNUSED(batch_id);
515 +#endif
516 + return;
517 +}
518 +
519 +int sql_queue_removed_alerts_to_aclk(RRDHOST *host)
520 +{
521 + CHECK_SQLITE_CONNECTION(db_meta);
522 +
523 + struct aclk_database_worker_config *wc = (struct aclk_database_worker_config *) host->dbsync_worker;
524 + if (unlikely(!wc)) {
525 + return 1;
526 + }
527 +
528 + BUFFER *sql = buffer_create(1024);
529 +
530 + buffer_sprintf(sql,"insert into aclk_alert_%s (alert_unique_id, date_created) " \
531 + "select unique_id alert_unique_id, strftime('%%s') date_created from health_log_%s where new_status = -2 and updated_by_id = 0 and unique_id not in (select alert_unique_id from aclk_alert_%s) order by unique_id asc on conflict (alert_unique_id) do nothing;", wc->uuid_str, wc->uuid_str, wc->uuid_str);
532 +
533 + db_execute(buffer_tostring(sql));
534 +
535 + buffer_free(sql);
536 +
537 + return 0;
538 +}
database/sqlite/sqlite_aclk_alert.h new
+17
@@ -0,0 +1,17 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#ifndef NETDATA_SQLITE_ACLK_ALERT_H
4 +#define NETDATA_SQLITE_ACLK_ALERT_H
5 +
6 +extern sqlite3 *db_meta;
7 +
8 +int aclk_add_alert_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
9 +void aclk_push_alert_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
10 +void aclk_send_alarm_health_log(char *node_id);
11 +void aclk_push_alarm_health_log(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
12 +void aclk_send_alarm_configuration (char *config_hash);
13 +int aclk_push_alert_config_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
14 +void aclk_start_alert_streaming(char *node_id, uint64_t batch_id, uint64_t start_seq_id);
15 +int sql_queue_removed_alerts_to_aclk(RRDHOST *host);
16 +
17 +#endif //NETDATA_SQLITE_ACLK_ALERT_H
database/sqlite/sqlite_aclk_chart.c
-9
@@ -5,15 +5,6 @@
5
6 #include "../../aclk/aclk_charts_api.h"
7
8 -#define CHECK_SQLITE_CONNECTION(db_meta) \
9 - if (unlikely(!db_meta)) { \
10 - if (default_rrd_memory_mode != RRD_MEMORY_MODE_DBENGINE) { \
11 - return 1; \
12 - } \
13 - error_report("Database has not been initialized"); \
14 - return 1; \
15 - }
16 -
8 static inline int sql_queue_chart_payload(struct aclk_database_worker_config *wc,
9 void *data, enum aclk_database_opcode opcode)
10 {
database/sqlite/sqlite_functions.h
+10
@@ -39,6 +39,16 @@ struct node_instance_list {
39
40 #define SQL_STORE_ACTIVE_DIMENSION \
41 "insert or replace into dimension_active (dim_id, date_created) values (@id, strftime('%s'));"
42 +
43 +#define CHECK_SQLITE_CONNECTION(db_meta) \
44 + if (unlikely(!db_meta)) { \
45 + if (default_rrd_memory_mode != RRD_MEMORY_MODE_DBENGINE) { \
46 + return 1; \
47 + } \
48 + error_report("Database has not been initialized"); \
49 + return 1; \
50 + }
51 +
52 extern int sql_init_database(void);
53 extern void sql_close_database(void);
54
health/health.h
+1
@@ -27,6 +27,7 @@ extern unsigned int default_health_enabled;
27 #define HEALTH_ENTRY_FLAG_EXEC_IN_PROGRESS 0x00000040
28
29 #define HEALTH_ENTRY_FLAG_SAVED 0x10000000
30 +#define HEALTH_ENTRY_FLAG_ACLK_QUEUED 0x20000000
31 #define HEALTH_ENTRY_FLAG_NO_CLEAR_NOTIFICATION 0x80000000
32
33 #ifndef HEALTH_LISTEN_PORT