@cryptotaxi247 / netdata-1 / commits / 2b224f36c

rename variables and functions in rrdpush

Costa Tsaousis (ktsaou) committed Sep 23, 2017 at 15:31 UTC 2b224f36c9fe8703b3321de1d8c2a8cfb5b950db
4 files changed +170 -169
src/rrd.h
+13 -13
@@ -424,18 +424,18 @@ struct rrdhost {
424 // ------------------------------------------------------------------------
425 // streaming of data to remote hosts - rrdpush
426
427 - int rrdpush_enabled:1; // 1 when this host sends metrics to another netdata
428 - volatile int rrdpush_connected:1; // 1 when the sender is ready to push metrics
429 - volatile int rrdpush_spawn:1; // 1 when the sender thread has been spawn
430 - volatile int rrdpush_error_shown:1; // 1 when we have logged a communication error
427 + int rrdpush_send_enabled:1; // 1 when this host sends metrics to another netdata
428 + volatile int rrdpush_send_connected:1; // 1 when the sender is ready to push metrics
429 + volatile int rrdpush_sender_spawn:1; // 1 when the sender thread has been spawn
430 + volatile int rrdpush_sender_error_shown:1; // 1 when we have logged a communication error
431 volatile int rrdpush_sender_join:1; // 1 when we have to join the sending thread
432 - char *rrdpush_destination; // where to send metrics to
433 - char *rrdpush_api_key; // the api key at the receiving netdata
434 - int rrdpush_socket; // the fd of the socket to the remote host, or -1
435 - pthread_t rrdpush_thread; // the sender thread
436 - netdata_mutex_t rrdpush_mutex; // exclusive access to rrdpush_buffer
437 - int rrdpush_pipe[2]; // collector to sender thread communication
438 - BUFFER *rrdpush_buffer; // collector fills it, sender sends them
432 + char *rrdpush_send_destination; // where to send metrics to
433 + char *rrdpush_send_api_key; // the api key at the receiving netdata
434 + int rrdpush_sender_socket; // the fd of the socket to the remote host, or -1
435 + pthread_t rrdpush_sender_thread; // the sender thread
436 + netdata_mutex_t rrdpush_sender_buffer_mutex; // exclusive access to rrdpush_sender_buffer
437 + int rrdpush_sender_pipe[2]; // collector to sender thread communication
438 + BUFFER *rrdpush_sender_buffer; // collector fills it, sender sends them
439
440
441 // ------------------------------------------------------------------------
@@ -561,8 +561,8 @@ extern void __rrd_check_wrlock(const char *file, const char *function, const uns
561 #else
562 #define rrdhost_check_rdlock(host) (void)0
563 #define rrdhost_check_wrlock(host) (void)0
564 -#define rrdset_check_rdlock(host) (void)0
565 -#define rrdset_check_wrlock(host) (void)0
564 +#define rrdset_check_rdlock(st) (void)0
565 +#define rrdset_check_wrlock(st) (void)0
566 #define rrd_check_rdlock() (void)0
567 #define rrd_check_wrlock() (void)0
568 #endif
src/rrdhost.c
+12 -12
@@ -120,15 +120,15 @@ RRDHOST *rrdhost_create(const char *hostname,
120 host->rrd_history_entries = align_entries_to_pagesize(memory_mode, entries);
121 host->rrd_memory_mode = memory_mode;
122 host->health_enabled = (memory_mode == RRD_MEMORY_MODE_NONE)? 0 : health_enabled;
123 - host->rrdpush_enabled = (rrdpush_enabled && rrdpush_destination && *rrdpush_destination && rrdpush_api_key && *rrdpush_api_key);
124 - host->rrdpush_destination = (host->rrdpush_enabled)?strdupz(rrdpush_destination):NULL;
125 - host->rrdpush_api_key = (host->rrdpush_enabled)?strdupz(rrdpush_api_key):NULL;
123 + host->rrdpush_send_enabled = (rrdpush_enabled && rrdpush_destination && *rrdpush_destination && rrdpush_api_key && *rrdpush_api_key);
124 + host->rrdpush_send_destination = (host->rrdpush_send_enabled)?strdupz(rrdpush_destination):NULL;
125 + host->rrdpush_send_api_key = (host->rrdpush_send_enabled)?strdupz(rrdpush_api_key):NULL;
126
127 - host->rrdpush_pipe[0] = -1;
128 - host->rrdpush_pipe[1] = -1;
129 - host->rrdpush_socket = -1;
127 + host->rrdpush_sender_pipe[0] = -1;
128 + host->rrdpush_sender_pipe[1] = -1;
129 + host->rrdpush_sender_socket = -1;
130
131 - netdata_mutex_init(&host->rrdpush_mutex);
131 + netdata_mutex_init(&host->rrdpush_sender_buffer_mutex);
132 netdata_rwlock_init(&host->rrdhost_rwlock);
133
134 rrdhost_init_hostname(host, hostname);
@@ -272,9 +272,9 @@ RRDHOST *rrdhost_create(const char *hostname,
272 , host->rrd_update_every
273 , rrd_memory_mode_name(host->rrd_memory_mode)
274 , host->rrd_history_entries
275 - , host->rrdpush_enabled?"enabled":"disabled"
276 - , host->rrdpush_destination?host->rrdpush_destination:""
277 - , host->rrdpush_api_key?host->rrdpush_api_key:""
275 + , host->rrdpush_send_enabled?"enabled":"disabled"
276 + , host->rrdpush_send_destination?host->rrdpush_send_destination:""
277 + , host->rrdpush_send_api_key?host->rrdpush_send_api_key:""
278 , host->health_enabled?"enabled":"disabled"
279 , host->cache_dir
280 , host->varlib_dir
@@ -508,8 +508,8 @@ void rrdhost_free(RRDHOST *host) {
508 freez((void *)host->os);
509 freez(host->cache_dir);
510 freez(host->varlib_dir);
511 - freez(host->rrdpush_api_key);
512 - freez(host->rrdpush_destination);
511 + freez(host->rrdpush_send_api_key);
512 + freez(host->rrdpush_send_destination);
513 freez(host->health_default_exec);
514 freez(host->health_default_recipient);
515 freez(host->health_log_filename);
src/rrdpush.c
+142 -141
@@ -59,27 +59,29 @@ int rrdpush_init() {
59 // this is for the first iterations of each chart
60 unsigned int remote_clock_resync_iterations = 60;
61
62 -#define rrdpush_lock(host) netdata_mutex_lock(&((host)->rrdpush_mutex))
63 -#define rrdpush_unlock(host) netdata_mutex_unlock(&((host)->rrdpush_mutex))
62 +#define rrdpush_buffer_lock(host) netdata_mutex_lock(&((host)->rrdpush_sender_buffer_mutex))
63 +#define rrdpush_buffer_unlock(host) netdata_mutex_unlock(&((host)->rrdpush_sender_buffer_mutex))
64
65 // checks if the current chart definition has been sent
66 static inline int need_to_send_chart_definition(RRDSET *st) {
67 + rrdset_check_rdlock(st);
68 +
69 if(unlikely(!(rrdset_flag_check(st, RRDSET_FLAG_EXPOSED_UPSTREAM))))
70 return 1;
71
72 RRDDIM *rd;
73 rrddim_foreach_read(rd, st)
72 - if(!rd->exposed)
74 + if(unlikely(!rd->exposed))
75 return 1;
76
77 return 0;
78 }
79
80 // sends the current chart definition
79 -static inline void send_chart_definition(RRDSET *st) {
81 +static inline void rrdpush_send_chart_definition_nolock(RRDSET *st) {
82 rrdset_flag_set(st, RRDSET_FLAG_EXPOSED_UPSTREAM);
83
82 - buffer_sprintf(st->rrdhost->rrdpush_buffer, "CHART \"%s\" \"%s\" \"%s\" \"%s\" \"%s\" \"%s\" \"%s\" %ld %d \"%s %s %s\"\n"
84 + buffer_sprintf(st->rrdhost->rrdpush_sender_buffer, "CHART \"%s\" \"%s\" \"%s\" \"%s\" \"%s\" \"%s\" \"%s\" %ld %d \"%s %s %s\"\n"
85 , st->id
86 , st->name
87 , st->title
@@ -96,7 +98,7 @@ static inline void send_chart_definition(RRDSET *st) {
98
99 RRDDIM *rd;
100 rrddim_foreach_read(rd, st) {
99 - buffer_sprintf(st->rrdhost->rrdpush_buffer, "DIMENSION \"%s\" \"%s\" \"%s\" " COLLECTED_NUMBER_FORMAT " " COLLECTED_NUMBER_FORMAT " \"%s %s\"\n"
101 + buffer_sprintf(st->rrdhost->rrdpush_sender_buffer, "DIMENSION \"%s\" \"%s\" \"%s\" " COLLECTED_NUMBER_FORMAT " " COLLECTED_NUMBER_FORMAT " \"%s %s\"\n"
102 , rd->id
103 , rd->name
104 , rrd_algorithm_name(rd->algorithm)
@@ -112,19 +114,20 @@ static inline void send_chart_definition(RRDSET *st) {
114 }
115
116 // sends the current chart dimensions
115 -static inline void send_chart_metrics(RRDSET *st) {
116 - buffer_sprintf(st->rrdhost->rrdpush_buffer, "BEGIN \"%s\" %llu\n", st->id, (st->upstream_resync_time > st->last_collected_time.tv_sec)?st->usec_since_last_update:0);
117 +static inline void rrdpush_send_chart_metrics_nolock(RRDSET *st) {
118 + buffer_sprintf(st->rrdhost->rrdpush_sender_buffer, "BEGIN \"%s\" %llu\n", st->id, (st->upstream_resync_time > st->last_collected_time.tv_sec)?st->usec_since_last_update:0);
119
120 RRDDIM *rd;
121 rrddim_foreach_read(rd, st) {
122 if(rd->updated && rd->exposed)
121 - buffer_sprintf(st->rrdhost->rrdpush_buffer, "SET \"%s\" = " COLLECTED_NUMBER_FORMAT "\n"
122 - , rd->id
123 - , rd->collected_value
123 + buffer_sprintf(st->rrdhost->rrdpush_sender_buffer
124 + , "SET \"%s\" = " COLLECTED_NUMBER_FORMAT "\n"
125 + , rd->id
126 + , rd->collected_value
127 );
128 }
129
127 - buffer_strcat(st->rrdhost->rrdpush_buffer, "END\n");
130 + buffer_strcat(st->rrdhost->rrdpush_sender_buffer, "END\n");
131 }
132
133 static void rrdpush_sender_thread_spawn(RRDHOST *host);
@@ -133,9 +136,9 @@ void rrdset_push_chart_definition(RRDSET *st) {
136 RRDHOST *host = st->rrdhost;
137
138 rrdset_rdlock(st);
136 - rrdpush_lock(host);
137 - send_chart_definition(st);
138 - rrdpush_unlock(host);
139 + rrdpush_buffer_lock(host);
140 + rrdpush_send_chart_definition_nolock(st);
141 + rrdpush_buffer_unlock(host);
142 rrdset_unlock(st);
143 }
144
@@ -145,35 +148,35 @@ void rrdset_done_push(RRDSET *st) {
148 if(unlikely(!rrdset_flag_check(st, RRDSET_FLAG_ENABLED)))
149 return;
150
148 - rrdpush_lock(host);
151 + rrdpush_buffer_lock(host);
152
150 - if(unlikely(host->rrdpush_enabled && !host->rrdpush_spawn))
153 + if(unlikely(host->rrdpush_send_enabled && !host->rrdpush_sender_spawn))
154 rrdpush_sender_thread_spawn(host);
155
153 - if(unlikely(!host->rrdpush_buffer || !host->rrdpush_connected)) {
154 - if(unlikely(!host->rrdpush_error_shown))
156 + if(unlikely(!host->rrdpush_sender_buffer || !host->rrdpush_send_connected)) {
157 + if(unlikely(!host->rrdpush_sender_error_shown))
158 error("STREAM %s [send]: not ready - discarding collected metrics.", host->hostname);
159
157 - host->rrdpush_error_shown = 1;
160 + host->rrdpush_sender_error_shown = 1;
161
159 - rrdpush_unlock(host);
162 + rrdpush_buffer_unlock(host);
163 return;
164 }
162 - else if(unlikely(host->rrdpush_error_shown)) {
165 + else if(unlikely(host->rrdpush_sender_error_shown)) {
166 info("STREAM %s [send]: sending metrics...", host->hostname);
164 - host->rrdpush_error_shown = 0;
167 + host->rrdpush_sender_error_shown = 0;
168 }
169
170 if(need_to_send_chart_definition(st))
168 - send_chart_definition(st);
171 + rrdpush_send_chart_definition_nolock(st);
172
170 - send_chart_metrics(st);
173 + rrdpush_send_chart_metrics_nolock(st);
174
175 // signal the sender there are more data
173 - if(host->rrdpush_pipe[PIPE_WRITE] != -1 && write(host->rrdpush_pipe[PIPE_WRITE], " ", 1) == -1)
176 + if(host->rrdpush_sender_pipe[PIPE_WRITE] != -1 && write(host->rrdpush_sender_pipe[PIPE_WRITE], " ", 1) == -1)
177 error("STREAM %s [send]: cannot write to internal pipe", host->hostname);
178
176 - rrdpush_unlock(host);
179 + rrdpush_buffer_unlock(host);
180 }
181
182 // ----------------------------------------------------------------------------
@@ -202,57 +205,25 @@ static void rrdpush_sender_thread_reset_all_charts(RRDHOST *host) {
205 }
206
207 static inline void rrdpush_sender_thread_data_flush(RRDHOST *host) {
205 - rrdpush_lock(host);
208 + rrdpush_buffer_lock(host);
209
207 - if(buffer_strlen(host->rrdpush_buffer))
208 - error("STREAM %s [send]: discarding %zu bytes of metrics already in the buffer.", host->hostname, buffer_strlen(host->rrdpush_buffer));
210 + if(buffer_strlen(host->rrdpush_sender_buffer))
211 + error("STREAM %s [send]: discarding %zu bytes of metrics already in the buffer.", host->hostname, buffer_strlen(host->rrdpush_sender_buffer));
212
210 - buffer_flush(host->rrdpush_buffer);
213 + buffer_flush(host->rrdpush_sender_buffer);
214
215 rrdpush_sender_thread_reset_all_charts(host);
216
214 - rrdpush_unlock(host);
215 -}
216 -
217 -static void rrdpush_sender_thread_cleanup_locked_all(RRDHOST *host) {
218 - host->rrdpush_connected = 0;
219 -
220 - if(host->rrdpush_socket != -1) {
221 - close(host->rrdpush_socket);
222 - host->rrdpush_socket = -1;
223 - }
224 -
225 - // close the pipe
226 - if(host->rrdpush_pipe[PIPE_READ] != -1) {
227 - close(host->rrdpush_pipe[PIPE_READ]);
228 - host->rrdpush_pipe[PIPE_READ] = -1;
229 - }
230 -
231 - if(host->rrdpush_pipe[PIPE_WRITE] != -1) {
232 - close(host->rrdpush_pipe[PIPE_WRITE]);
233 - host->rrdpush_pipe[PIPE_WRITE] = -1;
234 - }
235 -
236 - buffer_free(host->rrdpush_buffer);
237 - host->rrdpush_buffer = NULL;
238 -
239 - if(!host->rrdpush_sender_join) {
240 - info("STREAM %s [send]: sending thread detaches itself.", host->hostname);
241 - pthread_detach(pthread_self());
242 - }
243 -
244 - host->rrdpush_spawn = 0;
245 -
246 - info("STREAM %s [send]: sending thread now exits.", host->hostname);
217 + rrdpush_buffer_unlock(host);
218 }
219
220 void rrdpush_sender_thread_stop(RRDHOST *host) {
250 - rrdpush_lock(host);
221 + rrdpush_buffer_lock(host);
222 rrdhost_wrlock(host);
223
224 pthread_t thr = 0;
225
255 - if(host->rrdpush_spawn) {
226 + if(host->rrdpush_sender_spawn) {
227 info("STREAM %s [send]: signaling sending thread to stop...", host->hostname);
228
229 // signal the thread that we want to join it
@@ -260,16 +231,16 @@ void rrdpush_sender_thread_stop(RRDHOST *host) {
231
232 // copy the thread id, so that we will be waiting for the right one
233 // even if a new one has been spawn
263 - thr = host->rrdpush_thread;
234 + thr = host->rrdpush_sender_thread;
235
236 // signal it to cancel
266 - int ret = pthread_cancel(host->rrdpush_thread);
237 + int ret = pthread_cancel(host->rrdpush_sender_thread);
238 if(ret != 0)
239 error("STREAM %s [send]: pthread_cancel() returned error.", host->hostname);
240 }
241
242 rrdhost_unlock(host);
272 - rrdpush_unlock(host);
243 + rrdpush_buffer_unlock(host);
244
245 if(thr != 0) {
246 info("STREAM %s [send]: waiting for the sending thread to stop...", host->hostname);
@@ -286,18 +257,48 @@ void rrdpush_sender_thread_stop(RRDHOST *host) {
257 static void rrdpush_sender_thread_cleanup_callback(void *ptr) {
258 RRDHOST *host = (RRDHOST *)ptr;
259
289 - rrdpush_lock(host);
260 + rrdpush_buffer_lock(host);
261 rrdhost_wrlock(host);
262 +
263 info("STREAM %s [send]: sending thread cleans up...", host->hostname);
292 - rrdpush_sender_thread_cleanup_locked_all(host);
264 + host->rrdpush_send_connected = 0;
265 +
266 + if(host->rrdpush_sender_socket != -1) {
267 + close(host->rrdpush_sender_socket);
268 + host->rrdpush_sender_socket = -1;
269 + }
270 +
271 + // close the pipe
272 + if(host->rrdpush_sender_pipe[PIPE_READ] != -1) {
273 + close(host->rrdpush_sender_pipe[PIPE_READ]);
274 + host->rrdpush_sender_pipe[PIPE_READ] = -1;
275 + }
276 +
277 + if(host->rrdpush_sender_pipe[PIPE_WRITE] != -1) {
278 + close(host->rrdpush_sender_pipe[PIPE_WRITE]);
279 + host->rrdpush_sender_pipe[PIPE_WRITE] = -1;
280 + }
281 +
282 + buffer_free(host->rrdpush_sender_buffer);
283 + host->rrdpush_sender_buffer = NULL;
284 +
285 + if(!host->rrdpush_sender_join) {
286 + info("STREAM %s [send]: sending thread detaches itself.", host->hostname);
287 + pthread_detach(pthread_self());
288 + }
289 +
290 + host->rrdpush_sender_spawn = 0;
291 +
292 + info("STREAM %s [send]: sending thread now exits.", host->hostname);
293 +
294 rrdhost_unlock(host);
294 - rrdpush_unlock(host);
295 + rrdpush_buffer_unlock(host);
296 }
297
298 void *rrdpush_sender_thread(void *ptr) {
299 RRDHOST *host = (RRDHOST *)ptr;
300
300 - if(!host->rrdpush_enabled || !host->rrdpush_destination || !*host->rrdpush_destination || !host->rrdpush_api_key || !*host->rrdpush_api_key) {
301 + if(!host->rrdpush_send_enabled || !host->rrdpush_send_destination || !*host->rrdpush_send_destination || !host->rrdpush_send_api_key || !*host->rrdpush_send_api_key) {
302 error("STREAM %s [send]: thread created (task id %d), but host has streaming disabled.", host->hostname, gettid());
303 pthread_exit(NULL);
304 return NULL;
@@ -319,9 +320,9 @@ void *rrdpush_sender_thread(void *ptr) {
320 char connected_to[CONNECTED_TO_SIZE + 1] = "";
321
322 // initialize rrdpush globals
322 - host->rrdpush_buffer = buffer_create(1);
323 - host->rrdpush_connected = 0;
324 - if(pipe(host->rrdpush_pipe) == -1) fatal("STREAM %s [send]: cannot create required pipe.", host->hostname);
323 + host->rrdpush_sender_buffer = buffer_create(1);
324 + host->rrdpush_send_connected = 0;
325 + if(pipe(host->rrdpush_sender_pipe) == -1) fatal("STREAM %s [send]: cannot create required pipe.", host->hostname);
326
327 // initialize local variables
328 size_t begin = 0;
@@ -343,29 +344,29 @@ void *rrdpush_sender_thread(void *ptr) {
344
345 pthread_cleanup_push(rrdpush_sender_thread_cleanup_callback, host);
346
346 - for(; host->rrdpush_enabled && !netdata_exit ;) {
347 + for(; host->rrdpush_send_enabled && !netdata_exit ;) {
348 // check for outstanding cancellation requests
349 pthread_testcancel();
350
351 debug(D_STREAM, "STREAM: Checking if we need to timeout the connection...");
351 - if(host->rrdpush_socket != -1 && now_monotonic_sec() - last_sent_t > timeout) {
352 + if(host->rrdpush_sender_socket != -1 && now_monotonic_sec() - last_sent_t > timeout) {
353 error("STREAM %s [send to %s]: could not send metrics for %d seconds - closing connection - we have sent %zu bytes on this connection.", host->hostname, connected_to, timeout, sent_connection);
353 - close(host->rrdpush_socket);
354 - host->rrdpush_socket = -1;
354 + close(host->rrdpush_sender_socket);
355 + host->rrdpush_sender_socket = -1;
356 }
357
357 - if(unlikely(host->rrdpush_socket == -1)) {
358 + if(unlikely(host->rrdpush_sender_socket == -1)) {
359 debug(D_STREAM, "STREAM: Attempting to connect...");
360
360 - // stop appending data into rrdpush_buffer
361 + // stop appending data into rrdpush_sender_buffer
362 // they will be lost, so there is no point to do it
362 - host->rrdpush_connected = 0;
363 + host->rrdpush_send_connected = 0;
364
364 - info("STREAM %s [send to %s]: connecting...", host->hostname, host->rrdpush_destination);
365 - host->rrdpush_socket = connect_to_one_of(host->rrdpush_destination, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
365 + info("STREAM %s [send to %s]: connecting...", host->hostname, host->rrdpush_send_destination);
366 + host->rrdpush_sender_socket = connect_to_one_of(host->rrdpush_send_destination, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
367
367 - if(unlikely(host->rrdpush_socket == -1)) {
368 - error("STREAM %s [send to %s]: failed to connect", host->hostname, host->rrdpush_destination);
368 + if(unlikely(host->rrdpush_sender_socket == -1)) {
369 + error("STREAM %s [send to %s]: failed to connect", host->hostname, host->rrdpush_send_destination);
370 sleep(reconnect_delay);
371 continue;
372 }
@@ -378,7 +379,7 @@ void *rrdpush_sender_thread(void *ptr) {
379 "STREAM key=%s&hostname=%s&registry_hostname=%s&machine_guid=%s&update_every=%d&os=%s&tags=%s HTTP/1.1\r\n"
380 "User-Agent: netdata-push-service/%s\r\n"
381 "Accept: */*\r\n\r\n"
381 - , host->rrdpush_api_key
382 + , host->rrdpush_send_api_key
383 , host->hostname
384 , host->registry_hostname
385 , host->machine_guid
@@ -388,9 +389,9 @@ void *rrdpush_sender_thread(void *ptr) {
389 , program_version
390 );
391
391 - if(send_timeout(host->rrdpush_socket, http, strlen(http), 0, timeout) == -1) {
392 - close(host->rrdpush_socket);
393 - host->rrdpush_socket = -1;
392 + if(send_timeout(host->rrdpush_sender_socket, http, strlen(http), 0, timeout) == -1) {
393 + close(host->rrdpush_sender_socket);
394 + host->rrdpush_sender_socket = -1;
395 error("STREAM %s [send to %s]: failed to send http header to netdata", host->hostname, connected_to);
396 sleep(reconnect_delay);
397 continue;
@@ -398,17 +399,17 @@ void *rrdpush_sender_thread(void *ptr) {
399
400 info("STREAM %s [send to %s]: waiting response from remote netdata...", host->hostname, connected_to);
401
401 - if(recv_timeout(host->rrdpush_socket, http, HTTP_HEADER_SIZE, 0, timeout) == -1) {
402 - close(host->rrdpush_socket);
403 - host->rrdpush_socket = -1;
402 + if(recv_timeout(host->rrdpush_sender_socket, http, HTTP_HEADER_SIZE, 0, timeout) == -1) {
403 + close(host->rrdpush_sender_socket);
404 + host->rrdpush_sender_socket = -1;
405 error("STREAM %s [send to %s]: failed to initialize communication", host->hostname, connected_to);
406 sleep(reconnect_delay);
407 continue;
408 }
409
410 if(strncmp(http, START_STREAMING_PROMPT, strlen(START_STREAMING_PROMPT))) {
410 - close(host->rrdpush_socket);
411 - host->rrdpush_socket = -1;
411 + close(host->rrdpush_sender_socket);
412 + host->rrdpush_sender_socket = -1;
413 error("STREAM %s [send to %s]: server is not replying properly.", host->hostname, connected_to);
414 sleep(reconnect_delay);
415 continue;
@@ -417,28 +418,28 @@ void *rrdpush_sender_thread(void *ptr) {
418 info("STREAM %s [send to %s]: established communication - ready to send metrics...", host->hostname, connected_to);
419 last_sent_t = now_monotonic_sec();
420
420 - if(sock_setnonblock(host->rrdpush_socket) < 0)
421 + if(sock_setnonblock(host->rrdpush_sender_socket) < 0)
422 error("STREAM %s [send to %s]: cannot set non-blocking mode for socket.", host->hostname, connected_to);
423
423 - if(sock_enlarge_out(host->rrdpush_socket) < 0)
424 + if(sock_enlarge_out(host->rrdpush_sender_socket) < 0)
425 error("STREAM %s [send to %s]: cannot enlarge the socket buffer.", host->hostname, connected_to);
426
427 rrdpush_sender_thread_data_flush(host);
428 sent_connection = 0;
429
429 - // allow appending data into rrdpush_buffer
430 - host->rrdpush_connected = 1;
430 + // allow appending data into rrdpush_sender_buffer
431 + host->rrdpush_send_connected = 1;
432
432 - debug(D_STREAM, "STREAM: Connected on fd %d...", host->rrdpush_socket);
433 + debug(D_STREAM, "STREAM: Connected on fd %d...", host->rrdpush_sender_socket);
434 }
435
435 - ifd->fd = host->rrdpush_pipe[PIPE_READ];
436 + ifd->fd = host->rrdpush_sender_pipe[PIPE_READ];
437 ifd->events = POLLIN;
438 ifd->revents = 0;
439
439 - ofd->fd = host->rrdpush_socket;
440 + ofd->fd = host->rrdpush_sender_socket;
441 ofd->revents = 0;
441 - if(ofd->fd != -1 && begin < buffer_strlen(host->rrdpush_buffer)) {
442 + if(ofd->fd != -1 && begin < buffer_strlen(host->rrdpush_sender_buffer)) {
443 debug(D_STREAM, "STREAM: Requesting data output on streaming socket %d...", ofd->fd);
444 ofd->events = POLLOUT;
445 fdmax = 2;
@@ -449,13 +450,13 @@ void *rrdpush_sender_thread(void *ptr) {
450 fdmax = 1;
451 }
452
452 - debug(D_STREAM, "STREAM: Waiting for poll() events (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
453 + debug(D_STREAM, "STREAM: Waiting for poll() events (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_sender_buffer));
454 if(netdata_exit) break;
455 int retval = poll(fds, fdmax, 1000);
456 if(netdata_exit) break;
457
458 if(unlikely(retval == -1)) {
458 - debug(D_STREAM, "STREAM: poll() failed (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
459 + debug(D_STREAM, "STREAM: poll() failed (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_sender_buffer));
460
461 if(errno == EAGAIN || errno == EINTR) {
462 debug(D_STREAM, "STREAM: poll() failed with EAGAIN or EINTR...");
@@ -463,22 +464,22 @@ void *rrdpush_sender_thread(void *ptr) {
464 }
465
466 error("STREAM %s [send to %s]: failed to poll().", host->hostname, connected_to);
466 - close(host->rrdpush_socket);
467 - host->rrdpush_socket = -1;
467 + close(host->rrdpush_sender_socket);
468 + host->rrdpush_sender_socket = -1;
469 break;
470 }
471 else if(likely(retval)) {
472 if (ifd->revents & POLLIN || ifd->revents & POLLPRI) {
472 - debug(D_STREAM, "STREAM: Data added to send buffer (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
473 + debug(D_STREAM, "STREAM: Data added to send buffer (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_sender_buffer));
474
475 char buffer[1000 + 1];
475 - if (read(host->rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
476 + if (read(host->rrdpush_sender_pipe[PIPE_READ], buffer, 1000) == -1)
477 error("STREAM %s [send to %s]: cannot read from internal pipe.", host->hostname, connected_to);
478 }
479
480 if (ofd->revents & POLLOUT) {
480 - if (begin < buffer_strlen(host->rrdpush_buffer)) {
481 - debug(D_STREAM, "STREAM: Sending data (current buffer length %zu bytes, begin = %zu)...", buffer_strlen(host->rrdpush_buffer), begin);
481 + if (begin < buffer_strlen(host->rrdpush_sender_buffer)) {
482 + debug(D_STREAM, "STREAM: Sending data (current buffer length %zu bytes, begin = %zu)...", buffer_strlen(host->rrdpush_sender_buffer), begin);
483
484 // BEGIN RRDPUSH LOCKED SESSION
485
@@ -491,16 +492,16 @@ void *rrdpush_sender_thread(void *ptr) {
492 error("STREAM %s [send]: cannot set pthread cancel state to DISABLE.", host->hostname);
493
494 debug(D_STREAM, "STREAM: Getting exclusive lock on host...");
494 - rrdpush_lock(host);
495 + rrdpush_buffer_lock(host);
496
496 - debug(D_STREAM, "STREAM: Sending data, starting from %zu, size %zu...", begin, buffer_strlen(host->rrdpush_buffer));
497 - ssize_t ret = send(host->rrdpush_socket, &host->rrdpush_buffer->buffer[begin], buffer_strlen(host->rrdpush_buffer) - begin, MSG_DONTWAIT);
497 + debug(D_STREAM, "STREAM: Sending data, starting from %zu, size %zu...", begin, buffer_strlen(host->rrdpush_sender_buffer));
498 + ssize_t ret = send(host->rrdpush_sender_socket, &host->rrdpush_sender_buffer->buffer[begin], buffer_strlen(host->rrdpush_sender_buffer) - begin, MSG_DONTWAIT);
499 if (unlikely(ret == -1)) {
500 if (errno != EAGAIN && errno != EINTR && errno != EWOULDBLOCK) {
501 debug(D_STREAM, "STREAM: Send failed - closing socket...");
502 error("STREAM %s [send to %s]: failed to send metrics - closing connection - we have sent %zu bytes on this connection.", host->hostname, connected_to, sent_connection);
502 - close(host->rrdpush_socket);
503 - host->rrdpush_socket = -1;
503 + close(host->rrdpush_sender_socket);
504 + host->rrdpush_sender_socket = -1;
505 }
506 else {
507 debug(D_STREAM, "STREAM: Send failed - will retry...");
@@ -508,20 +509,20 @@ void *rrdpush_sender_thread(void *ptr) {
509 }
510 else if (likely(ret > 0)) {
511 // DEBUG - dump the scring to see it
511 - //char c = host->rrdpush_buffer->buffer[begin + ret];
512 - //host->rrdpush_buffer->buffer[begin + ret] = '\0';
513 - //debug(D_STREAM, "STREAM: sent from %zu to %zd:\n%s\n", begin, ret, &host->rrdpush_buffer->buffer[begin]);
514 - //host->rrdpush_buffer->buffer[begin + ret] = c;
512 + //char c = host->rrdpush_sender_buffer->buffer[begin + ret];
513 + //host->rrdpush_sender_buffer->buffer[begin + ret] = '\0';
514 + //debug(D_STREAM, "STREAM: sent from %zu to %zd:\n%s\n", begin, ret, &host->rrdpush_sender_buffer->buffer[begin]);
515 + //host->rrdpush_sender_buffer->buffer[begin + ret] = c;
516
517 sent_connection += ret;
518 sent_bytes += ret;
519 begin += ret;
520
520 - if (begin == buffer_strlen(host->rrdpush_buffer)) {
521 + if (begin == buffer_strlen(host->rrdpush_sender_buffer)) {
522 // we send it all
523
524 debug(D_STREAM, "STREAM: Sent %zd bytes (the whole buffer)...", ret);
524 - buffer_flush(host->rrdpush_buffer);
525 + buffer_flush(host->rrdpush_sender_buffer);
526 begin = 0;
527 }
528 else {
@@ -534,12 +535,12 @@ void *rrdpush_sender_thread(void *ptr) {
535 debug(D_STREAM, "STREAM: send() returned %zd - closing the socket...", ret);
536 error("STREAM %s [send to %s]: failed to send metrics (send() returned %zd) - closing connection - we have sent %zu bytes on this connection.",
537 host->hostname, connected_to, ret, sent_connection);
537 - close(host->rrdpush_socket);
538 - host->rrdpush_socket = -1;
538 + close(host->rrdpush_sender_socket);
539 + host->rrdpush_sender_socket = -1;
540 }
541
542 debug(D_STREAM, "STREAM: Releasing exclusive lock on host...");
542 - rrdpush_unlock(host);
543 + rrdpush_buffer_unlock(host);
544
545 if (pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
546 error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
@@ -554,20 +555,20 @@ void *rrdpush_sender_thread(void *ptr) {
555 if(unlikely(ofd->revents & POLLERR)) {
556 debug(D_STREAM, "STREAM: Send failed (POLLERR) - closing socket...");
557 error("STREAM %s [send to %s]: connection reports errors (POLLERR), closing it - we have sent %zu bytes on this connection.", host->hostname, connected_to, sent_connection);
557 - close(host->rrdpush_socket);
558 - host->rrdpush_socket = -1;
558 + close(host->rrdpush_sender_socket);
559 + host->rrdpush_sender_socket = -1;
560 }
561 else if(unlikely(ofd->revents & POLLHUP)) {
562 debug(D_STREAM, "STREAM: Send failed (POLLHUP) - closing socket...");
563 error("STREAM %s [send to %s]: connection closed by remote end (POLLHUP) - we have sent %zu bytes on this connection.", host->hostname, connected_to, sent_connection);
563 - close(host->rrdpush_socket);
564 - host->rrdpush_socket = -1;
564 + close(host->rrdpush_sender_socket);
565 + host->rrdpush_sender_socket = -1;
566 }
567 else if(unlikely(ofd->revents & POLLNVAL)) {
568 debug(D_STREAM, "STREAM: Send failed (POLLNVAL) - closing socket...");
569 error("STREAM %s [send to %s]: connection is invalid (POLLNVAL), closing it - we have sent %zu bytes on this connection.", host->hostname, connected_to, sent_connection);
569 - close(host->rrdpush_socket);
570 - host->rrdpush_socket = -1;
570 + close(host->rrdpush_sender_socket);
571 + host->rrdpush_sender_socket = -1;
572 }
573 }
574 else {
@@ -575,13 +576,13 @@ void *rrdpush_sender_thread(void *ptr) {
576 }
577
578 // protection from overflow
578 - if(buffer_strlen(host->rrdpush_buffer) > max_size) {
579 - debug(D_STREAM, "STREAM: Buffer is too big (%zu bytes), bigger than the max (%zu) - flushing it...", buffer_strlen(host->rrdpush_buffer), max_size);
579 + if(buffer_strlen(host->rrdpush_sender_buffer) > max_size) {
580 + debug(D_STREAM, "STREAM: Buffer is too big (%zu bytes), bigger than the max (%zu) - flushing it...", buffer_strlen(host->rrdpush_sender_buffer), max_size);
581 errno = 0;
581 - error("STREAM %s [send to %s]: too many data pending - buffer is %zu bytes long, %zu unsent - we have sent %zu bytes in total, %zu on this connection. Closing connection to flush the data.", host->hostname, connected_to, host->rrdpush_buffer->len, host->rrdpush_buffer->len - begin, sent_bytes, sent_connection);
582 - if(host->rrdpush_socket != -1) {
583 - close(host->rrdpush_socket);
584 - host->rrdpush_socket = -1;
582 + error("STREAM %s [send to %s]: too many data pending - buffer is %zu bytes long, %zu unsent - we have sent %zu bytes in total, %zu on this connection. Closing connection to flush the data.", host->hostname, connected_to, host->rrdpush_sender_buffer->len, host->rrdpush_sender_buffer->len - begin, sent_bytes, sent_connection);
583 + if(host->rrdpush_sender_socket != -1) {
584 + close(host->rrdpush_sender_socket);
585 + host->rrdpush_sender_socket = -1;
586 }
587 }
588 }
@@ -806,11 +807,11 @@ static void *rrdpush_receiver_thread(void *ptr) {
807 static void rrdpush_sender_thread_spawn(RRDHOST *host) {
808 rrdhost_wrlock(host);
809
809 - if(!host->rrdpush_spawn) {
810 - if(pthread_create(&host->rrdpush_thread, NULL, rrdpush_sender_thread, (void *) host))
810 + if(!host->rrdpush_sender_spawn) {
811 + if(pthread_create(&host->rrdpush_sender_thread, NULL, rrdpush_sender_thread, (void *) host))
812 error("STREAM %s [send]: failed to create new thread for client.", host->hostname);
813 else
813 - host->rrdpush_spawn = 1;
814 + host->rrdpush_sender_spawn = 1;
815 }
816
817 rrdhost_unlock(host);
src/rrdset.c
+3 -3
@@ -175,7 +175,7 @@ inline void rrdset_is_obsolete(RRDSET *st) {
175
176 // the chart will not get more updates (data collection)
177 // so, we have to push its definition now
178 - if(unlikely(st->rrdhost->rrdpush_enabled))
178 + if(unlikely(st->rrdhost->rrdpush_send_enabled))
179 rrdset_push_chart_definition(st);
180 }
181 }
@@ -1029,7 +1029,7 @@ void rrdset_done(RRDSET *st) {
1029 if(unlikely(netdata_exit)) return;
1030
1031 if(unlikely(st->rrd_memory_mode == RRD_MEMORY_MODE_NONE)) {
1032 - if(unlikely(st->rrdhost->rrdpush_enabled))
1032 + if(unlikely(st->rrdhost->rrdpush_send_enabled))
1033 rrdset_done_push_exclusive(st);
1034
1035 return;
@@ -1154,7 +1154,7 @@ void rrdset_done(RRDSET *st) {
1154 }
1155 st->counter_done++;
1156
1157 - if(unlikely(st->rrdhost->rrdpush_enabled))
1157 + if(unlikely(st->rrdhost->rrdpush_send_enabled))
1158 rrdset_done_push(st);
1159
1160 #ifdef NETDATA_INTERNAL_CHECKS