@cryptotaxi247 / netdata-1 / commits / 3b8d4c21e

Adjust the dimension liveness status check (#12933)

* Mark a chart to be exposed only if dimension is created or metadata changes * Add a calculate liveness for the dimension for collected to non collected (live -> stale) and vice versa * queue_dimension_to_aclk will have the rrdset and either 0 or last collected time If 0 then it will be marked as live else it will be marked as stale and last collected time will be sent to the cloud * Add an extra parameter to indicate if the payload check should be done in the database or it has been done already * Queue dimension sets dimension liveness and queues the exact payload to store in the database * Fix compilation error when --disable-cloud is specified

Stelios Fragkakis committed May 17, 2022 at 16:58 UTC 3b8d4c21e5dd7abbc0b8b5c5b5b0bc826c229abc
6 files changed +132 -65
database/rrd.h
+3 -1
@@ -1274,7 +1274,9 @@ extern void rrddim_isnot_obsolete(RRDSET *st, RRDDIM *rd);
1274
1275 extern collected_number rrddim_set_by_pointer(RRDSET *st, RRDDIM *rd, collected_number value);
1276 extern collected_number rrddim_set(RRDSET *st, const char *id, collected_number value);
1277 -
1277 +#if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
1278 +extern time_t calc_dimension_liveness(RRDDIM *rd, time_t now);
1279 +#endif
1280 extern long align_entries_to_pagesize(RRD_MEMORY_MODE mode, long entries);
1281
1282 // ----------------------------------------------------------------------------
database/rrddim.c
+27 -3
@@ -135,15 +135,31 @@ void rrdcalc_link_to_rrddim(RRDDIM *rd, RRDSET *st, RRDHOST *host) {
135 }
136 }
137
138 +// Return either
139 +// 0 : Dimension is live
140 +// last collected time : Dimension is not live
141 +
142 +#if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
143 +time_t calc_dimension_liveness(RRDDIM *rd, time_t now)
144 +{
145 + time_t last_updated = rd->last_collected_time.tv_sec;
146 + int live;
147 + if (rd->state->aclk_live_status == 1)
148 + live =
149 + ((now - last_updated) <
150 + MIN(rrdset_free_obsolete_time, RRDSET_MINIMUM_DIM_OFFLINE_MULTIPLIER * rd->update_every));
151 + else
152 + live = ((now - last_updated) < RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every);
153 + return live ? 0 : last_updated;
154 +}
155 +#endif
156 +
157 RRDDIM *rrddim_add_custom(RRDSET *st, const char *id, const char *name, collected_number multiplier,
158 collected_number divisor, RRD_ALGORITHM algorithm, RRD_MEMORY_MODE memory_mode)
159 {
160 RRDHOST *host = st->rrdhost;
161 rrdset_wrlock(st);
162
144 - rrdset_flag_set(st, RRDSET_FLAG_SYNC_CLOCK);
145 - rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
146 -
163 RRDDIM *rd = rrddim_find(st, id);
164 if(unlikely(rd)) {
165 debug(D_RRD_CALLS, "Cannot create rrd dimension '%s/%s', it already exists.", st->id, name?name:"<NONAME>");
@@ -168,11 +184,19 @@ RRDDIM *rrddim_add_custom(RRDSET *st, const char *id, const char *name, collecte
184 debug(D_METADATALOG, "DIMENSION [%s] metadata updated", rd->id);
185 (void)sql_store_dimension(&rd->state->metric_uuid, rd->rrdset->chart_uuid, rd->id, rd->name, rd->multiplier, rd->divisor,
186 rd->algorithm);
187 +#if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
188 + queue_dimension_to_aclk(rd, calc_dimension_liveness(rd, now_realtime_sec()));
189 +#endif
190 + rrdset_flag_set(st, RRDSET_FLAG_SYNC_CLOCK);
191 + rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
192 }
193 rrdset_unlock(st);
194 return rd;
195 }
196
197 + rrdset_flag_set(st, RRDSET_FLAG_SYNC_CLOCK);
198 + rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
199 +
200 char filename[FILENAME_MAX + 1];
201 char fullfilename[FILENAME_MAX + 1];
202
database/rrdhost.c
+1 -1
@@ -1489,7 +1489,7 @@ restart_after_removal:
1489 }
1490 #if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
1491 else
1492 - queue_dimension_to_aclk(rd);
1492 + queue_dimension_to_aclk(rd, rd->last_collected_time.tv_sec);
1493 #endif
1494 }
1495 last = rd;
database/rrdset.c
+3 -7
@@ -1798,12 +1798,8 @@ after_second_database_work:
1798
1799 #if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
1800 if (likely(!st->state->is_ar_chart)) {
1801 - if (!rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN) && likely(rrdset_flag_check(st, RRDSET_FLAG_ACLK))) {
1802 - int live =
1803 - ((mark - rd->last_collected_time.tv_sec) < RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every);
1804 - if (unlikely(live != rd->state->aclk_live_status))
1805 - queue_dimension_to_aclk(rd);
1806 - }
1801 + if (!rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN) && likely(rrdset_flag_check(st, RRDSET_FLAG_ACLK)))
1802 + queue_dimension_to_aclk(rd, calc_dimension_liveness(rd, mark));
1803 }
1804 #endif
1805 if(unlikely(!rd->updated))
@@ -1906,7 +1902,7 @@ after_second_database_work:
1902 } else {
1903 /* Do not delete this dimension */
1904 #if defined(ENABLE_ACLK) && defined(ENABLE_NEW_CLOUD_PROTOCOL)
1909 - queue_dimension_to_aclk(rd);
1905 + queue_dimension_to_aclk(rd, calc_dimension_liveness(rd, mark));
1906 #endif
1907 last = rd;
1908 rd = rd->next;
database/sqlite/sqlite_aclk_chart.c
+87 -52
@@ -58,19 +58,31 @@ bind_fail:
58 return send_status;
59 }
60
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, time_t *send_status)
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,
68 + time_t *send_status,
69 + int check_sent)
70 {
71 static __thread sqlite3_stmt *res_chart = NULL;
72 int rc;
73 time_t date_submitted;
74
68 - date_submitted = payload_sent(wc->uuid_str, uuid, payload, payload_size);
69 - if (send_status)
70 - *send_status = date_submitted;
71 - if (date_submitted)
75 + if (unlikely(!payload))
76 return 0;
77
78 + if (check_sent) {
79 + date_submitted = payload_sent(wc->uuid_str, uuid, payload, payload_size);
80 + if (send_status)
81 + *send_status = date_submitted;
82 + if (date_submitted)
83 + return 0;
84 + }
85 +
86 if (unlikely(!res_chart)) {
87 char sql[ACLK_SYNC_QUERY_SIZE];
88 snprintfz(sql,ACLK_SYNC_QUERY_SIZE-1,
@@ -160,7 +172,7 @@ int aclk_add_chart_event(struct aclk_database_worker_config *wc, struct aclk_dat
172 size_t size;
173 char *payload = generate_chart_instance_updated(&size, &chart_payload);
174 if (likely(payload))
163 - rc = aclk_add_chart_payload(wc, st->chart_uuid, claim_id, ACLK_PAYLOAD_CHART, (void *) payload, size, NULL);
175 + rc = aclk_add_chart_payload(wc, st->chart_uuid, claim_id, ACLK_PAYLOAD_CHART, (void *) payload, size, NULL, 1);
176 freez(payload);
177 chart_instance_updated_destroy(&chart_payload);
178 }
@@ -198,7 +210,7 @@ static inline int aclk_upd_dimension_event(struct aclk_database_worker_config *w
210 dim_payload.last_timestamp.tv_sec = last_time;
211 char *payload = generate_chart_dimension_updated(&size, &dim_payload);
212 if (likely(payload))
201 - rc = aclk_add_chart_payload(wc, dim_uuid, claim_id, ACLK_PAYLOAD_DIMENSION, (void *)payload, size, send_status);
213 + rc = aclk_add_chart_payload(wc, dim_uuid, claim_id, ACLK_PAYLOAD_DIMENSION, (void *)payload, size, send_status, 1);
214 freez(payload);
215 return rc;
216 }
@@ -272,39 +284,22 @@ bind_fail:
284
285 int aclk_add_dimension_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
286 {
275 - int rc = 0;
287 + int rc = 1;
288 CHECK_SQLITE_CONNECTION(db_meta);
289
278 - char *claim_id = is_agent_claimed();
279 -
280 - RRDDIM *rd = cmd.data;
281 -
282 - if (likely(claim_id)) {
283 - time_t send_status = 0;
284 - time_t now = now_realtime_sec();
285 -
286 - time_t first_t = rd->state->query_ops.oldest_time(rd);
287 - time_t last_t = rd->state->query_ops.latest_time(rd);
288 -
289 - int live = ((now - last_t) < (RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every));
290 + struct aclk_chart_dimension_data *aclk_cd_data = cmd.data;
291
291 - rc = aclk_upd_dimension_event(
292 - wc,
293 - claim_id,
294 - &rd->state->metric_uuid,
295 - rd->id,
296 - rd->name,
297 - rd->rrdset->id,
298 - first_t,
299 - live ? 0 : last_t,
300 - &send_status);
292 + char *claim_id = is_agent_claimed();
293 + if (!claim_id)
294 + goto cleanup;
295
302 - if (!send_status)
303 - rd->state->aclk_live_status = live;
296 + rc = aclk_add_chart_payload(wc, &aclk_cd_data->uuid, claim_id, ACLK_PAYLOAD_DIMENSION,
297 + (void *) aclk_cd_data->payload, aclk_cd_data->payload_size, NULL, 0);
298
305 - freez(claim_id);
306 - }
307 - rrddim_flag_clear(rd, RRDDIM_FLAG_ACLK);
299 + freez(claim_id);
300 +cleanup:
301 + freez(aclk_cd_data->payload);
302 + freez(aclk_cd_data);
303 return rc;
304 }
305
@@ -1091,16 +1086,63 @@ void sql_get_last_chart_sequence(struct aclk_database_worker_config *wc)
1086 return;
1087 }
1088
1094 -void queue_dimension_to_aclk(RRDDIM *rd)
1089 +void queue_dimension_to_aclk(RRDDIM *rd, time_t last_updated)
1090 {
1096 - if (rrddim_flag_check(rd, RRDDIM_FLAG_ACLK))
1091 + int live = !last_updated;
1092 +
1093 + if (likely(rd->state->aclk_live_status == live))
1094 + return;
1095 +
1096 + rd->state->aclk_live_status = live;
1097 +
1098 + struct aclk_database_worker_config *wc = rd->rrdset->rrdhost->dbsync_worker;
1099 + if (unlikely(!wc))
1100 return;
1101
1099 - rrddim_flag_set(rd, RRDDIM_FLAG_ACLK);
1100 - int rc = sql_queue_chart_payload((struct aclk_database_worker_config *) rd->rrdset->rrdhost->dbsync_worker,
1101 - rd, ACLK_DATABASE_ADD_DIMENSION);
1102 - if (unlikely(rc))
1103 - rrddim_flag_clear(rd, RRDDIM_FLAG_ACLK);
1102 + char *claim_id = is_agent_claimed();
1103 + if (unlikely(!claim_id))
1104 + return;
1105 +
1106 + struct chart_dimension_updated dim_payload;
1107 + memset(&dim_payload, 0, sizeof(dim_payload));
1108 + dim_payload.node_id = wc->node_id;
1109 + dim_payload.claim_id = claim_id;
1110 + dim_payload.name = rd->name;
1111 + dim_payload.id = rd->id;
1112 + dim_payload.chart_id = rd->rrdset->id;
1113 + dim_payload.created_at.tv_sec = rd->state->query_ops.oldest_time(rd);
1114 + dim_payload.last_timestamp.tv_sec = last_updated;
1115 +
1116 + size_t size = 0;
1117 + char *payload = generate_chart_dimension_updated(&size, &dim_payload);
1118 +
1119 + freez(claim_id);
1120 + if (unlikely(!payload))
1121 + return;
1122 +
1123 + time_t date_submitted = payload_sent(wc->uuid_str, &rd->state->metric_uuid, payload, size);
1124 + if (date_submitted) {
1125 + freez(payload);
1126 + return;
1127 + }
1128 +
1129 + struct aclk_chart_dimension_data *aclk_cd_data = mallocz(sizeof(*aclk_cd_data));
1130 + uuid_copy(aclk_cd_data->uuid, rd->state->metric_uuid);
1131 + aclk_cd_data->payload = payload;
1132 + aclk_cd_data->payload_size = size;
1133 +
1134 + struct aclk_database_cmd cmd;
1135 + memset(&cmd, 0, sizeof(cmd));
1136 +
1137 + cmd.opcode = ACLK_DATABASE_ADD_DIMENSION;
1138 + cmd.data = aclk_cd_data;
1139 + int rc = aclk_database_enq_cmd_noblock(wc, &cmd);
1140 +
1141 + if (unlikely(rc)) {
1142 + freez(aclk_cd_data->payload);
1143 + freez(aclk_cd_data);
1144 + rd->state->aclk_live_status = !live;
1145 + }
1146 return;
1147 }
1148
@@ -1270,15 +1312,8 @@ void sql_check_chart_liveness(RRDSET *st) {
1312
1313 debug(D_ACLK_SYNC,"Check chart liveness [%s] scanning dimensions", st->name);
1314 rrddim_foreach_read(rd, st) {
1273 - if (!rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN)) {
1274 - int live = (mark - rd->last_collected_time.tv_sec) < RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER * rd->update_every;
1275 - if (unlikely(live != rd->state->aclk_live_status)) {
1276 - debug(D_ACLK_SYNC,"Dimension change [%s] on [%s] from live %d --> %d", rd->id, rd->rrdset->name, rd->state->aclk_live_status, live);
1277 - queue_dimension_to_aclk(rd);
1278 - }
1279 - else
1280 - debug(D_ACLK_SYNC,"Dimension check [%s] on [%s] liveness matches", rd->id, st->name);
1281 - }
1315 + if (!rrddim_flag_check(rd, RRDDIM_FLAG_HIDDEN))
1316 + queue_dimension_to_aclk(rd, calc_dimension_liveness(rd, mark));
1317 }
1318 rrdset_unlock(st);
1319 }
database/sqlite/sqlite_aclk_chart.h
+11 -1
@@ -16,10 +16,20 @@ extern sqlite3 *db_meta;
16 #define RRDSET_MINIMUM_DIM_LIVE_MULTIPLIER (3)
17 #endif
18
19 +#ifndef RRDSET_MINIMUM_DIM_OFFLINE_MULTIPLIER
20 +#define RRDSET_MINIMUM_DIM_OFFLINE_MULTIPLIER (30)
21 +#endif
22 +
23 #ifndef ACLK_MAX_DIMENSION_CLEANUP
24 #define ACLK_MAX_DIMENSION_CLEANUP (500)
25 #endif
26
27 +struct aclk_chart_dimension_data {
28 + uuid_t uuid;
29 + char *payload;
30 + size_t payload_size;
31 +};
32 +
33 struct aclk_chart_sync_stats {
34 int updates;
35 uint64_t batch_id;
@@ -37,7 +47,7 @@ struct aclk_chart_sync_stats {
47 };
48
49 extern int queue_chart_to_aclk(RRDSET *st);
40 -extern void queue_dimension_to_aclk(RRDDIM *rd);
50 +extern void queue_dimension_to_aclk(RRDDIM *rd, time_t last_updated);
51 extern void sql_create_aclk_table(RRDHOST *host, uuid_t *host_uuid, uuid_t *node_id);
52 int aclk_add_chart_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
53 int aclk_add_dimension_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);