@cryptotaxi247 / netdata-1 / commits / eca8792ee

detect streaming cannot send data and re-connect; fixes #2154

Costa Tsaousis (ktsaou) committed May 7, 2017 at 02:43 UTC eca8792eea95d80e5fa42cbadf666a0824630990
1 file changed +55 -38
src/rrdpush.c
+55 -38
@@ -274,6 +274,7 @@ void *rrdpush_sender_thread(void *ptr) {
274 .tv_usec = 0
275 };
276
277 + time_t last_sent_t = 0;
278 struct pollfd fds[2], *ifd, *ofd;
279 nfds_t fdmax;
280
@@ -281,6 +282,11 @@ void *rrdpush_sender_thread(void *ptr) {
282 ofd = &fds[1];
283
284 for(; host->rrdpush_enabled && !netdata_exit ;) {
285 + if(host->rrdpush_socket != -1 && now_monotonic_sec() - last_sent_t > timeout) {
286 + 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);
287 + close(host->rrdpush_socket);
288 + host->rrdpush_socket = -1;
289 + }
290
291 if(unlikely(host->rrdpush_socket == -1)) {
292 // stop appending data into rrdpush_buffer
@@ -338,10 +344,14 @@ void *rrdpush_sender_thread(void *ptr) {
344 }
345
346 info("STREAM %s [send to %s]: established communication - sending metrics...", host->hostname, connected_to);
347 + last_sent_t = now_monotonic_sec();
348
349 if(sock_setnonblock(host->rrdpush_socket) < 0)
350 error("STREAM %s [send to %s]: cannot set non-blocking mode for socket.", host->hostname, connected_to);
351
352 + if(sock_enlarge_out(host->rrdpush_socket) < 0)
353 + error("STREAM %s [send to %s]: cannot enlarge the socket buffer.", host->hostname, connected_to);
354 +
355 rrdpush_sender_thread_data_flush(host);
356 sent_connection = 0;
357
@@ -377,58 +387,65 @@ void *rrdpush_sender_thread(void *ptr) {
387 host->rrdpush_socket = -1;
388 break;
389 }
380 - else if(unlikely(!retval)) {
381 - // timeout
382 - continue;
383 - }
390 + else if(likely(retval)) {
391 + if (ifd->revents & POLLIN) {
392 + char buffer[1000 + 1];
393 + if (read(host->rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
394 + error("STREAM %s [send to %s]: cannot read from internal pipe.", host->hostname, connected_to);
395 + }
396
385 - if(ifd->revents & POLLIN) {
386 - char buffer[1000 + 1];
387 - if(read(host->rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
388 - error("STREAM %s [send to %s]: cannot read from internal pipe.", host->hostname, connected_to);
389 - }
397 + if (ofd->revents & POLLOUT && begin < buffer_strlen(host->rrdpush_buffer)) {
398
391 - if(ofd->revents & POLLOUT && begin < buffer_strlen(host->rrdpush_buffer)) {
399 + // BEGIN RRDPUSH LOCKED SESSION
400
393 - // BEGIN RRDPUSH LOCKED SESSION
401 + // during this session, data collectors
402 + // will not be able to append data to our buffer
403 + // but the socket is in non-blocking mode
404 + // so, we will not block at send()
405
395 - // during this session, data collectors
396 - // will not be able to append data to our buffer
397 - // but the socket is in non-blocking mode
398 - // so, we will not block at send()
406 + if (pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL) != 0)
407 + error("STREAM %s [send]: cannot set pthread cancel state to DISABLE.", host->hostname);
408
400 - if(pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL) != 0)
401 - error("STREAM %s [send]: cannot set pthread cancel state to DISABLE.", host->hostname);
409 + rrdpush_lock(host);
410
403 - rrdpush_lock(host);
411 + ssize_t ret = send(host->rrdpush_socket, &host->rrdpush_buffer->buffer[begin],
412 + buffer_strlen(host->rrdpush_buffer) - begin, MSG_DONTWAIT);
413 + if (unlikely(ret == -1)) {
414 + if (errno != EAGAIN && errno != EINTR && errno != EWOULDBLOCK) {
415 + 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);
416 + close(host->rrdpush_socket);
417 + host->rrdpush_socket = -1;
418 + }
419 + }
420 + else if(likely(ret > 0)) {
421 + sent_connection += ret;
422 + sent_bytes += ret;
423 + begin += ret;
424
405 - ssize_t ret = send(host->rrdpush_socket, &host->rrdpush_buffer->buffer[begin], buffer_strlen(host->rrdpush_buffer) - begin, MSG_DONTWAIT);
406 - if(ret == -1) {
407 - if(errno != EAGAIN && errno != EINTR) {
408 - 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);
425 + if (begin == buffer_strlen(host->rrdpush_buffer)) {
426 + // we send it all
427 +
428 + buffer_flush(host->rrdpush_buffer);
429 + begin = 0;
430 + }
431 +
432 + last_sent_t = now_monotonic_sec();
433 + }
434 + else {
435 + error("STREAM %s [send to %s]: failed to send metrics (send() returned %zd) - closing connection - we have sent %zu bytes on this connection.", host->hostname, connected_to, ret, sent_connection);
436 close(host->rrdpush_socket);
437 host->rrdpush_socket = -1;
438 }
412 - }
413 - else {
414 - sent_connection += ret;
415 - sent_bytes += ret;
416 - begin += ret;
417 - if(begin == buffer_strlen(host->rrdpush_buffer)) {
418 - // we send it all
419 -
420 - buffer_flush(host->rrdpush_buffer);
421 - begin = 0;
422 - }
423 - }
439
425 - rrdpush_unlock(host);
440 + rrdpush_unlock(host);
441
427 - if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
428 - error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
442 + if (pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
443 + error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
444
430 - // END RRDPUSH LOCKED SESSION
445 + // END RRDPUSH LOCKED SESSION
446 + }
447 }
448 + // else timeout
449
450 // protection from overflow
451 if(host->rrdpush_buffer->len > max_size) {