@cryptotaxi247 / netdata-1 / commits / b2267e4de

disconnect the sending thread only when there are not any slaves pushing metrics

Costa Tsaousis (ktsaou) committed Sep 23, 2017 at 03:13 UTC b2267e4decc7301c098c4901055e40b97bc8e11c
1 file changed +15 -9
src/rrdpush.c
+15 -9
@@ -160,7 +160,7 @@ void rrdset_done_push(RRDSET *st) {
160 return;
161 }
162 else if(unlikely(host->rrdpush_error_shown)) {
163 - info("STREAM %s [send]: ready - sending metrics...", host->hostname);
163 + info("STREAM %s [send]: sending metrics...", host->hostname);
164 host->rrdpush_error_shown = 0;
165 }
166
@@ -236,12 +236,14 @@ static void rrdpush_sender_thread_cleanup_locked_all(RRDHOST *host) {
236 buffer_free(host->rrdpush_buffer);
237 host->rrdpush_buffer = NULL;
238
239 - if(!host->rrdpush_sender_join)
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
244 - pthread_exit(NULL);
246 + info("STREAM %s [send]: sending thread now exits.", host->hostname);
247 }
248
249 void rrdpush_sender_thread_stop(RRDHOST *host) {
@@ -251,7 +253,7 @@ void rrdpush_sender_thread_stop(RRDHOST *host) {
253 pthread_t thr = 0;
254
255 if(host->rrdpush_spawn) {
254 - info("STREAM %s [send]: stopping sending thread...", host->hostname);
256 + info("STREAM %s [send]: signaling sending thread to stop...", host->hostname);
257
258 // signal the thread that we want to join it
259 host->rrdpush_sender_join = 1;
@@ -270,11 +272,14 @@ void rrdpush_sender_thread_stop(RRDHOST *host) {
272 rrdpush_unlock(host);
273
274 if(thr != 0) {
273 - info("STREAM %s [send]: waiting for sending thread to stop...", host->hostname);
275 + info("STREAM %s [send]: waiting for the sending thread to stop...", host->hostname);
276 +
277 void *result;
278 int ret = pthread_join(thr, &result);
279 if(ret != 0)
280 error("STREAM %s [send]: pthread_join() returned error.", host->hostname);
281 +
282 + info("STREAM %s [send]: sending thread has exited.", host->hostname);
283 }
284 }
285
@@ -283,7 +288,7 @@ static void rrdpush_sender_thread_cleanup_callback(void *ptr) {
288
289 rrdpush_lock(host);
290 rrdhost_wrlock(host);
286 - info("STREAM %s [send]: sending thread self-exits.", host->hostname);
291 + info("STREAM %s [send]: sending thread cleans up...", host->hostname);
292 rrdpush_sender_thread_cleanup_locked_all(host);
293 rrdhost_unlock(host);
294 rrdpush_unlock(host);
@@ -409,7 +414,7 @@ void *rrdpush_sender_thread(void *ptr) {
414 continue;
415 }
416
412 - info("STREAM %s [send to %s]: established communication - sending metrics...", host->hostname, connected_to);
417 + info("STREAM %s [send to %s]: established communication - ready to send metrics...", host->hostname, connected_to);
418 last_sent_t = now_monotonic_sec();
419
420 if(sock_setnonblock(host->rrdpush_socket) < 0)
@@ -736,7 +741,7 @@ static int rrdpush_receive(int fd, const char *key, const char *hostname, const
741 size_t count = pluginsd_process(host, &cd, fp, 1);
742
743 log_stream_connection(client_ip, client_port, key, host->machine_guid, host->hostname, "DISCONNECTED");
739 - error("STREAM %s [receive from [%s]:%s]: disconnected (completed updates %zu).", host->hostname, client_ip, client_port, count);
744 + error("STREAM %s [receive from [%s]:%s]: disconnected (completed %zu updates).", host->hostname, client_ip, client_port, count);
745
746 rrdhost_wrlock(host);
747 host->senders_disconnected_time = now_realtime_sec();
@@ -748,7 +753,8 @@ static int rrdpush_receive(int fd, const char *key, const char *hostname, const
753 }
754 rrdhost_unlock(host);
755
751 - rrdpush_sender_thread_stop(host);
756 + if(host->connected_senders == 0)
757 + rrdpush_sender_thread_stop(host);
758
759 // cleanup
760 fclose(fp);