@cryptotaxi247 / netdata-1 / commits / 80d83b7bd

api v2 nodes for streaming statuses (#15162)

* api v2 nodes for streaming statuses * remove test * move parts of the output * in api/v2/data return 5 values per point when aggregation=percentage and raw option is given; return final values when aggregation=percentage is not the final grouping

Costa Tsaousis committed Jun 8, 2023 at 16:33 UTC 80d83b7bd1eca5872ed3ac5c34eb8bcb5fbd56e8
11 files changed +325 -110
collectors/plugins.d/pluginsd_parser.c
+6 -3
@@ -1467,9 +1467,11 @@ PARSER_RC pluginsd_replay_end(char **words, size_t num_words, void *user)
1467 time_t started = st->rrdhost->receiver->replication_first_time_t;
1468 time_t current = ((PARSER_USER_OBJECT *) user)->replay.end_time;
1469
1470 - if(started && current > started)
1470 + if(started && current > started) {
1471 + host->rrdpush_receiver_replication_percent = (NETDATA_DOUBLE) (current - started) * 100.0 / (NETDATA_DOUBLE) (now - started);
1472 worker_set_metric(WORKER_RECEIVER_JOB_REPLICATION_COMPLETION,
1472 - (NETDATA_DOUBLE)(current - started) * 100.0 / (NETDATA_DOUBLE)(now - started));
1473 + host->rrdpush_receiver_replication_percent);
1474 + }
1475 }
1476
1477 ((PARSER_USER_OBJECT *) user)->replay.start_time = 0;
@@ -1509,7 +1511,8 @@ PARSER_RC pluginsd_replay_end(char **words, size_t num_words, void *user)
1511
1512 pluginsd_set_chart_from_parent(user, NULL, PLUGINSD_KEYWORD_REPLAY_END);
1513
1512 - worker_set_metric(WORKER_RECEIVER_JOB_REPLICATION_COMPLETION, 100.0);
1514 + host->rrdpush_receiver_replication_percent = 100.0;
1515 + worker_set_metric(WORKER_RECEIVER_JOB_REPLICATION_COMPLETION, host->rrdpush_receiver_replication_percent);
1516
1517 return PARSER_RC_OK;
1518 }
database/contexts/api_v2.c
+3 -56
@@ -354,63 +354,10 @@ static ssize_t rrdcontext_to_json_v2_add_host(void *data, RRDHOST *host, bool qu
354
355 time_t now = now_realtime_sec();
356 buffer_json_member_add_object(wb, "status");
357 -
358 - size_t receiver_hops = host->system_info ? host->system_info->hops : (host == localhost) ? 0 : 1;
359 - buffer_json_member_add_object(wb, "collection");
360 - buffer_json_member_add_uint64(wb, "hops", receiver_hops);
361 - buffer_json_member_add_boolean(wb, "online", host == localhost || !rrdhost_flag_check(host, RRDHOST_FLAG_ORPHAN | RRDHOST_FLAG_RRDPUSH_RECEIVER_DISCONNECTED));
362 - buffer_json_member_add_boolean(wb, "replicating", rrdhost_receiver_replicating_charts(host));
363 - if(host != localhost && host->receiver) {
364 - buffer_json_member_add_object(wb, "source");
365 -
366 - char buf[1024 + 1];
367 - snprintfz(buf, 1024, "%s:%s", host->receiver->client_ip ? host->receiver->client_ip : "", host->receiver->client_port ? host->receiver->client_port : "");
368 - buffer_json_member_add_string(wb, "connection", buf);
369 - stream_capabilities_to_json_array(wb, host->receiver->capabilities, "capabilities");
370 -
371 - buffer_json_object_close(wb);
372 - }
373 - buffer_json_object_close(wb); // collection
374 -
375 - bool sender_connected = rrdhost_flag_check(host, RRDHOST_FLAG_RRDPUSH_SENDER_CONNECTED);
376 - buffer_json_member_add_object(wb, "streaming");
377 - buffer_json_member_add_uint64(wb, "hops", host->sender ? host->sender->hops : receiver_hops + 1);
378 - buffer_json_member_add_boolean(wb, "online", sender_connected);
379 - buffer_json_member_add_boolean(wb, "replicating", sender_connected && rrdhost_sender_replicating_charts(host));
380 -
381 - if(host->sender) {
382 - buffer_json_member_add_object(wb, "destination");
383 - buffer_json_member_add_string(wb, "connected_to", sender_connected ? host->sender->connected_to : "");
384 - stream_capabilities_to_json_array(wb, sender_connected ? host->sender->capabilities : 0, "capabilities");
385 -
386 - buffer_json_member_add_array(wb, "candidates");
387 - struct rrdpush_destinations *d;
388 - for(d = host->destinations ; d ; d = d->next) {
389 - buffer_json_add_array_item_object(wb);
390 -
391 - if(d->ssl) {
392 - char buf[1024 + 1];
393 - snprintfz(buf, 1024, "%s:SSL", string2str(d->destination));
394 - buffer_json_member_add_string(wb, "destination", buf);
395 - }
396 - else
397 - buffer_json_member_add_string(wb, "destination", string2str(d->destination));
398 -
399 - buffer_json_member_add_time_t(wb, "last_check", d->last_attempt);
400 - buffer_json_member_add_time_t(wb, "last_check_secs_ago", now - d->last_attempt);
401 - buffer_json_member_add_string(wb, "last_error", d->last_error);
402 - buffer_json_member_add_string(wb, "last_handshake", stream_handshake_error_to_string(d->last_handshake));
403 - buffer_json_member_add_time_t(wb, "next_check", d->postpone_reconnection_until);
404 - buffer_json_member_add_time_t(wb, "next_check_in_secs", (d->postpone_reconnection_until > now) ? d->postpone_reconnection_until - now : 0);
405 - buffer_json_object_close(wb);
406 - }
407 - buffer_json_array_close(wb);
408 -
409 - buffer_json_object_close(wb); // destination
357 + {
358 + rrdhost_receiver_to_json(wb, host, "collection", now);
359 + rrdhost_sender_to_json(wb, host, "streaming", now);
360 }
411 -
412 - buffer_json_object_close(wb); // streaming
413 -
361 buffer_json_object_close(wb); // status
362 }
363
database/contexts/rrdcontext.h
+11 -2
@@ -531,15 +531,24 @@ static inline bool query_has_group_by_aggregation_percentage(QUERY_TARGET *qt) {
531 // we need to send back "raw" output with "count"
532 // otherwise, we need to send back "raw" output with "hidden"
533
534 + bool last_is_percentage = false;
535 +
536 for(int g = 0; g < MAX_QUERY_GROUP_BY_PASSES ;g++) {
537 + if(qt->request.group_by[g].group_by == RRDR_GROUP_BY_NONE)
538 + break;
539 +
540 if(qt->request.group_by[g].group_by & RRDR_GROUP_BY_PERCENTAGE_OF_INSTANCE)
541 + // backwards compatibility
542 return false;
543
544 if(qt->request.group_by[g].aggregation == RRDR_GROUP_BY_FUNCTION_PERCENTAGE)
539 - return true;
545 + last_is_percentage = true;
546 +
547 + else
548 + last_is_percentage = false;
549 }
550
542 - return false;
551 + return last_is_percentage;
552 }
553
554 static inline bool query_target_has_percentage_of_group(QUERY_TARGET *qt) {
database/rrd.h
+2
@@ -1150,9 +1150,11 @@ struct rrdhost {
1150 struct rrdpush_destinations *destination; // the current destination from the above list
1151 SIMPLE_PATTERN *rrdpush_send_charts_matching; // pattern to match the charts to be sent
1152
1153 + const char *rrdpush_last_receiver_exit_reason;
1154 time_t rrdpush_seconds_to_replicate; // max time we want to replicate from the child
1155 time_t rrdpush_replication_step; // seconds per replication step
1156 size_t rrdpush_receiver_replicating_charts; // the number of charts currently being replicated from a child
1157 + NETDATA_DOUBLE rrdpush_receiver_replication_percent; // the % of replication completion
1158
1159 // the following are state information for the threading
1160 // streaming metrics from this netdata to an upstream netdata
database/rrdhost.c
+1
@@ -340,6 +340,7 @@ int is_legacy = 1;
340
341 host->rrdpush_seconds_to_replicate = rrdpush_seconds_to_replicate;
342 host->rrdpush_replication_step = rrdpush_replication_step;
343 + host->rrdpush_receiver_replication_percent = 100.0;
344
345 switch(memory_mode) {
346 default:
streaming/receiver.c
+68
@@ -423,6 +423,58 @@ static void rrdpush_receiver_replication_reset(RRDHOST *host) {
423 rrdhost_receiver_replicating_charts_zero(host);
424 }
425
426 +void rrdhost_receiver_to_json(BUFFER *wb, RRDHOST *host, const char *key, time_t now __maybe_unused) {
427 + size_t receiver_hops = host->system_info ? host->system_info->hops : (host == localhost) ? 0 : 1;
428 +
429 + netdata_mutex_lock(&host->receiver_lock);
430 +
431 + buffer_json_member_add_object(wb, key);
432 + buffer_json_member_add_uint64(wb, "hops", receiver_hops);
433 +
434 + bool online = host == localhost || !rrdhost_flag_check(host, RRDHOST_FLAG_ORPHAN | RRDHOST_FLAG_RRDPUSH_RECEIVER_DISCONNECTED);
435 + buffer_json_member_add_boolean(wb, "online", online);
436 +
437 + if(host->child_connect_time || host->child_disconnected_time) {
438 + time_t since = MAX(host->child_connect_time, host->child_disconnected_time);
439 + buffer_json_member_add_time_t(wb, "since", since);
440 + buffer_json_member_add_time_t(wb, "age", now - since);
441 + }
442 +
443 + if(!online && host->rrdpush_last_receiver_exit_reason)
444 + buffer_json_member_add_string(wb, "reason", host->rrdpush_last_receiver_exit_reason);
445 +
446 + if(host != localhost && host->receiver) {
447 + buffer_json_member_add_object(wb, "replication");
448 + {
449 + size_t instances = rrdhost_receiver_replicating_charts(host);
450 + buffer_json_member_add_boolean(wb, "in_progress", instances);
451 + buffer_json_member_add_double(wb, "completion", host->rrdpush_receiver_replication_percent);
452 + buffer_json_member_add_uint64(wb, "instances", instances);
453 + }
454 + buffer_json_object_close(wb); // replication
455 +
456 + buffer_json_member_add_object(wb, "source");
457 + {
458 +
459 + char buf[1024 + 1];
460 + SOCKET_PEERS peers = socket_peers(host->receiver->fd);
461 + bool ssl = SSL_connection(&host->receiver->ssl);
462 +
463 + snprintfz(buf, 1024, "[%s]:%d%s", peers.local.ip, peers.local.port, ssl ? ":SSL" : "");
464 + buffer_json_member_add_string(wb, "local", buf);
465 +
466 + snprintfz(buf, 1024, "[%s]:%d%s", peers.peer.ip, peers.peer.port, ssl ? ":SSL" : "");
467 + buffer_json_member_add_string(wb, "remote", buf);
468 +
469 + stream_capabilities_to_json_array(wb, host->receiver->capabilities, "capabilities");
470 + }
471 + buffer_json_object_close(wb); // source
472 + }
473 + buffer_json_object_close(wb); // collection
474 +
475 + netdata_mutex_unlock(&host->receiver_lock);
476 +}
477 +
478 static bool rrdhost_set_receiver(RRDHOST *host, struct receiver_state *rpt) {
479 bool signal_rrdcontext = false;
480 bool set_this = false;
@@ -496,6 +548,7 @@ static void rrdhost_clear_receiver(struct receiver_state *rpt) {
548
549 rrdhost_flag_set(host, RRDHOST_FLAG_ORPHAN);
550 host->receiver = NULL;
551 + host->rrdpush_last_receiver_exit_reason = rpt->exit.reason;
552 }
553
554 netdata_mutex_unlock(&host->receiver_lock);
@@ -546,6 +599,18 @@ bool stop_streaming_receiver(RRDHOST *host, const char *reason) {
599 return ret;
600 }
601
602 +static void rrdpush_send_error_on_taken_over_connection(struct receiver_state *rpt, const char *msg) {
603 + send_timeout(
604 +#ifdef ENABLE_HTTPS
605 + &rpt->ssl,
606 +#endif
607 + rpt->fd,
608 + (char *)msg,
609 + strlen(msg),
610 + 0,
611 + 5);
612 +}
613 +
614 void rrdpush_receive_log_status(struct receiver_state *rpt, const char *msg, const char *status) {
615
616 log_stream_connection(rpt->client_ip, rpt->client_port,
@@ -677,11 +742,13 @@ static void rrdpush_receive(struct receiver_state *rpt)
742
743 if(!host) {
744 rrdpush_receive_log_status(rpt, "failed to find/create host structure", "INTERNAL ERROR DROPPING CONNECTION");
745 + rrdpush_send_error_on_taken_over_connection(rpt, START_STREAMING_ERROR_INTERNAL_ERROR);
746 goto cleanup;
747 }
748
749 if (unlikely(rrdhost_flag_check(host, RRDHOST_FLAG_PENDING_CONTEXT_LOAD))) {
750 rrdpush_receive_log_status(rpt, "host is initializing", "INITIALIZATION IN PROGRESS RETRY LATER");
751 + rrdpush_send_error_on_taken_over_connection(rpt, START_STREAMING_ERROR_INITIALIZATION);
752 goto cleanup;
753 }
754
@@ -690,6 +757,7 @@ static void rrdpush_receive(struct receiver_state *rpt)
757
758 if(!rrdhost_set_receiver(host, rpt)) {
759 rrdpush_receive_log_status(rpt, "host is already served by another receiver", "DUPLICATE RECEIVER DROPPING CONNECTION");
760 + rrdpush_send_error_on_taken_over_connection(rpt, START_STREAMING_ERROR_ALREADY_STREAMING);
761 goto cleanup;
762 }
763 }
streaming/replication.c
+10 -1
@@ -274,6 +274,12 @@ static void replication_query_finalize(BUFFER *wb, struct replication_query *q,
274 replication_queries.queries_finished += queries;
275 replication_queries.points_read += q->points_read;
276 replication_queries.points_generated += q->points_generated;
277 +
278 + if(q->st && q->st->rrdhost->sender) {
279 + struct sender_state *s = q->st->rrdhost->sender;
280 + s->replication.latest_completed_before_t = q->query.before;
281 + }
282 +
283 netdata_spinlock_unlock(&replication_queries.spinlock);
284 }
285
@@ -644,7 +650,7 @@ bool replication_response_execute_and_finalize(struct replication_query *q, size
650 buffer_fast_strcat(wb, "\n", 1);
651
652 worker_is_busy(WORKER_JOB_BUFFER_COMMIT);
647 - sender_commit(host->sender, wb);
653 + sender_commit(host->sender, wb, STREAM_TRAFFIC_TYPE_REPLICATION);
654 worker_is_busy(WORKER_JOB_CLEANUP);
655
656 if(enable_streaming) {
@@ -1466,6 +1472,9 @@ void replication_add_request(struct sender_state *sender, const char *chart_id,
1472 .not_indexed_preprocessing = false,
1473 };
1474
1475 + if(!sender->replication.oldest_request_after_t || rq.after < sender->replication.oldest_request_after_t)
1476 + sender->replication.oldest_request_after_t = rq.after;
1477 +
1478 if(start_streaming && rrdpush_sender_get_buffer_used_percent(sender) <= STREAMING_START_MAX_SENDER_BUFFER_PERCENTAGE_ALLOWED)
1479 replication_execute_request(&rq, false);
1480
streaming/rrdpush.c
+20 -17
@@ -375,7 +375,7 @@ bool rrdset_push_chart_definition_now(RRDSET *st) {
375
376 BUFFER *wb = sender_start(host->sender);
377 rrdpush_send_chart_definition(wb, st);
378 - sender_commit(host->sender, wb);
378 + sender_commit(host->sender, wb, STREAM_TRAFFIC_TYPE_METADATA);
379 sender_thread_buffer_free();
380
381 return true;
@@ -443,7 +443,7 @@ void rrdset_push_metrics_finished(RRDSET_STREAM_BUFFER *rsb, RRDSET *st) {
443 buffer_fast_strcat(rsb->wb, PLUGINSD_KEYWORD_END_V2 "\n", sizeof(PLUGINSD_KEYWORD_END_V2) - 1 + 1);
444 }
445
446 - sender_commit(st->rrdhost->sender, rsb->wb);
446 + sender_commit(st->rrdhost->sender, rsb->wb, STREAM_TRAFFIC_TYPE_DATA);
447
448 *rsb = (RRDSET_STREAM_BUFFER){ .wb = NULL, };
449 }
@@ -483,7 +483,7 @@ RRDSET_STREAM_BUFFER rrdset_push_metric_initialize(RRDSET *st, time_t wall_clock
483 if(unlikely(!exposed_upstream)) {
484 BUFFER *wb = sender_start(host->sender);
485 replication_in_progress = rrdpush_send_chart_definition(wb, st);
486 - sender_commit(host->sender, wb);
486 + sender_commit(host->sender, wb, STREAM_TRAFFIC_TYPE_METADATA);
487 }
488
489 if(replication_in_progress)
@@ -514,7 +514,7 @@ void rrdpush_send_host_labels(RRDHOST *host) {
514 rrdlabels_walkthrough_read(host->rrdlabels, send_labels_callback, wb);
515 buffer_sprintf(wb, "OVERWRITE %s\n", "labels");
516
517 - sender_commit(host->sender, wb);
517 + sender_commit(host->sender, wb, STREAM_TRAFFIC_TYPE_METADATA);
518
519 sender_thread_buffer_free();
520 }
@@ -533,7 +533,7 @@ void rrdpush_claimed_id(RRDHOST *host)
533 buffer_sprintf(wb, "CLAIMED_ID %s %s\n", host->machine_guid, (host->aclk_state.claimed_id ? host->aclk_state.claimed_id : "NULL") );
534
535 rrdhost_aclk_state_unlock(host);
536 - sender_commit(host->sender, wb);
536 + sender_commit(host->sender, wb, STREAM_TRAFFIC_TYPE_METADATA);
537
538 sender_thread_buffer_free();
539 }
@@ -706,7 +706,7 @@ int rrdpush_receiver_permission_denied(struct web_client *w) {
706 // we always respond with the same message and error code
707 // to prevent an attacker from gaining info about the error
708 buffer_flush(w->response.data);
709 - buffer_sprintf(w->response.data, "You are not permitted to access this. Check the logs for more info.");
709 + buffer_strcat(w->response.data, START_STREAMING_ERROR_NOT_PERMITTED);
710 return HTTP_RESP_UNAUTHORIZED;
711 }
712
@@ -714,7 +714,7 @@ int rrdpush_receiver_too_busy_now(struct web_client *w) {
714 // we always respond with the same message and error code
715 // to prevent an attacker from gaining info about the error
716 buffer_flush(w->response.data);
717 - buffer_sprintf(w->response.data, "The server is too busy now to accept this request. Try later.");
717 + buffer_strcat(w->response.data, START_STREAMING_ERROR_BUSY_TRY_LATER);
718 return HTTP_RESP_SERVICE_UNAVAILABLE;
719 }
720
@@ -1146,7 +1146,7 @@ int rrdpush_receiver_thread_spawn(struct web_client *w, char *decoded_query_stri
1146
1147 // Have not set WEB_CLIENT_FLAG_DONT_CLOSE_SOCKET - caller should clean up
1148 buffer_flush(w->response.data);
1149 - buffer_strcat(w->response.data, "This GUID is already streaming to this server");
1149 + buffer_strcat(w->response.data, START_STREAMING_ERROR_ALREADY_STREAMING);
1150 receiver_state_free(rpt);
1151 return HTTP_RESP_CONFLICT;
1152 }
@@ -1191,15 +1191,18 @@ static struct {
1191 { STREAM_HANDSHAKE_OK_V3, "OK_V3" },
1192 { STREAM_HANDSHAKE_OK_V2, "OK_V2" },
1193 { STREAM_HANDSHAKE_OK_V1, "OK_V1" },
1194 - { STREAM_HANDSHAKE_ERROR_BAD_HANDSHAKE, "ERROR_BAD_HANDSHAKE" },
1195 - { STREAM_HANDSHAKE_ERROR_LOCALHOST, "ERROR_LOCALHOST" },
1196 - { STREAM_HANDSHAKE_ERROR_ALREADY_CONNECTED, "ERROR_ALREADY_CONNECTED" },
1197 - { STREAM_HANDSHAKE_ERROR_DENIED, "ERROR_DENIED" },
1198 - { STREAM_HANDSHAKE_ERROR_SEND_TIMEOUT, "ERROR_SEND_TIMEOUT" },
1199 - { STREAM_HANDSHAKE_ERROR_RECEIVE_TIMEOUT, "ERROR_RECEIVE_TIMEOUT" },
1200 - { STREAM_HANDSHAKE_ERROR_INVALID_CERTIFICATE, "ERROR_INVALID_CERTIFICATE" },
1201 - { STREAM_HANDSHAKE_ERROR_SSL_ERROR, "ERROR_SSL_ERROR" },
1202 - { STREAM_HANDSHAKE_ERROR_CANT_CONNECT, "ERROR_CANT_CONNECT" },
1194 + { STREAM_HANDSHAKE_ERROR_BAD_HANDSHAKE, "BAD HANDSHAKE" },
1195 + { STREAM_HANDSHAKE_ERROR_LOCALHOST, "LOCALHOST" },
1196 + { STREAM_HANDSHAKE_ERROR_ALREADY_CONNECTED, "ALREADY CONNECTED" },
1197 + { STREAM_HANDSHAKE_ERROR_DENIED, "DENIED" },
1198 + { STREAM_HANDSHAKE_ERROR_SEND_TIMEOUT, "SEND TIMEOUT" },
1199 + { STREAM_HANDSHAKE_ERROR_RECEIVE_TIMEOUT, "RECEIVE TIMEOUT" },
1200 + { STREAM_HANDSHAKE_ERROR_INVALID_CERTIFICATE, "INVALID CERTIFICATE" },
1201 + { STREAM_HANDSHAKE_ERROR_SSL_ERROR, "SSL ERROR" },
1202 + { STREAM_HANDSHAKE_ERROR_CANT_CONNECT, "CANT CONNECT" },
1203 + { STREAM_HANDSHAKE_BUSY_TRY_LATER, "BUSY TRY LATER" },
1204 + { STREAM_HANDSHAKE_INTERNAL_ERROR, "INTERNAL ERROR" },
1205 + { STREAM_HANDSHAKE_INITIALIZATION, "INITIALIZING" },
1206 { 0, NULL },
1207 };
1208
streaming/rrdpush.h
+28 -2
@@ -72,6 +72,9 @@ STREAM_CAPABILITIES stream_our_capabilities();
72 #define START_STREAMING_ERROR_SAME_LOCALHOST "Don't hit me baby, you are trying to stream my localhost back"
73 #define START_STREAMING_ERROR_ALREADY_STREAMING "This GUID is already streaming to this server"
74 #define START_STREAMING_ERROR_NOT_PERMITTED "You are not permitted to access this. Check the logs for more info."
75 +#define START_STREAMING_ERROR_BUSY_TRY_LATER "The server is too busy now to accept this request. Try later."
76 +#define START_STREAMING_ERROR_INTERNAL_ERROR "The server encountered an internal error. Try later."
77 +#define START_STREAMING_ERROR_INITIALIZATION "The server is initializing. Try later."
78
79 typedef enum {
80 STREAM_HANDSHAKE_OK_V5 = 5, // COMPRESSION
@@ -87,10 +90,25 @@ typedef enum {
90 STREAM_HANDSHAKE_ERROR_RECEIVE_TIMEOUT = -6,
91 STREAM_HANDSHAKE_ERROR_INVALID_CERTIFICATE = -7,
92 STREAM_HANDSHAKE_ERROR_SSL_ERROR = -8,
90 - STREAM_HANDSHAKE_ERROR_CANT_CONNECT = -9
93 + STREAM_HANDSHAKE_ERROR_CANT_CONNECT = -9,
94 + STREAM_HANDSHAKE_BUSY_TRY_LATER = -10,
95 + STREAM_HANDSHAKE_INTERNAL_ERROR = -11,
96 + STREAM_HANDSHAKE_INITIALIZATION = -12,
97 } STREAM_HANDSHAKE;
98
99
100 +// ----------------------------------------------------------------------------
101 +
102 +typedef enum __attribute__((packed)) {
103 + STREAM_TRAFFIC_TYPE_REPLICATION,
104 + STREAM_TRAFFIC_TYPE_FUNCTIONS,
105 + STREAM_TRAFFIC_TYPE_METADATA,
106 + STREAM_TRAFFIC_TYPE_DATA,
107 +
108 + // terminator
109 + STREAM_TRAFFIC_TYPE_MAX,
110 +} STREAM_TRAFFIC_TYPE;
111 +
112 // ----------------------------------------------------------------------------
113
114 typedef struct {
@@ -148,6 +166,7 @@ struct sender_state {
166 size_t sent_bytes_on_this_connection;
167 size_t send_attempts;
168 time_t last_traffic_seen_t;
169 + time_t last_state_since_t; // the timestamp of the last state (online/offline) change
170 size_t not_connected_loops;
171 // Metrics are collected asynchronously by collector threads calling rrdset_done_push(). This can also trigger
172 // the lazy creation of the sender thread - both cases (buffer access and thread creation) are guarded here.
@@ -157,6 +176,8 @@ struct sender_state {
176 int read_len;
177 STREAM_CAPABILITIES capabilities;
178
179 + size_t sent_bytes_on_this_connection_per_type[STREAM_TRAFFIC_TYPE_MAX];
180 +
181 int rrdpush_sender_pipe[2]; // collector to sender thread signaling
182 int rrdpush_sender_socket;
183
@@ -176,6 +197,8 @@ struct sender_state {
197
198 struct {
199 DICTIONARY *requests; // de-duplication of replication requests, per chart
200 + time_t oldest_request_after_t; // the timestamp of the oldest replication request
201 + time_t latest_completed_before_t; // the timestamp of the latest replication request
202
203 struct {
204 size_t pending_requests; // the currently outstanding replication requests
@@ -306,7 +329,7 @@ void rrdpush_destinations_init(RRDHOST *host);
329 void rrdpush_destinations_free(RRDHOST *host);
330
331 BUFFER *sender_start(struct sender_state *s);
309 -void sender_commit(struct sender_state *s, BUFFER *wb);
332 +void sender_commit(struct sender_state *s, BUFFER *wb, STREAM_TRAFFIC_TYPE type);
333 int rrdpush_init();
334 bool rrdpush_receiver_needs_dbengine();
335 int configured_as_parent();
@@ -368,6 +391,9 @@ bool stop_streaming_receiver(RRDHOST *host, const char *reason);
391
392 void sender_thread_buffer_free(void);
393
394 +void rrdhost_receiver_to_json(BUFFER *wb, RRDHOST *host, const char *key, time_t now __maybe_unused);
395 +void rrdhost_sender_to_json(BUFFER *wb, RRDHOST *host, const char *key, time_t now __maybe_unused);
396 +
397 #include "replication.h"
398
399 #endif //NETDATA_RRDPUSH_H
streaming/sender.c
+159 -11
@@ -84,7 +84,7 @@ static inline void deactivate_compression(struct sender_state *s) {
84 #define SENDER_BUFFER_ADAPT_TO_TIMES_MAX_SIZE 3
85
86 // Collector thread finishing a transmission
87 -void sender_commit(struct sender_state *s, BUFFER *wb) {
87 +void sender_commit(struct sender_state *s, BUFFER *wb, STREAM_TRAFFIC_TYPE type) {
88
89 if(unlikely(wb != sender_thread_buffer))
90 fatal("STREAMING: sender is trying to commit a buffer that is not this thread's buffer.");
@@ -163,6 +163,8 @@ void sender_commit(struct sender_state *s, BUFFER *wb) {
163
164 if(cbuffer_add_unsafe(s->buffer, dst, dst_len))
165 s->flags |= SENDER_FLAG_OVERFLOW;
166 + else
167 + s->sent_bytes_on_this_connection_per_type[type] += dst_len;
168
169 src = src + size_to_compress;
170 src_len -= size_to_compress;
@@ -170,9 +172,13 @@ void sender_commit(struct sender_state *s, BUFFER *wb) {
172 }
173 else if(cbuffer_add_unsafe(s->buffer, src, src_len))
174 s->flags |= SENDER_FLAG_OVERFLOW;
175 + else
176 + s->sent_bytes_on_this_connection_per_type[type] += src_len;
177 #else
178 if(cbuffer_add_unsafe(s->buffer, src, src_len))
179 s->flags |= SENDER_FLAG_OVERFLOW;
180 + else
181 + s->sent_bytes_on_this_connection_per_type[type] += src_len;
182 #endif
183
184 replication_recalculate_buffer_used_ratio_unsafe(s);
@@ -204,7 +210,7 @@ void rrdpush_sender_send_this_host_variable_now(RRDHOST *host, const RRDVAR_ACQU
210 if(rrdhost_can_send_definitions_to_parent(host)) {
211 BUFFER *wb = sender_start(host->sender);
212 rrdpush_sender_add_host_variable_to_buffer(wb, rva);
207 - sender_commit(host->sender, wb);
213 + sender_commit(host->sender, wb, STREAM_TRAFFIC_TYPE_METADATA);
214 sender_thread_buffer_free();
215 }
216 }
@@ -233,7 +239,7 @@ static void rrdpush_sender_thread_send_custom_host_variables(RRDHOST *host) {
239 };
240 int ret = rrdvar_walkthrough_read(host->rrdvars, rrdpush_sender_thread_custom_host_variables_callback, &tmp);
241 (void)ret;
236 - sender_commit(host->sender, wb);
242 + sender_commit(host->sender, wb, STREAM_TRAFFIC_TYPE_METADATA);
243 sender_thread_buffer_free();
244
245 debug(D_STREAM, "RRDVAR sent %d VARIABLES", ret);
@@ -426,6 +432,33 @@ struct {
432 .worker_job_id = WORKER_SENDER_JOB_DISCONNECT_BAD_HANDSHAKE,
433 .postpone_reconnect_seconds = 1 * 60, // 1 minute
434 },
435 + {
436 + .response = START_STREAMING_ERROR_BUSY_TRY_LATER,
437 + .length = sizeof(START_STREAMING_ERROR_BUSY_TRY_LATER) - 1,
438 + .version = STREAM_HANDSHAKE_BUSY_TRY_LATER,
439 + .dynamic = false,
440 + .error = "remote server is currently busy, we should try later",
441 + .worker_job_id = WORKER_SENDER_JOB_DISCONNECT_BAD_HANDSHAKE,
442 + .postpone_reconnect_seconds = 2 * 60, // 2 minutes
443 + },
444 + {
445 + .response = START_STREAMING_ERROR_INTERNAL_ERROR,
446 + .length = sizeof(START_STREAMING_ERROR_INTERNAL_ERROR) - 1,
447 + .version = STREAM_HANDSHAKE_INTERNAL_ERROR,
448 + .dynamic = false,
449 + .error = "remote server is encountered an internal error, we should try later",
450 + .worker_job_id = WORKER_SENDER_JOB_DISCONNECT_BAD_HANDSHAKE,
451 + .postpone_reconnect_seconds = 5 * 60, // 5 minutes
452 + },
453 + {
454 + .response = START_STREAMING_ERROR_INITIALIZATION,
455 + .length = sizeof(START_STREAMING_ERROR_INITIALIZATION) - 1,
456 + .version = STREAM_HANDSHAKE_INITIALIZATION,
457 + .dynamic = false,
458 + .error = "remote server is initializing, we should try later",
459 + .worker_job_id = WORKER_SENDER_JOB_DISCONNECT_BAD_HANDSHAKE,
460 + .postpone_reconnect_seconds = 2 * 60, // 2 minute
461 + },
462
463 // terminator
464 {
@@ -750,6 +783,10 @@ static bool attempt_to_connect(struct sender_state *state)
783 {
784 state->send_attempts = 0;
785
786 + // reset the bytes we have sent for this session
787 + state->sent_bytes_on_this_connection = 0;
788 + memset(state->sent_bytes_on_this_connection_per_type, 0, sizeof(state->sent_bytes_on_this_connection_per_type));
789 +
790 if(rrdpush_sender_thread_connect_to_parent(state->host, state->default_port, state->timeout, state)) {
791 // reset the buffer, to properly send charts and metrics
792 rrdpush_sender_on_connect(state->host);
@@ -760,9 +797,6 @@ static bool attempt_to_connect(struct sender_state *state)
797 // make sure the next reconnection will be immediate
798 state->not_connected_loops = 0;
799
763 - // reset the bytes we have sent for this session
764 - state->sent_bytes_on_this_connection = 0;
765 -
800 // let the data collection threads know we are ready
801 rrdhost_flag_set(state->host, RRDHOST_FLAG_RRDPUSH_SENDER_CONNECTED);
802
@@ -776,9 +810,6 @@ static bool attempt_to_connect(struct sender_state *state)
810 // increase the failed connections counter
811 state->not_connected_loops++;
812
779 - // reset the number of bytes sent
780 - state->sent_bytes_on_this_connection = 0;
781 -
813 // slow re-connection on repeating errors
814 usec_t now_ut = now_monotonic_usec();
815 usec_t end_ut = now_ut + USEC_PER_SEC * state->reconnect_delay;
@@ -899,7 +930,7 @@ void stream_execute_function_callback(BUFFER *func_wb, int code, void *data) {
930 buffer_fast_strcat(wb, buffer_tostring(func_wb), buffer_strlen(func_wb));
931 pluginsd_function_result_end_to_buffer(wb);
932
902 - sender_commit(s, wb);
933 + sender_commit(s, wb, STREAM_TRAFFIC_TYPE_FUNCTIONS);
934 sender_thread_buffer_free();
935
936 internal_error(true, "STREAM %s [send to %s] FUNCTION transaction %s sending back response (%zu bytes, %llu usec).",
@@ -1067,6 +1098,120 @@ void rrdpush_signal_sender_to_wake_up(struct sender_state *s) {
1098 }
1099 }
1100
1101 +static NETDATA_DOUBLE rrdhost_sender_replication_completion(RRDHOST *host, time_t now, size_t *instances) {
1102 + size_t charts = rrdhost_sender_replicating_charts(host);
1103 + NETDATA_DOUBLE completion;
1104 + if(!charts || !host->sender->replication.oldest_request_after_t)
1105 + completion = 100.0;
1106 + else if(!host->sender->replication.latest_completed_before_t || host->sender->replication.latest_completed_before_t < host->sender->replication.oldest_request_after_t)
1107 + completion = 0.0;
1108 + else {
1109 + time_t total = now - host->sender->replication.oldest_request_after_t;
1110 + time_t current = host->sender->replication.latest_completed_before_t - host->sender->replication.oldest_request_after_t;
1111 + completion = (NETDATA_DOUBLE) current * 100.0 / (NETDATA_DOUBLE) total;
1112 + }
1113 +
1114 + *instances = charts;
1115 +
1116 + return completion;
1117 +}
1118 +
1119 +void rrdhost_sender_to_json(BUFFER *wb, RRDHOST *host, const char *key, time_t now __maybe_unused) {
1120 + bool online = rrdhost_flag_check(host, RRDHOST_FLAG_RRDPUSH_SENDER_CONNECTED);
1121 + buffer_json_member_add_object(wb, key);
1122 +
1123 + if(host->sender)
1124 + buffer_json_member_add_uint64(wb, "hops", host->sender->hops);
1125 +
1126 + buffer_json_member_add_boolean(wb, "online", online);
1127 +
1128 + if(host->sender && host->sender->last_state_since_t) {
1129 + buffer_json_member_add_time_t(wb, "since", host->sender->last_state_since_t);
1130 + buffer_json_member_add_time_t(wb, "age", now - host->sender->last_state_since_t);
1131 + }
1132 +
1133 + if(!online && host->sender && host->sender->exit.reason)
1134 + buffer_json_member_add_string(wb, "reason", host->sender->exit.reason);
1135 +
1136 + buffer_json_member_add_object(wb, "replication");
1137 + {
1138 + size_t instances;
1139 + NETDATA_DOUBLE completion = rrdhost_sender_replication_completion(host, now, &instances);
1140 + buffer_json_member_add_boolean(wb, "in_progress", instances);
1141 + buffer_json_member_add_double(wb, "completion", completion);
1142 + buffer_json_member_add_uint64(wb, "instances", instances);
1143 + }
1144 + buffer_json_object_close(wb);
1145 +
1146 + if(host->sender) {
1147 + netdata_mutex_lock(&host->sender->mutex);
1148 +
1149 + buffer_json_member_add_object(wb, "destination");
1150 + {
1151 + char buf[1024 + 1];
1152 + if(online && host->sender->rrdpush_sender_socket != -1) {
1153 + SOCKET_PEERS peers = socket_peers(host->sender->rrdpush_sender_socket);
1154 + bool ssl = SSL_connection(&host->sender->ssl);
1155 +
1156 + snprintfz(buf, 1024, "[%s]:%d%s", peers.local.ip, peers.local.port, ssl ? ":SSL" : "");
1157 + buffer_json_member_add_string(wb, "local", buf);
1158 +
1159 + snprintfz(buf, 1024, "[%s]:%d%s", peers.peer.ip, peers.peer.port, ssl ? ":SSL" : "");
1160 + buffer_json_member_add_string(wb, "remote", buf);
1161 +
1162 + stream_capabilities_to_json_array(wb, online ? host->sender->capabilities : 0,
1163 + "capabilities");
1164 +
1165 + buffer_json_member_add_object(wb, "traffic");
1166 + {
1167 + bool compression = false;
1168 +#ifdef ENABLE_COMPRESSION
1169 + compression = (stream_has_capability(host->sender, STREAM_CAP_COMPRESSION) && host->sender->compressor);
1170 +#endif
1171 + buffer_json_member_add_boolean(wb, "compression", compression);
1172 + buffer_json_member_add_uint64(wb, "data", host->sender->sent_bytes_on_this_connection_per_type[STREAM_TRAFFIC_TYPE_DATA]);
1173 + buffer_json_member_add_uint64(wb, "metadata", host->sender->sent_bytes_on_this_connection_per_type[STREAM_TRAFFIC_TYPE_METADATA]);
1174 + buffer_json_member_add_uint64(wb, "functions", host->sender->sent_bytes_on_this_connection_per_type[STREAM_TRAFFIC_TYPE_FUNCTIONS]);
1175 + buffer_json_member_add_uint64(wb, "replication", host->sender->sent_bytes_on_this_connection_per_type[STREAM_TRAFFIC_TYPE_REPLICATION]);
1176 + }
1177 + buffer_json_object_close(wb); // traffic
1178 + }
1179 +
1180 + buffer_json_member_add_array(wb, "candidates");
1181 + struct rrdpush_destinations *d;
1182 + for (d = host->destinations; d; d = d->next) {
1183 + buffer_json_add_array_item_object(wb);
1184 + {
1185 +
1186 + if (d->ssl) {
1187 + snprintfz(buf, 1024, "%s:SSL", string2str(d->destination));
1188 + buffer_json_member_add_string(wb, "destination", buf);
1189 + }
1190 + else
1191 + buffer_json_member_add_string(wb, "destination", string2str(d->destination));
1192 +
1193 + buffer_json_member_add_time_t(wb, "last_check", d->last_attempt);
1194 + buffer_json_member_add_time_t(wb, "age", now - d->last_attempt);
1195 + buffer_json_member_add_string(wb, "last_error", d->last_error);
1196 + buffer_json_member_add_string(wb, "last_handshake",
1197 + stream_handshake_error_to_string(d->last_handshake));
1198 + buffer_json_member_add_time_t(wb, "next_check", d->postpone_reconnection_until);
1199 + buffer_json_member_add_time_t(wb, "next_in",
1200 + (d->postpone_reconnection_until > now) ?
1201 + d->postpone_reconnection_until - now : 0);
1202 + }
1203 + buffer_json_object_close(wb); // each candidate
1204 + }
1205 + buffer_json_array_close(wb); // candidates
1206 + }
1207 + buffer_json_object_close(wb); // destination
1208 +
1209 + netdata_mutex_unlock(&host->sender->mutex);
1210 + }
1211 +
1212 + buffer_json_object_close(wb); // streaming
1213 +}
1214 +
1215 static bool rrdhost_set_sender(RRDHOST *host) {
1216 if(unlikely(!host->sender)) return false;
1217
@@ -1076,6 +1221,8 @@ static bool rrdhost_set_sender(RRDHOST *host) {
1221 rrdhost_flag_clear(host, RRDHOST_FLAG_RRDPUSH_SENDER_CONNECTED | RRDHOST_FLAG_RRDPUSH_SENDER_READY_4_METRICS);
1222 rrdhost_flag_set(host, RRDHOST_FLAG_RRDPUSH_SENDER_SPAWN);
1223 host->sender->tid = gettid();
1224 + host->sender->last_state_since_t = now_realtime_sec();
1225 + host->sender->exit.reason = NULL;
1226 ret = true;
1227 }
1228 netdata_mutex_unlock(&host->sender->mutex);
@@ -1091,8 +1238,8 @@ static void rrdhost_clear_sender___while_having_sender_mutex(RRDHOST *host) {
1238 if(host->sender->tid == gettid()) {
1239 host->sender->tid = 0;
1240 host->sender->exit.shutdown = false;
1094 - host->sender->exit.reason = NULL;
1241 rrdhost_flag_clear(host, RRDHOST_FLAG_RRDPUSH_SENDER_SPAWN | RRDHOST_FLAG_RRDPUSH_SENDER_CONNECTED | RRDHOST_FLAG_RRDPUSH_SENDER_READY_4_METRICS);
1242 + host->sender->last_state_since_t = now_realtime_sec();
1243 }
1244
1245 rrdpush_reset_destinations_postpone_time(host);
@@ -1291,6 +1438,7 @@ void *rrdpush_sender_thread(void *ptr) {
1438 now_s = s->last_traffic_seen_t = now_monotonic_sec();
1439 rrdpush_claimed_id(s->host);
1440 rrdpush_send_host_labels(s->host);
1441 + s->replication.oldest_request_after_t = 0;
1442
1443 rrdhost_flag_set(s->host, RRDHOST_FLAG_RRDPUSH_SENDER_READY_4_METRICS);
1444 info("STREAM %s [send to %s]: enabling metrics streaming...", rrdhost_hostname(s->host), s->connected_to);
web/api/formatters/json/json.c
+17 -18
@@ -248,8 +248,8 @@ void rrdr2json_v2(RRDR *r, BUFFER *wb) {
248 QUERY_TARGET *qt = r->internal.qt;
249 RRDR_OPTIONS options = qt->window.options;
250
251 - bool send_fourth_number = query_target_aggregatable(qt);
252 - bool fourth_number_is_vh = send_fourth_number && r->vh && query_has_group_by_aggregation_percentage(qt);
251 + bool send_count = query_target_aggregatable(qt);
252 + bool send_hidden = send_count && r->vh && query_has_group_by_aggregation_percentage(qt);
253
254 buffer_json_member_add_object(wb, "result");
255
@@ -267,16 +267,17 @@ void rrdr2json_v2(RRDR *r, BUFFER *wb) {
267 buffer_json_array_close(wb); // labels
268
269 buffer_json_member_add_object(wb, "point");
270 - buffer_json_member_add_uint64(wb, "value", 0);
271 - buffer_json_member_add_uint64(wb, "arp", 1);
272 - buffer_json_member_add_uint64(wb, "pa", 2);
273 - if(send_fourth_number) {
274 - if(fourth_number_is_vh)
275 - buffer_json_member_add_uint64(wb, "hidden", 3);
276 - else
277 - buffer_json_member_add_uint64(wb, "count", 3);
270 + {
271 + size_t point_count = 0;
272 + buffer_json_member_add_uint64(wb, "value", point_count++);
273 + buffer_json_member_add_uint64(wb, "arp", point_count++);
274 + buffer_json_member_add_uint64(wb, "pa", point_count++);
275 + if (send_count)
276 + buffer_json_member_add_uint64(wb, "count", point_count++);
277 + if (send_hidden)
278 + buffer_json_member_add_uint64(wb, "hidden", point_count++);
279 }
279 - buffer_json_object_close(wb);
280 + buffer_json_object_close(wb); // point
281
282 buffer_json_member_add_array(wb, "data");
283 if(i) {
@@ -290,7 +291,7 @@ void rrdr2json_v2(RRDR *r, BUFFER *wb) {
291 // for each line in the array
292 for (i = start; i != end; i += step) {
293 NETDATA_DOUBLE *cn = &r->v[ i * r->d ];
293 - NETDATA_DOUBLE *ch = fourth_number_is_vh ? &r->vh[i * r->d ] : NULL;
294 + NETDATA_DOUBLE *ch = send_hidden ? &r->vh[i * r->d ] : NULL;
295 RRDR_VALUE_FLAGS *co = &r->o[ i * r->d ];
296 NETDATA_DOUBLE *ar = &r->ar[ i * r->d ];
297 uint32_t *gbc = &r->gbc [ i * r->d ];
@@ -330,12 +331,10 @@ void rrdr2json_v2(RRDR *r, BUFFER *wb) {
331 buffer_json_add_array_item_uint64(wb, o);
332
333 // add the count
333 - if(send_fourth_number) {
334 - if(fourth_number_is_vh)
335 - buffer_json_add_array_item_double(wb, ch[d]);
336 - else
337 - buffer_json_add_array_item_uint64(wb, gbc[d]);
338 - }
334 + if(send_count)
335 + buffer_json_add_array_item_uint64(wb, gbc[d]);
336 + if(send_hidden)
337 + buffer_json_add_array_item_double(wb, ch[d]);
338
339 buffer_json_array_close(wb); // point
340 }