@cryptotaxi247 / netdata-1 / commits / a19795e85

do not resend charts upstream when chart variables are being updated (#13946)

* do not resend charts upstream when chart variables are being updated * re-stream archived hosts that are now being collected

Costa Tsaousis committed Nov 3, 2022 at 12:13 UTC a19795e85fd1d026171661c7f97bde8f9f7d0b1a
7 files changed +41 -40
collectors/plugins.d/pluginsd_parser.c
+17 -4
@@ -287,6 +287,12 @@ PARSER_RC pluginsd_chart_definition_end(char **words, size_t num_words, void *us
287 "REPLAY: received " PLUGINSD_KEYWORD_CHART_DEFINITION_END " with malformed timings (first time %llu, last time %llu).",
288 (unsigned long long)first_entry_child, (unsigned long long)last_entry_child);
289
290 +// internal_error(
291 +// true,
292 +// "REPLAY host '%s', chart '%s': received " PLUGINSD_KEYWORD_CHART_DEFINITION_END " first time %llu, last time %llu.",
293 +// rrdhost_hostname(host), rrdset_id(st),
294 +// (unsigned long long)first_entry_child, (unsigned long long)last_entry_child);
295 +
296 rrdset_flag_clear(st, RRDSET_FLAG_RECEIVER_REPLICATION_FINISHED);
297
298 bool ok = replicate_chart_request(send_to_plugin, user_object->parser, host, st, first_entry_child, last_entry_child, 0, 0);
@@ -1098,12 +1104,12 @@ PARSER_RC pluginsd_replay_end(char **words, size_t num_words, void *user)
1104 }
1105
1106 time_t update_every_child = str2l(get_word(words, num_words, 1));
1101 - time_t first_entry_child = str2l(get_word(words, num_words, 2));
1102 - time_t last_entry_child = str2l(get_word(words, num_words, 3));
1107 + time_t first_entry_child = (time_t)str2ull(get_word(words, num_words, 2));
1108 + time_t last_entry_child = (time_t)str2ull(get_word(words, num_words, 3));
1109
1110 bool start_streaming = (strcmp(get_word(words, num_words, 4), "true") == 0);
1105 - time_t first_entry_requested = str2l(get_word(words, num_words, 5));
1106 - time_t last_entry_requested = str2l(get_word(words, num_words, 6));
1111 + time_t first_entry_requested = (time_t)str2ull(get_word(words, num_words, 5));
1112 + time_t last_entry_requested = (time_t)str2ull(get_word(words, num_words, 6));
1113
1114 PARSER_USER_OBJECT *user_object = user;
1115
@@ -1116,6 +1122,13 @@ PARSER_RC pluginsd_replay_end(char **words, size_t num_words, void *user)
1122 return PARSER_RC_ERROR;
1123 }
1124
1125 +// internal_error(true,
1126 +// "REPLAY: host '%s', chart '%s': received " PLUGINSD_KEYWORD_REPLAY_END " child first_t = %llu, last_t = %llu, start_streaming = %s, requested first_t = %llu, last_t = %llu",
1127 +// rrdhost_hostname(host), rrdset_id(st),
1128 +// (unsigned long long)first_entry_child, (unsigned long long)last_entry_child,
1129 +// start_streaming?"true":"false",
1130 +// (unsigned long long)first_entry_requested, (unsigned long long)last_entry_requested);
1131 +
1132 ((PARSER_USER_OBJECT *) user)->st = NULL;
1133 ((PARSER_USER_OBJECT *) user)->count++;
1134
database/rrd.h
+4
@@ -510,9 +510,11 @@ typedef enum rrdset_flags {
510 RRDSET_FLAG_OBSOLETE = (1 << 3), // this is marked by the collector/module as obsolete
511 RRDSET_FLAG_EXPORTING_SEND = (1 << 4), // if set, this chart should be sent to Prometheus web API and external databases
512 RRDSET_FLAG_EXPORTING_IGNORE = (1 << 5), // if set, this chart should not be sent to Prometheus web API and external databases
513 +
514 RRDSET_FLAG_UPSTREAM_SEND = (1 << 6), // if set, this chart should be sent upstream (streaming)
515 RRDSET_FLAG_UPSTREAM_IGNORE = (1 << 7), // if set, this chart should not be sent upstream (streaming)
516 RRDSET_FLAG_UPSTREAM_EXPOSED = (1 << 8), // if set, we have sent this chart definition to netdata parent (streaming)
517 +
518 RRDSET_FLAG_STORE_FIRST = (1 << 9), // if set, do not eliminate the first collection during interpolation
519 RRDSET_FLAG_HETEROGENEOUS = (1 << 10), // if set, the chart is not homogeneous (dimensions in it have multiple algorithms, multipliers or dividers)
520 RRDSET_FLAG_HOMOGENEOUS_CHECK = (1 << 11), // if set, the chart should be checked to determine if the dimensions are homogeneous
@@ -532,6 +534,8 @@ typedef enum rrdset_flags {
534
535 RRDSET_FLAG_SENDER_REPLICATION_FINISHED = (1 << 23), // the sending side has completed replication
536 RRDSET_FLAG_RECEIVER_REPLICATION_FINISHED = (1 << 24), // the receiving side has completed replication
537 +
538 + RRDSET_FLAG_UPSTREAM_SEND_VARIABLES = (1 << 25), // a custom variable has been updated and needs to be exposed to parent
539 } RRDSET_FLAGS;
540
541 #define rrdset_flag_check(st, flag) (__atomic_load_n(&((st)->flags), __ATOMIC_SEQ_CST) & (flag))
database/rrdhost.c
+1
@@ -1044,6 +1044,7 @@ void stop_streaming_sender(RRDHOST *host)
1044 dictionary_destroy(host->sender->replication_requests);
1045 freez(host->sender);
1046 host->sender = NULL;
1047 + rrdhost_flag_clear(host, RRDHOST_FLAG_RRDPUSH_SENDER_INITIALIZED);
1048 }
1049
1050 void stop_streaming_receiver(RRDHOST *host)
database/rrdsetvar.c
+3 -3
@@ -268,14 +268,14 @@ void rrdsetvar_custom_chart_variable_set(RRDSET *st, const RRDSETVAR_ACQUIRED *r
268 NETDATA_DOUBLE *v = rs->value;
269 if(*v != value) {
270 *v = value;
271 -
272 - // mark the chart to be sent upstream
273 - rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_EXPOSED);
271 + rrdset_flag_set(st, RRDSET_FLAG_UPSTREAM_SEND_VARIABLES);
272 }
273 }
274 }
275
276 void rrdsetvar_print_to_streaming_custom_chart_variables(RRDSET *st, BUFFER *wb) {
277 + rrdset_flag_clear(st, RRDSET_FLAG_UPSTREAM_SEND_VARIABLES);
278 +
279 // send the chart local custom variables
280 RRDSETVAR *rs;
281 dfe_start_read(st->rrdsetvar_root_index, rs) {
streaming/rrdpush.c
+14 -18
@@ -161,16 +161,7 @@ int rrdpush_init() {
161 // this is for the first iterations of each chart
162 unsigned int remote_clock_resync_iterations = 60;
163
164 -static inline bool should_send_chart_matching(RRDSET *st) {
165 - // get all the flags we need to check, with one atomic operation
166 - RRDSET_FLAGS flags = rrdset_flag_check(st,
167 - RRDSET_FLAG_UPSTREAM_SEND
168 - | RRDSET_FLAG_UPSTREAM_IGNORE
169 - | RRDSET_FLAG_ANOMALY_RATE_CHART
170 - | RRDSET_FLAG_ANOMALY_DETECTION
171 - | RRDSET_FLAG_RECEIVER_REPLICATION_FINISHED
172 - );
173 -
164 +static inline bool should_send_chart_matching(RRDSET *st, RRDSET_FLAGS flags) {
165 if(!(flags & RRDSET_FLAG_RECEIVER_REPLICATION_FINISHED))
166 return false;
167
@@ -220,8 +211,6 @@ int configured_as_parent() {
211 return is_parent;
212 }
213
223 -#define need_to_send_chart_definition(st) (!rrdset_flag_check(st, RRDSET_FLAG_UPSTREAM_EXPOSED))
224 -
214 // chart labels
215 static int send_clabels_callback(const char *name, const char *value, RRDLABEL_SRC ls, void *data) {
216 BUFFER *wb = (BUFFER *)data;
@@ -334,7 +323,7 @@ static inline void rrdpush_send_chart_definition(BUFFER *wb, RRDSET *st) {
323 }
324
325 // sends the current chart dimensions
337 -void rrdpush_send_chart_metrics(BUFFER *wb, RRDSET *st, struct sender_state *s) {
326 +static void rrdpush_send_chart_metrics(BUFFER *wb, RRDSET *st, struct sender_state *s, RRDSET_FLAGS flags) {
327 buffer_fast_strcat(wb, "BEGIN \"", 7);
328 buffer_fast_strcat(wb, rrdset_id(st), string_strlen(st->id));
329 buffer_fast_strcat(wb, "\" ", 2);
@@ -365,6 +354,10 @@ void rrdpush_send_chart_metrics(BUFFER *wb, RRDSET *st, struct sender_state *s)
354 }
355 }
356 rrddim_foreach_done(rd);
357 +
358 + if(unlikely(flags & RRDSET_FLAG_UPSTREAM_SEND_VARIABLES))
359 + rrdsetvar_print_to_streaming_custom_chart_variables(st, wb);
360 +
361 buffer_fast_strcat(wb, "END\n", 4);
362 }
363
@@ -374,7 +367,8 @@ static void rrdpush_sender_thread_spawn(RRDHOST *host);
367 bool rrdset_push_chart_definition_now(RRDSET *st) {
368 RRDHOST *host = st->rrdhost;
369
377 - if(unlikely(!rrdhost_can_send_definitions_to_parent(host) || !should_send_chart_matching(st)))
370 + if(unlikely(!rrdhost_can_send_definitions_to_parent(host)
371 + || !should_send_chart_matching(st, __atomic_load_n(&st->flags, __ATOMIC_SEQ_CST))))
372 return false;
373
374 BUFFER *wb = sender_start(host->sender);
@@ -412,16 +406,18 @@ void rrdset_done_push(RRDSET *st) {
406 rrdhost_flag_clear(host, RRDHOST_FLAG_RRDPUSH_SENDER_LOGGED_STATUS);
407 }
408
415 - if(unlikely(!should_send_chart_matching(st)))
409 + RRDSET_FLAGS rrdset_flags = __atomic_load_n(&st->flags, __ATOMIC_SEQ_CST);
410 +
411 + if(unlikely(!should_send_chart_matching(st, rrdset_flags)))
412 return;
413
414 BUFFER *wb = sender_start(host->sender);
415
420 - if(unlikely(need_to_send_chart_definition(st)))
416 + if(unlikely(!(rrdset_flags & RRDSET_FLAG_UPSTREAM_EXPOSED)))
417 rrdpush_send_chart_definition(wb, st);
418
423 - if (rrdset_flag_check(st, RRDSET_FLAG_SENDER_REPLICATION_FINISHED))
424 - rrdpush_send_chart_metrics(wb, st, host->sender);
419 + if (likely(rrdset_flags & RRDSET_FLAG_SENDER_REPLICATION_FINISHED))
420 + rrdpush_send_chart_metrics(wb, st, host->sender, rrdset_flags);
421
422 sender_commit(host->sender, wb);
423 }
streaming/rrdpush.h
-2
@@ -272,6 +272,4 @@ void log_sender_capabilities(struct sender_state *s);
272 STREAM_CAPABILITIES convert_stream_version_to_capabilities(int32_t version);
273 int32_t stream_capabilities_to_vn(uint32_t caps);
274
275 -void rrdpush_send_chart_metrics(BUFFER *wb, RRDSET *st, struct sender_state *s);
276 -
275 #endif //NETDATA_RRDPUSH_H
streaming/sender.c
+2 -13
@@ -1254,11 +1254,8 @@ void *rrdpush_sender_thread(void *ptr) {
1254 rrdpush_claimed_id(s->host);
1255 rrdpush_send_host_labels(s->host);
1256
1257 - // TO PUSH METRICS WITH DEFINITIONS:
1258 - //if(unlikely(s->rrdpush_sender_socket != -1 && __atomic_load_n(&s->host->rrdpush_sender_connected, __ATOMIC_SEQ_CST))) {
1259 - // thread_data->sending_definitions_status = SENDING_DEFINITIONS_DONE;
1260 - // rrdhost_flag_set(s->host, RRDHOST_FLAG_STREAM_COLLECTED_METRICS);
1261 - //}
1257 + rrdhost_flag_set(s->host, RRDHOST_FLAG_RRDPUSH_SENDER_READY_4_METRICS);
1258 + info("STREAM %s [send to %s]: enabling metrics streaming...", rrdhost_hostname(s->host), s->connected_to);
1259
1260 continue;
1261 }
@@ -1280,14 +1277,6 @@ void *rrdpush_sender_thread(void *ptr) {
1277
1278 if(outstanding)
1279 s->send_attempts++;
1283 - else {
1284 - if(unlikely(rrdhost_flag_check(s->host, RRDHOST_FLAG_RRDPUSH_SENDER_CONNECTED) &&
1285 - !rrdhost_flag_check(s->host, RRDHOST_FLAG_RRDPUSH_SENDER_READY_4_METRICS))) {
1286 - // let the data collection threads know we are ready to push metrics
1287 - rrdhost_flag_set(s->host, RRDHOST_FLAG_RRDPUSH_SENDER_READY_4_METRICS);
1288 - info("STREAM %s [send to %s]: enabling metrics streaming...", rrdhost_hostname(s->host), s->connected_to);
1289 - }
1290 - }
1280
1281 if(unlikely(s->rrdpush_sender_pipe[PIPE_READ] == -1)) {
1282 if(!rrdpush_sender_pipe_close(s->host, s->rrdpush_sender_pipe, true)) {