Fix issue with charts not properly synchronized with the cloud (#12451)
* Add function to check a specific chart * If a chart is not obsoleted, check if the liveness needs to be updated * Calculate liveness based on a (constant * update_every) for each dimension * Scan all dimensions when the retention message is constructed and update liveness if needed * If initial state, set to computed live * Set computed live state to dimension * Add a maximum dimension cleanup on startup to prevent message flood * Schedule chart updates if charts streaming is enabled * Adjust live state for dimension * The query executed will have a valid dimension uuid only if memory mode is dbengine
Stelios Fragkakis committed
Apr 1, 2022 at 18:12 UTC
e816ee49237bfbdb289d7b0bba2bf2a8b94b0a5e
4 files changed
+99
-37
database/rrdhost.c
+4
@@ -1504,6 +1504,10 @@ restart_after_removal:
1504
rrdset_free(st);
1505
goto restart_after_removal;
1506
}
1507
+#if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
1508
+ else
1509
+ sql_check_chart_liveness(st);
1510
+#endif
1511
}
1512
}
1513
database/rrdset.c
+2
-2
@@ -1833,8 +1833,8 @@ after_second_database_work:
1833
#if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
1834
if (likely(!st->state->is_ar_chart)) {
1835
if (!rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN)) {
1836
- int live = ((mark - rd->last_collected_time.tv_sec) <
1837
- MAX(RRDSET_MINIMUM_LIVE_MULTIPLIER * rd->update_every, rrdset_free_obsolete_time));
1836
+ int live =
1837
+ ((mark - rd->last_collected_time.tv_sec) < RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every);
1838
if (unlikely(live != rd->state->aclk_live_status)) {
1839
if (likely(rrdset_flag_check(st, RRDSET_FLAG_ACLK))) {
1840
if (likely(!queue_dimension_to_aclk(rd))) {
database/sqlite/sqlite_aclk_chart.c
+86
-33
@@ -58,18 +58,15 @@ bind_fail:
58
return send_status;
59
}
60
61
-static int aclk_add_chart_payload(
62
- struct aclk_database_worker_config *wc,
63
- uuid_t *uuid,
64
- char *claim_id,
65
- ACLK_PAYLOAD_TYPE payload_type,
66
- void *payload,
67
- size_t payload_size)
61
+static int aclk_add_chart_payload(struct aclk_database_worker_config *wc, uuid_t *uuid, char *claim_id,
62
+ ACLK_PAYLOAD_TYPE payload_type, void *payload, size_t payload_size, int *send_status)
63
{
64
static __thread sqlite3_stmt *res_chart = NULL;
65
int rc;
66
67
rc = payload_sent(wc->uuid_str, uuid, payload, payload_size);
68
+ if (send_status)
69
+ *send_status = rc;
70
if (rc == 1)
71
return 0;
72
@@ -162,22 +159,16 @@ int aclk_add_chart_event(struct aclk_database_worker_config *wc, struct aclk_dat
159
size_t size;
160
char *payload = generate_chart_instance_updated(&size, &chart_payload);
161
if (likely(payload))
165
- rc = aclk_add_chart_payload(wc, st->chart_uuid, claim_id, ACLK_PAYLOAD_CHART, (void *)payload, size);
162
+ rc = aclk_add_chart_payload(wc, st->chart_uuid, claim_id, ACLK_PAYLOAD_CHART, (void *) payload, size, NULL);
163
freez(payload);
164
chart_instance_updated_destroy(&chart_payload);
165
}
166
return rc;
167
}
168
172
-static inline int aclk_upd_dimension_event(
173
- struct aclk_database_worker_config *wc,
174
- char *claim_id,
175
- uuid_t *dim_uuid,
176
- const char *dim_id,
177
- const char *dim_name,
178
- const char *chart_type_id,
179
- time_t first_time,
180
- time_t last_time)
169
+static inline int aclk_upd_dimension_event(struct aclk_database_worker_config *wc, char *claim_id, uuid_t *dim_uuid,
170
+ const char *dim_id, const char *dim_name, const char *chart_type_id, time_t first_time, time_t last_time,
171
+ int *send_status)
172
{
173
int rc = 0;
174
size_t size;
@@ -190,13 +181,11 @@ static inline int aclk_upd_dimension_event(
181
182
#ifdef NETDATA_INTERNAL_CHECKS
183
if (!first_time)
193
- info(
194
- "Host %s (node %s) deleting dimension id=[%s] name=[%s] chart=[%s]",
195
- wc->host_guid,
196
- wc->node_id,
197
- dim_id,
198
- dim_name,
199
- chart_type_id);
184
+ info("Host %s (node %s) deleting dimension id=[%s] name=[%s] chart=[%s]",
185
+ wc->host_guid, wc->node_id, dim_id, dim_name, chart_type_id);
186
+ if (last_time)
187
+ info("Host %s (node %s) stopped collecting dimension id=[%s] name=[%s] chart=[%s] %ld seconds ago at %ld",
188
+ wc->host_guid, wc->node_id, dim_id, dim_name, chart_type_id, now_realtime_sec() - last_time, last_time);
189
#endif
190
191
dim_payload.node_id = wc->node_id;
@@ -208,7 +197,7 @@ static inline int aclk_upd_dimension_event(
197
dim_payload.last_timestamp.tv_sec = last_time;
198
char *payload = generate_chart_dimension_updated(&size, &dim_payload);
199
if (likely(payload))
211
- rc = aclk_add_chart_payload(wc, dim_uuid, claim_id, ACLK_PAYLOAD_DIMENSION, (void *)payload, size);
200
+ rc = aclk_add_chart_payload(wc, dim_uuid, claim_id, ACLK_PAYLOAD_DIMENSION, (void *)payload, size, send_status);
201
freez(payload);
202
return rc;
203
}
@@ -252,7 +241,7 @@ void aclk_process_dimension_deletion(struct aclk_database_worker_config *wc, str
241
242
unsigned count = 0;
243
while (sqlite3_step(res) == SQLITE_ROW) {
255
- (void)aclk_upd_dimension_event(
244
+ (void) aclk_upd_dimension_event(
245
wc,
246
claim_id,
247
(uuid_t *)sqlite3_column_text(res, 3),
@@ -260,7 +249,8 @@ void aclk_process_dimension_deletion(struct aclk_database_worker_config *wc, str
249
(const char *)sqlite3_column_text(res, 1),
250
(const char *)sqlite3_column_text(res, 2),
251
0,
263
- 0);
252
+ 0,
253
+ NULL);
254
count++;
255
}
256
@@ -289,12 +279,13 @@ int aclk_add_dimension_event(struct aclk_database_worker_config *wc, struct aclk
279
RRDDIM *rd = cmd.data;
280
281
if (likely(claim_id)) {
282
+ int send_status = 0;
283
time_t now = now_realtime_sec();
284
285
time_t first_t = rd->state->query_ops.oldest_time(rd);
286
time_t last_t = rd->state->query_ops.latest_time(rd);
287
297
- int live = ((now - last_t) < MAX(RRDSET_MINIMUM_LIVE_MULTIPLIER * rd->update_every, rrdset_free_obsolete_time));
288
+ int live = ((now - last_t) < (RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every));
289
290
rc = aclk_upd_dimension_event(
291
wc,
@@ -304,7 +295,11 @@ int aclk_add_dimension_event(struct aclk_database_worker_config *wc, struct aclk
295
rd->name,
296
rd->rrdset->id,
297
first_t,
307
- live ? 0 : last_t);
298
+ live ? 0 : last_t,
299
+ &send_status);
300
+
301
+ if (!send_status)
302
+ rd->state->aclk_live_status = live;
303
304
freez(claim_id);
305
}
@@ -891,6 +886,8 @@ void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_d
886
time_t first_entry_t;
887
time_t last_entry_t;
888
uint32_t update_every = 0;
889
+ uint32_t dimension_update_count = 0;
890
+ int send_status;
891
892
struct retention_updated rotate_data;
893
@@ -906,7 +903,7 @@ void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_d
903
rotate_data.claim_id = claim_id;
904
rotate_data.node_id = strdupz(wc->node_id);
905
909
- // time_t now = now_realtime_sec();
906
+ time_t now = now_realtime_sec();
907
while (sqlite3_step(res) == SQLITE_ROW) {
908
if (!update_every || update_every != (uint32_t)sqlite3_column_int(res, 1)) {
909
if (update_every) {
@@ -944,6 +941,24 @@ void aclk_update_retention(struct aclk_database_worker_config *wc, struct aclk_d
941
942
if (likely(!rc && first_entry_t))
943
start_time = MIN(start_time, first_entry_t);
944
+
945
+ if (memory_mode == RRD_MEMORY_MODE_DBENGINE && wc->chart_updates) {
946
+ int live = ((now - last_entry_t) < (RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * update_every));
947
+ if ((!live || !first_entry_t) && (dimension_update_count < ACLK_MAX_DIMENSION_CLEANUP)) {
948
+ (void)aclk_upd_dimension_event(
949
+ wc,
950
+ claim_id,
951
+ (uuid_t *)sqlite3_column_blob(res, 0),
952
+ (const char *)(const char *)sqlite3_column_text(res, 3),
953
+ (const char *)(const char *)sqlite3_column_text(res, 4),
954
+ (const char *)(const char *)sqlite3_column_text(res, 2),
955
+ first_entry_t,
956
+ live ? 0 : last_entry_t,
957
+ &send_status);
958
+ if (!send_status)
959
+ dimension_update_count++;
960
+ }
961
+ }
962
}
963
if (update_every) {
964
debug(D_ACLK_SYNC, "Update %s for %u oldest time = %ld", wc->host_guid, update_every, start_time);
@@ -1053,8 +1068,7 @@ void aclk_send_dimension_update(RRDDIM *rd)
1068
time_t last_entry_t = rrddim_last_entry_t(rd);
1069
1070
time_t now = now_realtime_sec();
1056
- int live = ((now - rd->last_collected_time.tv_sec) <
1057
- MAX(RRDSET_MINIMUM_LIVE_MULTIPLIER * rd->update_every, rrdset_free_obsolete_time));
1071
+ int live = ((now - rd->last_collected_time.tv_sec) < (RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every));
1072
1073
if (!live || rd->state->aclk_live_status != live || !first_entry_t) {
1074
(void)aclk_upd_dimension_event(
@@ -1065,7 +1079,8 @@ void aclk_send_dimension_update(RRDDIM *rd)
1079
rd->name,
1080
rd->rrdset->id,
1081
first_entry_t,
1068
- live ? 0 : last_entry_t);
1082
+ live ? 0 : last_entry_t,
1083
+ NULL);
1084
1085
if (!first_entry_t)
1086
debug(
@@ -1180,6 +1195,44 @@ struct aclk_chart_sync_stats *aclk_get_chart_sync_stats(RRDHOST *host)
1195
buffer_free(sql);
1196
return aclk_statistics;
1197
}
1198
+
1199
+void sql_check_chart_liveness(RRDSET *st) {
1200
+ RRDDIM *rd;
1201
+
1202
+ if (unlikely(st->state->is_ar_chart))
1203
+ return;
1204
+
1205
+ rrdset_rdlock(st);
1206
+ if (unlikely(!rrdset_flag_check(st, RRDSET_FLAG_ACLK))) {
1207
+ if (likely(st->dimensions && st->counter_done && !queue_chart_to_aclk(st))) {
1208
+ debug(D_ACLK_SYNC,"Check chart liveness [%s] submit chart definition", st->name);
1209
+ rrdset_flag_set(st, RRDSET_FLAG_ACLK);
1210
+ }
1211
+ }
1212
+ else
1213
+ debug(D_ACLK_SYNC,"Check chart liveness [%s] chart definition already submitted", st->name);
1214
+ time_t mark = now_realtime_sec();
1215
+
1216
+ debug(D_ACLK_SYNC,"Check chart liveness [%s] scanning dimensions", st->name);
1217
+ rrddim_foreach_read(rd, st) {
1218
+ if (!rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN)) {
1219
+ int live = (mark - rd->last_collected_time.tv_sec) < RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every;
1220
+ if (unlikely(live != rd->state->aclk_live_status)) {
1221
+ if (likely(rrdset_flag_check(st, RRDSET_FLAG_ACLK))) {
1222
+ if (likely(!queue_dimension_to_aclk(rd))) {
1223
+ debug(D_ACLK_SYNC,"Dimension change [%s] on [%s] from live %d --> %d", rd->id, rd->rrdset->name, rd->state->aclk_live_status, live);
1224
+ rd->state->aclk_live_status = live;
1225
+ rrddim_flag_set(rd, RRDDIM_FLAG_ACLK);
1226
+ }
1227
+ }
1228
+ }
1229
+ else
1230
+ debug(D_ACLK_SYNC,"Dimension check [%s] on [%s] liveness matches", rd->id, st->name);
1231
+ }
1232
+ }
1233
+ rrdset_unlock(st);
1234
+}
1235
+
1236
#endif //ENABLE_NEW_CLOUD_PROTOCOL
1237
1238
// ST is read locked
database/sqlite/sqlite_aclk_chart.h
+7
-2
@@ -12,8 +12,12 @@ typedef enum payload_type {
12
13
extern sqlite3 *db_meta;
14
15
-#ifndef RRDSET_MINIMUM_LIVE_MULTIPLIER
16
-#define RRDSET_MINIMUM_LIVE_MULTIPLIER (1.5)
15
+#ifndef RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER
16
+#define RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER (3)
17
+#endif
18
+
19
+#ifndef ACLK_MAX_DIMENSION_CLEANUP
20
+#define ACLK_MAX_DIMENSION_CLEANUP (500)
21
#endif
22
23
struct aclk_chart_sync_stats {
@@ -52,4 +56,5 @@ void aclk_process_dimension_deletion(struct aclk_database_worker_config *wc, str
56
uint32_t sql_get_pending_count(struct aclk_database_worker_config *wc);
57
void aclk_send_dimension_update(RRDDIM *rd);
58
struct aclk_chart_sync_stats *aclk_get_chart_sync_stats(RRDHOST *host);
59
+void sql_check_chart_liveness(RRDSET *st);
60
#endif //NETDATA_SQLITE_ACLK_CHART_H