@cryptotaxi247 / netdata-1 / commits / 7c51e14c4

PULSE: network traffic (#19419)

* pulse now tracks network traffic for web server, statsd, streaming and aclk * show gaps on the network traffic chart when aclk is not connected * fix contexts shutdown * log nodes info every 10 seconds

Costa Tsaousis committed Jan 16, 2025 at 22:35 UTC 7c51e14c4004911ac17f8c230ca761c1206d82f2
16 files changed +273 -54
CMakeLists.txt
+2
@@ -1168,6 +1168,8 @@ set(DAEMON_FILES
1168 src/daemon/config/netdata-conf-profile.c
1169 src/daemon/config/netdata-conf-profile.h
1170 src/daemon/pulse/pulse-daemon-memory-system.c
1171 + src/daemon/pulse/pulse-network.c
1172 + src/daemon/pulse/pulse-network.h
1173 )
1174
1175 set(H2O_FILES
src/aclk/aclk.c
+11 -1
@@ -24,12 +24,17 @@ int aclk_pubacks_per_conn = 0; // How many PubAcks we got since MQTT conn est.
24 int aclk_rcvd_cloud_msgs = 0;
25 int aclk_connection_counter = 0;
26
27 +mqtt_wss_client mqttwss_client;
28 +
29 static bool aclk_connected = false;
30 static inline void aclk_set_connected(void) {
31 __atomic_store_n(&aclk_connected, true, __ATOMIC_RELAXED);
32 }
33 static inline void aclk_set_disconnected(void) {
34 __atomic_store_n(&aclk_connected, false, __ATOMIC_RELAXED);
35 +
36 + if(mqttwss_client)
37 + mqtt_wss_reset_stats(mqttwss_client);
38 }
39
40 inline bool aclk_online(void) {
@@ -64,7 +69,12 @@ float last_backoff_value = 0;
69
70 time_t aclk_block_until = 0;
71
67 -mqtt_wss_client mqttwss_client;
72 +struct mqtt_wss_stats aclk_statistics(void) {
73 + if(mqttwss_client)
74 + return mqtt_wss_get_stats(mqttwss_client);
75 + else
76 + return (struct mqtt_wss_stats) { 0 };
77 +}
78
79 struct aclk_shared_state aclk_shared_state = {
80 .mqtt_shutdown_msg_id = -1,
src/aclk/aclk.h
+2
@@ -94,4 +94,6 @@ char *aclk_state_json(void);
94 void add_aclk_host_labels(void);
95 void aclk_queue_node_info(RRDHOST *host, bool immediate);
96
97 +struct mqtt_wss_stats aclk_statistics(void);
98 +
99 #endif /* ACLK_H */
src/aclk/mqtt_websockets/mqtt_wss_client.c
+7 -1
@@ -989,12 +989,18 @@ struct mqtt_wss_stats mqtt_wss_get_stats(mqtt_wss_client client)
989 struct mqtt_wss_stats current;
990 spinlock_lock(&client->stat_lock);
991 current = client->stats;
992 - memset(&client->stats, 0, sizeof(client->stats));
992 spinlock_unlock(&client->stat_lock);
993 mqtt_ng_get_stats(client->mqtt, &current.mqtt);
994 return current;
995 }
996
997 +void mqtt_wss_reset_stats(mqtt_wss_client client)
998 +{
999 + spinlock_lock(&client->stat_lock);
1000 + memset(&client->stats, 0, sizeof(client->stats));
1001 + spinlock_unlock(&client->stat_lock);
1002 +}
1003 +
1004 int mqtt_wss_set_topic_alias(mqtt_wss_client client, const char *topic)
1005 {
1006 return mqtt_ng_set_topic_alias(client->mqtt, topic);
src/aclk/mqtt_websockets/mqtt_wss_client.h
+1
@@ -141,6 +141,7 @@ struct mqtt_wss_stats {
141 };
142
143 struct mqtt_wss_stats mqtt_wss_get_stats(mqtt_wss_client client);
144 +void mqtt_wss_reset_stats(mqtt_wss_client client);
145
146 #ifdef MQTT_WSS_DEBUG
147 #include <openssl/ssl.h>
src/collectors/statsd.plugin/statsd.c
+8 -1
@@ -965,6 +965,8 @@ static int statsd_rcv_callback(POLLINFO *pi, nd_poll_event_t *events) {
965 d->len += rc;
966 statsd.tcp_socket_reads++;
967 statsd.tcp_bytes_read += rc;
968 +
969 + pulse_statsd_received_bytes(rc);
970 }
971
972 if(likely(d->len > 0)) {
@@ -1016,12 +1018,15 @@ static int statsd_rcv_callback(POLLINFO *pi, nd_poll_event_t *events) {
1018 statsd.udp_socket_reads++;
1019 statsd.udp_packets_received += rc;
1020
1019 - size_t i;
1021 + size_t i, total_size = 0;
1022 for (i = 0; i < (size_t)rc; ++i) {
1023 size_t len = (size_t)d->msgs[i].msg_len;
1024 statsd.udp_bytes_read += len;
1025 + total_size += len;
1026 statsd_process(d->msgs[i].msg_hdr.msg_iov->iov_base, len, 0);
1027 }
1028 +
1029 + pulse_statsd_received_bytes(total_size);
1030 }
1031 } while (rc != -1);
1032
@@ -1043,6 +1048,8 @@ static int statsd_rcv_callback(POLLINFO *pi, nd_poll_event_t *events) {
1048 statsd.udp_packets_received++;
1049 statsd.udp_bytes_read += rc;
1050 statsd_process(d->buffer, (size_t) rc, 0);
1051 +
1052 + pulse_statsd_received_bytes(rc);
1053 }
1054 } while (rc != -1);
1055 #endif
src/daemon/pulse/pulse-http-api.c
+6 -40
@@ -13,8 +13,6 @@ static struct web_statistics {
13 PAD64(uint64_t) web_requests;
14 PAD64(uint64_t) web_usec;
15 PAD64(uint64_t) web_usec_max;
16 - PAD64(uint64_t) bytes_received;
17 - PAD64(uint64_t) bytes_sent;
16
17 PAD64(uint64_t) content_size_uncompressed;
18 PAD64(uint64_t) content_size_compressed;
@@ -29,8 +27,8 @@ void pulse_web_client_disconnected(void) {
27 }
28
29 void pulse_web_request_completed(uint64_t dt,
32 - uint64_t bytes_received,
33 - uint64_t bytes_sent,
30 + uint64_t bytes_received __maybe_unused,
31 + uint64_t bytes_sent __maybe_unused,
32 uint64_t content_size,
33 uint64_t compressed_content_size) {
34 uint64_t old_web_usec_max = live_stats.web_usec_max;
@@ -39,8 +37,8 @@ void pulse_web_request_completed(uint64_t dt,
37
38 __atomic_fetch_add(&live_stats.web_requests, 1, __ATOMIC_RELAXED);
39 __atomic_fetch_add(&live_stats.web_usec, dt, __ATOMIC_RELAXED);
42 - __atomic_fetch_add(&live_stats.bytes_received, bytes_received, __ATOMIC_RELAXED);
43 - __atomic_fetch_add(&live_stats.bytes_sent, bytes_sent, __ATOMIC_RELAXED);
40 +// __atomic_fetch_add(&live_stats.bytes_received, bytes_received, __ATOMIC_RELAXED);
41 +// __atomic_fetch_add(&live_stats.bytes_sent, bytes_sent, __ATOMIC_RELAXED);
42 __atomic_fetch_add(&live_stats.content_size_uncompressed, content_size, __ATOMIC_RELAXED);
43 __atomic_fetch_add(&live_stats.content_size_compressed, compressed_content_size, __ATOMIC_RELAXED);
44 }
@@ -50,8 +48,8 @@ static inline void pulse_web_copy(struct web_statistics *gs, uint8_t options) {
48 gs->web_requests = __atomic_load_n(&live_stats.web_requests, __ATOMIC_RELAXED);
49 gs->web_usec = __atomic_load_n(&live_stats.web_usec, __ATOMIC_RELAXED);
50 gs->web_usec_max = __atomic_load_n(&live_stats.web_usec_max, __ATOMIC_RELAXED);
53 - gs->bytes_received = __atomic_load_n(&live_stats.bytes_received, __ATOMIC_RELAXED);
54 - gs->bytes_sent = __atomic_load_n(&live_stats.bytes_sent, __ATOMIC_RELAXED);
51 +// gs->bytes_received = __atomic_load_n(&live_stats.bytes_received, __ATOMIC_RELAXED);
52 +// gs->bytes_sent = __atomic_load_n(&live_stats.bytes_sent, __ATOMIC_RELAXED);
53 gs->content_size_uncompressed = __atomic_load_n(&live_stats.content_size_uncompressed, __ATOMIC_RELAXED);
54 gs->content_size_compressed = __atomic_load_n(&live_stats.content_size_compressed, __ATOMIC_RELAXED);
55
@@ -125,38 +123,6 @@ void pulse_web_do(bool extended) {
123
124 // ----------------------------------------------------------------
125
128 - {
129 - static RRDSET *st_bytes = NULL;
130 - static RRDDIM *rd_in = NULL,
131 - *rd_out = NULL;
132 -
133 - if (unlikely(!st_bytes)) {
134 - st_bytes = rrdset_create_localhost(
135 - "netdata"
136 - , "net"
137 - , NULL
138 - , "HTTP API"
139 - , "netdata.http_api_traffic"
140 - , "Netdata Web API Network Traffic"
141 - , "kilobits/s"
142 - , "netdata"
143 - , "pulse"
144 - , 130400
145 - , localhost->rrd_update_every
146 - , RRDSET_TYPE_AREA
147 - );
148 -
149 - rd_in = rrddim_add(st_bytes, "in", NULL, 8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
150 - rd_out = rrddim_add(st_bytes, "out", NULL, -8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
151 - }
152 -
153 - rrddim_set_by_pointer(st_bytes, rd_in, (collected_number) gs.bytes_received);
154 - rrddim_set_by_pointer(st_bytes, rd_out, (collected_number) gs.bytes_sent);
155 - rrdset_done(st_bytes);
156 - }
157 -
158 - // ----------------------------------------------------------------
159 -
126 {
127 static unsigned long long old_web_requests = 0, old_web_usec = 0;
128 static collected_number average_response_time = -1;
src/daemon/pulse/pulse-network.c new
+189
@@ -0,0 +1,189 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#define PULSE_INTERNALS 1
4 +#include "pulse-network.h"
5 +
6 +#define PULSE_NETWORK_CHART_TITLE "Netdata Network Traffic"
7 +#define PULSE_NETWORK_CHART_FAMILY "Network Traffic"
8 +#define PULSE_NETWORK_CHART_CONTEXT "netdata.network"
9 +#define PULSE_NETWORK_CHART_UNITS "kilobits/s"
10 +#define PULSE_NETWORK_CHART_PRIORITY 130150
11 +
12 +static struct network_statistics {
13 + bool extended;
14 + PAD64(uint64_t) api_bytes_received;
15 + PAD64(uint64_t) api_bytes_sent;
16 + PAD64(uint64_t) statsd_bytes_received;
17 + PAD64(uint64_t) statsd_bytes_sent;
18 + PAD64(uint64_t) stream_bytes_received;
19 + PAD64(uint64_t) stream_bytes_sent;
20 +} live_stats = { 0 };
21 +
22 +void pulse_web_server_received_bytes(size_t bytes) {
23 + __atomic_add_fetch(&live_stats.api_bytes_received, bytes, __ATOMIC_RELAXED);
24 +}
25 +
26 +void pulse_web_server_sent_bytes(size_t bytes) {
27 + __atomic_add_fetch(&live_stats.api_bytes_sent, bytes, __ATOMIC_RELAXED);
28 +}
29 +
30 +void pulse_statsd_received_bytes(size_t bytes) {
31 + __atomic_add_fetch(&live_stats.statsd_bytes_received, bytes, __ATOMIC_RELAXED);
32 +}
33 +
34 +void pulse_statsd_sent_bytes(size_t bytes) {
35 + __atomic_add_fetch(&live_stats.statsd_bytes_sent, bytes, __ATOMIC_RELAXED);
36 +}
37 +
38 +void pulse_stream_received_bytes(size_t bytes) {
39 + __atomic_add_fetch(&live_stats.stream_bytes_received, bytes, __ATOMIC_RELAXED);
40 +}
41 +
42 +void pulse_stream_sent_bytes(size_t bytes) {
43 + __atomic_add_fetch(&live_stats.stream_bytes_sent, bytes, __ATOMIC_RELAXED);
44 +}
45 +
46 +static inline void pulse_network_copy(struct network_statistics *gs) {
47 + gs->api_bytes_received = __atomic_load_n(&live_stats.api_bytes_received, __ATOMIC_RELAXED);
48 + gs->api_bytes_sent = __atomic_load_n(&live_stats.api_bytes_sent, __ATOMIC_RELAXED);
49 +
50 + gs->statsd_bytes_received = __atomic_load_n(&live_stats.statsd_bytes_received, __ATOMIC_RELAXED);
51 + gs->statsd_bytes_sent = __atomic_load_n(&live_stats.statsd_bytes_sent, __ATOMIC_RELAXED);
52 +
53 + gs->stream_bytes_received = __atomic_load_n(&live_stats.stream_bytes_received, __ATOMIC_RELAXED);
54 + gs->stream_bytes_sent = __atomic_load_n(&live_stats.stream_bytes_sent, __ATOMIC_RELAXED);
55 +}
56 +
57 +void pulse_network_do(bool extended __maybe_unused) {
58 + static struct network_statistics gs;
59 + pulse_network_copy(&gs);
60 +
61 + if(gs.api_bytes_received || gs.api_bytes_sent) {
62 + static RRDSET *st_bytes = NULL;
63 + static RRDDIM *rd_in = NULL,
64 + *rd_out = NULL;
65 +
66 + if (unlikely(!st_bytes)) {
67 + st_bytes = rrdset_create_localhost(
68 + "netdata"
69 + , "network_api"
70 + , NULL
71 + , PULSE_NETWORK_CHART_FAMILY
72 + , PULSE_NETWORK_CHART_CONTEXT
73 + , PULSE_NETWORK_CHART_TITLE
74 + , PULSE_NETWORK_CHART_UNITS
75 + , "netdata"
76 + , "pulse"
77 + , PULSE_NETWORK_CHART_PRIORITY
78 + , localhost->rrd_update_every
79 + , RRDSET_TYPE_AREA
80 + );
81 +
82 + rrdlabels_add(st_bytes->rrdlabels, "endpoint", "web-server", RRDLABEL_SRC_AUTO);
83 +
84 + rd_in = rrddim_add(st_bytes, "in", NULL, 8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
85 + rd_out = rrddim_add(st_bytes, "out", NULL, -8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
86 + }
87 +
88 + rrddim_set_by_pointer(st_bytes, rd_in, (collected_number) gs.api_bytes_received);
89 + rrddim_set_by_pointer(st_bytes, rd_out, (collected_number) gs.api_bytes_sent);
90 + rrdset_done(st_bytes);
91 + }
92 +
93 + if(gs.statsd_bytes_received || gs.statsd_bytes_sent) {
94 + static RRDSET *st_bytes = NULL;
95 + static RRDDIM *rd_in = NULL,
96 + *rd_out = NULL;
97 +
98 + if (unlikely(!st_bytes)) {
99 + st_bytes = rrdset_create_localhost(
100 + "netdata"
101 + , "network_statsd"
102 + , NULL
103 + , PULSE_NETWORK_CHART_FAMILY
104 + , PULSE_NETWORK_CHART_CONTEXT
105 + , PULSE_NETWORK_CHART_TITLE
106 + , PULSE_NETWORK_CHART_UNITS
107 + , "netdata"
108 + , "pulse"
109 + , PULSE_NETWORK_CHART_PRIORITY
110 + , localhost->rrd_update_every
111 + , RRDSET_TYPE_AREA
112 + );
113 +
114 + rrdlabels_add(st_bytes->rrdlabels, "endpoint", "statsd", RRDLABEL_SRC_AUTO);
115 +
116 + rd_in = rrddim_add(st_bytes, "in", NULL, 8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
117 + rd_out = rrddim_add(st_bytes, "out", NULL, -8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
118 + }
119 +
120 + rrddim_set_by_pointer(st_bytes, rd_in, (collected_number) gs.statsd_bytes_received);
121 + rrddim_set_by_pointer(st_bytes, rd_out, (collected_number) gs.statsd_bytes_sent);
122 + rrdset_done(st_bytes);
123 + }
124 +
125 + if(gs.stream_bytes_received || gs.stream_bytes_sent) {
126 + static RRDSET *st_bytes = NULL;
127 + static RRDDIM *rd_in = NULL,
128 + *rd_out = NULL;
129 +
130 + if (unlikely(!st_bytes)) {
131 + st_bytes = rrdset_create_localhost(
132 + "netdata"
133 + , "network_streaming"
134 + , NULL
135 + , PULSE_NETWORK_CHART_FAMILY
136 + , PULSE_NETWORK_CHART_CONTEXT
137 + , PULSE_NETWORK_CHART_TITLE
138 + , PULSE_NETWORK_CHART_UNITS
139 + , "netdata"
140 + , "pulse"
141 + , PULSE_NETWORK_CHART_PRIORITY
142 + , localhost->rrd_update_every
143 + , RRDSET_TYPE_AREA
144 + );
145 +
146 + rrdlabels_add(st_bytes->rrdlabels, "endpoint", "streaming", RRDLABEL_SRC_AUTO);
147 +
148 + rd_in = rrddim_add(st_bytes, "in", NULL, 8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
149 + rd_out = rrddim_add(st_bytes, "out", NULL, -8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
150 + }
151 +
152 + rrddim_set_by_pointer(st_bytes, rd_in, (collected_number) gs.stream_bytes_received);
153 + rrddim_set_by_pointer(st_bytes, rd_out, (collected_number) gs.stream_bytes_sent);
154 + rrdset_done(st_bytes);
155 + }
156 +
157 + if(aclk_online()) {
158 + struct mqtt_wss_stats t = aclk_statistics();
159 + if (t.bytes_rx || t.bytes_tx) {
160 + static RRDSET *st_bytes = NULL;
161 + static RRDDIM *rd_in = NULL, *rd_out = NULL;
162 +
163 + if (unlikely(!st_bytes)) {
164 + st_bytes = rrdset_create_localhost(
165 + "netdata",
166 + "network_aclk",
167 + NULL,
168 + PULSE_NETWORK_CHART_FAMILY,
169 + PULSE_NETWORK_CHART_CONTEXT,
170 + PULSE_NETWORK_CHART_TITLE,
171 + PULSE_NETWORK_CHART_UNITS,
172 + "netdata",
173 + "pulse",
174 + PULSE_NETWORK_CHART_PRIORITY,
175 + localhost->rrd_update_every,
176 + RRDSET_TYPE_AREA);
177 +
178 + rrdlabels_add(st_bytes->rrdlabels, "endpoint", "aclk", RRDLABEL_SRC_AUTO);
179 +
180 + rd_in = rrddim_add(st_bytes, "in", NULL, 8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
181 + rd_out = rrddim_add(st_bytes, "out", NULL, -8, BITS_IN_A_KILOBIT, RRD_ALGORITHM_INCREMENTAL);
182 + }
183 +
184 + rrddim_set_by_pointer(st_bytes, rd_in, (collected_number)t.bytes_rx);
185 + rrddim_set_by_pointer(st_bytes, rd_out, (collected_number)t.bytes_tx);
186 + rrdset_done(st_bytes);
187 + }
188 + }
189 +}
src/daemon/pulse/pulse-network.h new
+19
@@ -0,0 +1,19 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +#ifndef NETDATA_PULSE_NETWORK_H
4 +#define NETDATA_PULSE_NETWORK_H
5 +
6 +#include "daemon/common.h"
7 +
8 +void pulse_network_do(bool extended);
9 +
10 +void pulse_web_server_received_bytes(size_t bytes);
11 +void pulse_web_server_sent_bytes(size_t bytes);
12 +
13 +void pulse_statsd_received_bytes(size_t bytes);
14 +void pulse_statsd_sent_bytes(size_t bytes);
15 +
16 +void pulse_stream_received_bytes(size_t bytes);
17 +void pulse_stream_sent_bytes(size_t bytes);
18 +
19 +#endif //NETDATA_PULSE_NETWORK_H
src/daemon/pulse/pulse.c
+6 -1
@@ -18,8 +18,9 @@
18 #define WORKER_JOB_MALLOC_TRACE 12
19 #define WORKER_JOB_REGISTRY 13
20 #define WORKER_JOB_ARAL 14
21 +#define WORKER_JOB_NETWORK 15
22
22 -#if WORKER_UTILIZATION_MAX_JOB_TYPES < 15
23 +#if WORKER_UTILIZATION_MAX_JOB_TYPES < 16
24 #error "WORKER_UTILIZATION_MAX_JOB_TYPES has to be at least 14"
25 #endif
26
@@ -44,6 +45,7 @@ static void pulse_register_workers(void) {
45 worker_register_job_name(WORKER_JOB_MALLOC_TRACE, "malloc_trace");
46 worker_register_job_name(WORKER_JOB_REGISTRY, "registry");
47 worker_register_job_name(WORKER_JOB_ARAL, "aral");
48 + worker_register_job_name(WORKER_JOB_NETWORK, "network");
49 }
50
51 static void pulse_cleanup(void *pptr)
@@ -99,6 +101,9 @@ void *pulse_thread_main(void *ptr) {
101 worker_is_busy(WORKER_JOB_QUERIES);
102 pulse_queries_do(pulse_extended_enabled);
103
104 + worker_is_busy(WORKER_JOB_NETWORK);
105 + pulse_network_do(pulse_extended_enabled);
106 +
107 worker_is_busy(WORKER_JOB_ML);
108 pulse_ml_do(pulse_extended_enabled);
109
src/daemon/pulse/pulse.h
+1
@@ -24,6 +24,7 @@ extern bool pulse_extended_enabled;
24 #include "pulse-workers.h"
25 #include "pulse-trace-allocations.h"
26 #include "pulse-aral.h"
27 +#include "pulse-network.h"
28
29 void *pulse_thread_main(void *ptr);
30 void *pulse_thread_sqlite3_main(void *ptr);
src/database/contexts/worker.c
+3 -3
@@ -407,7 +407,7 @@ static void rrdinstance_post_process_updates(RRDINSTANCE *ri, bool force, RRD_FL
407 if(dictionary_entries(ri->rrdmetrics) > 0) {
408 RRDMETRIC *rm;
409 dfe_start_read((DICTIONARY *)ri->rrdmetrics, rm) {
410 - if(unlikely(!service_running(SERVICE_CONTEXT))) break;
410 + if(unlikely(worker_jobs && !service_running(SERVICE_CONTEXT))) break;
411
412 RRD_FLAGS reason_to_pass = reason;
413 if(rrd_flag_check(ri, RRD_FLAG_UPDATE_REASON_UPDATE_RETENTION))
@@ -516,7 +516,7 @@ static void rrdcontext_post_process_updates(RRDCONTEXT *rc, bool force, RRD_FLAG
516 if(dictionary_entries(rc->rrdinstances) > 0) {
517 RRDINSTANCE *ri;
518 dfe_start_reentrant(rc->rrdinstances, ri) {
519 - if(unlikely(!service_running(SERVICE_CONTEXT))) break;
519 + if(unlikely(worker_jobs && !service_running(SERVICE_CONTEXT))) break;
520
521 RRD_FLAGS reason_to_pass = reason;
522 if(rrd_flag_check(rc, RRD_FLAG_UPDATE_REASON_UPDATE_RETENTION))
@@ -703,7 +703,7 @@ static void rrdcontext_dequeue_from_post_processing(RRDCONTEXT *rc) {
703
704 void rrdcontext_initial_processing_after_loading(RRDCONTEXT *rc) {
705 rrdcontext_dequeue_from_post_processing(rc);
706 - rrdcontext_post_process_updates(rc, false, RRD_FLAG_NONE, true);
706 + rrdcontext_post_process_updates(rc, false, RRD_FLAG_NONE, false);
707 }
708
709 void rrdcontext_delete_after_loading(RRDHOST *host, RRDCONTEXT *rc) {
src/database/sqlite/sqlite_aclk_node.c
+1 -1
@@ -217,7 +217,7 @@ void aclk_check_node_info_and_collectors(void)
217 context_pp_post = "')";
218 }
219
220 - nd_log_limit_static_thread_var(erl, 10, 100 * USEC_PER_MS);
220 + nd_log_limit_static_global_var(erl, 10, 100 * USEC_PER_MS);
221 nd_log_limit(&erl, NDLS_DAEMON, NDLP_INFO,
222 "NODES INFO: %zu nodes loading contexts%s%s%s, %zu receiving replication%s%s%s, %zu sending replication%s%s%s, %zu pending context post processing%s%s%s%s",
223 context_loading, context_loading_pre, context_loading_body, context_loading_post,
src/streaming/stream-receiver.c
+8 -5
@@ -167,6 +167,7 @@ static inline ssize_t receiver_read_uncompressed(struct receiver_state *r) {
167
168 r->thread.uncompressed.read_len += bytes;
169 r->thread.uncompressed.read_buffer[r->thread.uncompressed.read_len] = '\0';
170 + pulse_stream_received_bytes(bytes);
171 }
172
173 return bytes;
@@ -276,15 +277,16 @@ static inline ssize_t receiver_read_compressed(struct receiver_state *r) {
277 internal_fatal(r->thread.uncompressed.read_buffer[r->thread.uncompressed.read_len] != '\0',
278 "%s: read_buffer does not start with zero #2", __FUNCTION__ );
279
279 - ssize_t bytes_read = read_stream(r, r->thread.compressed.buf + r->thread.compressed.used,
280 + ssize_t bytes = read_stream(r, r->thread.compressed.buf + r->thread.compressed.used,
281 r->thread.compressed.size - r->thread.compressed.used);
282
282 - if(bytes_read > 0) {
283 - r->thread.compressed.used += bytes_read;
284 - worker_set_metric(WORKER_RECEIVER_JOB_BYTES_READ, (NETDATA_DOUBLE)bytes_read);
283 + if(bytes > 0) {
284 + r->thread.compressed.used += bytes;
285 + worker_set_metric(WORKER_RECEIVER_JOB_BYTES_READ, (NETDATA_DOUBLE)bytes);
286 + pulse_stream_received_bytes(bytes);
287 }
288
287 - return bytes_read;
289 + return bytes;
290 }
291
292 // --------------------------------------------------------------------------------------------------------------------
@@ -697,6 +699,7 @@ bool stream_receiver_send_data(struct stream_thread *sth, struct receiver_state
699
700 ssize_t rc = write_stream(rpt, chunk, outstanding);
701 if (likely(rc > 0)) {
702 + pulse_stream_sent_bytes(rc);
703 rpt->thread.last_traffic_ut = now_ut;
704 stream_circular_buffer_del_unsafe(scb, rc, now_ut);
705 if (!stats->bytes_outstanding) {
src/streaming/stream-sender.c
+2
@@ -622,6 +622,7 @@ bool stream_sender_send_data(struct stream_thread *sth, struct sender_state *s,
622
623 ssize_t rc = nd_sock_send_nowait(&s->sock, chunk, outstanding);
624 if (likely(rc > 0)) {
625 + pulse_stream_sent_bytes(rc);
626 stream_circular_buffer_del_unsafe(s->scb, rc, now_ut);
627 replication_sender_recalculate_buffer_used_ratio_unsafe(s);
628 s->thread.last_traffic_ut = now_ut;
@@ -712,6 +713,7 @@ bool stream_sender_receive_data(struct stream_thread *sth, struct sender_state *
713
714 s->thread.last_traffic_ut = now_ut;
715 sth->snd.bytes_received += rc;
716 + pulse_stream_received_bytes(rc);
717
718 worker_is_busy(WORKER_SENDER_JOB_EXECUTE);
719 stream_sender_execute_commands(s);
src/web/server/static/static-threaded.c
+7 -1
@@ -179,13 +179,17 @@ static int web_server_rcv_callback(POLLINFO *pi, nd_poll_event_t *events) {
179 bytes = web_client_receive(w);
180
181 if (likely(bytes > 0)) {
182 + pulse_web_server_received_bytes(bytes);
183 +
184 netdata_log_debug(D_WEB_CLIENT, "%llu: processing received data on fd %d.", w->id, fd);
185 worker_is_idle();
186 worker_is_busy(WORKER_JOB_PROCESS);
187 web_client_process_request_from_web_server(w);
188
189 if (unlikely(w->mode == HTTP_REQUEST_MODE_STREAM)) {
188 - web_client_send(w);
190 + ssize_t rc = web_client_send(w);
191 + if(rc > 0)
192 + pulse_web_server_sent_bytes(rc);
193 }
194 else if(unlikely(w->fd == fd && web_client_has_wait_receive(w)))
195 *events |= ND_POLL_READ;
@@ -229,6 +233,8 @@ static int web_server_snd_callback(POLLINFO *pi, nd_poll_event_t *events) {
233 goto cleanup;
234 }
235
236 + pulse_web_server_sent_bytes(ret);
237 +
238 if(unlikely(w->fd == fd && web_client_has_wait_receive(w)))
239 *events |= ND_POLL_READ;
240