@cryptotaxi247 / netdata-1 / commits / 87e8d838d

allow faster reconnect on first failure

Costa Tsaousis (ktsaou) committed Oct 1, 2017 at 23:22 UTC 87e8d838d958873c36ebb610cd08212a18f8c096
1 file changed +124 -86
src/rrdpush.c
+124 -86
@@ -263,6 +263,83 @@ static inline void rrdpush_sender_thread_close_socket(RRDHOST *host) {
263 }
264 }
265
266 +static int rrdpush_sender_thread_connect_to_master(RRDHOST *host, int default_port, int timeout, size_t *reconnects_counter, char *connected_to, size_t connected_to_size) {
267 + struct timeval tv = {
268 + .tv_sec = timeout,
269 + .tv_usec = 0
270 + };
271 +
272 + // make sure the socket is closed
273 + rrdpush_sender_thread_close_socket(host);
274 +
275 + debug(D_STREAM, "STREAM: Attempting to connect...");
276 + info("STREAM %s [send to %s]: connecting...", host->hostname, host->rrdpush_send_destination);
277 +
278 + host->rrdpush_sender_socket = connect_to_one_of(
279 + host->rrdpush_send_destination
280 + , default_port
281 + , &tv
282 + , reconnects_counter
283 + , connected_to
284 + , connected_to_size
285 + );
286 +
287 + if(unlikely(host->rrdpush_sender_socket == -1)) {
288 + error("STREAM %s [send to %s]: failed to connect", host->hostname, host->rrdpush_send_destination);
289 + return 0;
290 + }
291 +
292 + info("STREAM %s [send to %s]: initializing communication...", host->hostname, connected_to);
293 +
294 + #define HTTP_HEADER_SIZE 8192
295 + char http[HTTP_HEADER_SIZE + 1];
296 + snprintfz(http, HTTP_HEADER_SIZE,
297 + "STREAM key=%s&hostname=%s&registry_hostname=%s&machine_guid=%s&update_every=%d&os=%s&tags=%s HTTP/1.1\r\n"
298 + "User-Agent: netdata-push-service/%s\r\n"
299 + "Accept: */*\r\n\r\n"
300 + , host->rrdpush_send_api_key
301 + , host->hostname
302 + , host->registry_hostname
303 + , host->machine_guid
304 + , default_rrd_update_every
305 + , host->os
306 + , (host->tags)?host->tags:""
307 + , program_version
308 + );
309 +
310 + if(send_timeout(host->rrdpush_sender_socket, http, strlen(http), 0, timeout) == -1) {
311 + error("STREAM %s [send to %s]: failed to send HTTP header to remote netdata.", host->hostname, connected_to);
312 + rrdpush_sender_thread_close_socket(host);
313 + return 0;
314 + }
315 +
316 + info("STREAM %s [send to %s]: waiting response from remote netdata...", host->hostname, connected_to);
317 +
318 + if(recv_timeout(host->rrdpush_sender_socket, http, HTTP_HEADER_SIZE, 0, timeout) == -1) {
319 + error("STREAM %s [send to %s]: remote netdata does not respond.", host->hostname, connected_to);
320 + rrdpush_sender_thread_close_socket(host);
321 + return 0;
322 + }
323 +
324 + if(strncmp(http, START_STREAMING_PROMPT, strlen(START_STREAMING_PROMPT)) != 0) {
325 + error("STREAM %s [send to %s]: server is not replying properly (is it a netdata?).", host->hostname, connected_to);
326 + rrdpush_sender_thread_close_socket(host);
327 + return 0;
328 + }
329 +
330 + info("STREAM %s [send to %s]: established communication - ready to send metrics...", host->hostname, connected_to);
331 +
332 + if(sock_setnonblock(host->rrdpush_sender_socket) < 0)
333 + error("STREAM %s [send to %s]: cannot set non-blocking mode for socket.", host->hostname, connected_to);
334 +
335 + if(sock_enlarge_out(host->rrdpush_sender_socket) < 0)
336 + error("STREAM %s [send to %s]: cannot enlarge the socket buffer.", host->hostname, connected_to);
337 +
338 + debug(D_STREAM, "STREAM: Connected on fd %d...", host->rrdpush_sender_socket);
339 +
340 + return 1;
341 +}
342 +
343 static void rrdpush_sender_thread_cleanup_callback(void *ptr) {
344 RRDHOST *host = (RRDHOST *)ptr;
345
@@ -334,12 +411,8 @@ void *rrdpush_sender_thread(void *ptr) {
411 size_t begin = 0;
412 size_t reconnects_counter = 0;
413 size_t sent_bytes = 0;
337 - size_t sent_connection = 0;
414 + size_t sent_bytes_on_this_connection = 0;
415
339 - struct timeval tv = {
340 - .tv_sec = timeout,
341 - .tv_usec = 0
342 - };
416
417 time_t last_sent_t = 0;
418 struct pollfd fds[2], *ifd, *ofd;
@@ -348,90 +421,54 @@ void *rrdpush_sender_thread(void *ptr) {
421 ifd = &fds[0];
422 ofd = &fds[1];
423
424 + size_t not_connected_loops = 0;
425 +
426 pthread_cleanup_push(rrdpush_sender_thread_cleanup_callback, host);
427
428 for(; host->rrdpush_send_enabled && !netdata_exit ;) {
429 // check for outstanding cancellation requests
430 pthread_testcancel();
431
357 - if(host->rrdpush_sender_socket == -1)
358 - sleep(reconnect_delay);
359 -
360 - debug(D_STREAM, "STREAM: Checking if we need to timeout the connection...");
361 - if(host->rrdpush_sender_socket != -1 && now_monotonic_sec() - last_sent_t > timeout) {
362 - 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);
363 - rrdpush_sender_thread_close_socket(host);
364 - }
365 -
432 + // if we don't have socket open, lets wait a bit
433 if(unlikely(host->rrdpush_sender_socket == -1)) {
367 - debug(D_STREAM, "STREAM: Attempting to connect...");
368 -
369 - // stop appending data into rrdpush_sender_buffer
370 - // they will be lost, so there is no point to do it
371 - host->rrdpush_sender_connected = 0;
372 -
373 - info("STREAM %s [send to %s]: connecting...", host->hostname, host->rrdpush_send_destination);
374 - host->rrdpush_sender_socket = connect_to_one_of(host->rrdpush_send_destination, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
375 -
376 - if(unlikely(host->rrdpush_sender_socket == -1)) {
377 - error("STREAM %s [send to %s]: failed to connect", host->hostname, host->rrdpush_send_destination);
378 - continue;
434 + if(not_connected_loops == 0 && sent_bytes_on_this_connection > 0) {
435 + // fast re-connection on first disconnect
436 + sleep_usec(USEC_PER_MS * 500); // milliseconds
437 }
380 -
381 - info("STREAM %s [send to %s]: initializing communication...", host->hostname, connected_to);
382 -
383 - #define HTTP_HEADER_SIZE 8192
384 - char http[HTTP_HEADER_SIZE + 1];
385 - snprintfz(http, HTTP_HEADER_SIZE,
386 - "STREAM key=%s&hostname=%s&registry_hostname=%s&machine_guid=%s&update_every=%d&os=%s&tags=%s HTTP/1.1\r\n"
387 - "User-Agent: netdata-push-service/%s\r\n"
388 - "Accept: */*\r\n\r\n"
389 - , host->rrdpush_send_api_key
390 - , host->hostname
391 - , host->registry_hostname
392 - , host->machine_guid
393 - , default_rrd_update_every
394 - , host->os
395 - , (host->tags)?host->tags:""
396 - , program_version
397 - );
398 -
399 - if(send_timeout(host->rrdpush_sender_socket, http, strlen(http), 0, timeout) == -1) {
400 - error("STREAM %s [send to %s]: failed to send http header to netdata", host->hostname, connected_to);
401 - rrdpush_sender_thread_close_socket(host);
402 - continue;
438 + else {
439 + // slow re-connection on repeating errors
440 + sleep_usec(USEC_PER_SEC * reconnect_delay); // seconds
441 }
442
405 - info("STREAM %s [send to %s]: waiting response from remote netdata...", host->hostname, connected_to);
406 -
407 - if(recv_timeout(host->rrdpush_sender_socket, http, HTTP_HEADER_SIZE, 0, timeout) == -1) {
408 - error("STREAM %s [send to %s]: failed to initialize communication", host->hostname, connected_to);
409 - rrdpush_sender_thread_close_socket(host);
410 - continue;
411 - }
412 -
413 - if(strncmp(http, START_STREAMING_PROMPT, strlen(START_STREAMING_PROMPT))) {
414 - error("STREAM %s [send to %s]: server is not replying properly.", host->hostname, connected_to);
415 - rrdpush_sender_thread_close_socket(host);
416 - continue;
417 - }
443 + if(rrdpush_sender_thread_connect_to_master(host, default_port, timeout, &reconnects_counter, connected_to, CONNECTED_TO_SIZE)) {
444 + last_sent_t = now_monotonic_sec();
445
419 - info("STREAM %s [send to %s]: established communication - ready to send metrics...", host->hostname, connected_to);
420 - last_sent_t = now_monotonic_sec();
446 + // reset the buffer, to properly send charts and metrics
447 + rrdpush_sender_thread_data_flush(host);
448
422 - if(sock_setnonblock(host->rrdpush_sender_socket) < 0)
423 - error("STREAM %s [send to %s]: cannot set non-blocking mode for socket.", host->hostname, connected_to);
449 + // make sure the next reconnection will be immediate
450 + not_connected_loops = 0;
451
425 - if(sock_enlarge_out(host->rrdpush_sender_socket) < 0)
426 - error("STREAM %s [send to %s]: cannot enlarge the socket buffer.", host->hostname, connected_to);
452 + // reset the bytes we have sent for this session
453 + sent_bytes_on_this_connection = 0;
454
428 - rrdpush_sender_thread_data_flush(host);
429 - sent_connection = 0;
455 + // let the data collection threads know we are ready
456 + host->rrdpush_sender_connected = 1;
457 + }
458 + else {
459 + // increase the failed connections counter
460 + not_connected_loops++;
461
431 - // allow appending data into rrdpush_sender_buffer
432 - host->rrdpush_sender_connected = 1;
462 + // reset the number of bytes sent
463 + sent_bytes_on_this_connection = 0;
464 + }
465
434 - debug(D_STREAM, "STREAM: Connected on fd %d...", host->rrdpush_sender_socket);
466 + // loop through
467 + continue;
468 + }
469 + else if(unlikely(now_monotonic_sec() - last_sent_t > timeout)) {
470 + 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_bytes_on_this_connection);
471 + rrdpush_sender_thread_close_socket(host);
472 }
473
474 ifd->fd = host->rrdpush_sender_pipe[PIPE_READ];
@@ -452,21 +489,22 @@ void *rrdpush_sender_thread(void *ptr) {
489 }
490
491 debug(D_STREAM, "STREAM: Waiting for poll() events (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_sender_buffer));
455 - if(netdata_exit) break;
492 + if(unlikely(netdata_exit)) break;
493 int retval = poll(fds, fdmax, 1000);
457 - if(netdata_exit) break;
494 + if(unlikely(netdata_exit)) break;
495
496 if(unlikely(retval == -1)) {
497 debug(D_STREAM, "STREAM: poll() failed (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_sender_buffer));
498
499 if(errno == EAGAIN || errno == EINTR) {
500 debug(D_STREAM, "STREAM: poll() failed with EAGAIN or EINTR...");
464 - continue;
501 + }
502 + else {
503 + error("STREAM %s [send to %s]: failed to poll(). Closing socket.", host->hostname, connected_to);
504 + rrdpush_sender_thread_close_socket(host);
505 }
506
467 - error("STREAM %s [send to %s]: failed to poll().", host->hostname, connected_to);
468 - rrdpush_sender_thread_close_socket(host);
469 - break;
507 + continue;
508 }
509 else if(likely(retval)) {
510 if (ifd->revents & POLLIN || ifd->revents & POLLPRI) {
@@ -499,7 +537,7 @@ void *rrdpush_sender_thread(void *ptr) {
537 if (unlikely(ret == -1)) {
538 if (errno != EAGAIN && errno != EINTR && errno != EWOULDBLOCK) {
539 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);
540 + 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_bytes_on_this_connection);
541 rrdpush_sender_thread_close_socket(host);
542 }
543 else {
@@ -513,7 +551,7 @@ void *rrdpush_sender_thread(void *ptr) {
551 //debug(D_STREAM, "STREAM: sent from %zu to %zd:\n%s\n", begin, ret, &host->rrdpush_sender_buffer->buffer[begin]);
552 //host->rrdpush_sender_buffer->buffer[begin + ret] = c;
553
516 - sent_connection += ret;
554 + sent_bytes_on_this_connection += ret;
555 sent_bytes += ret;
556 begin += ret;
557
@@ -533,7 +571,7 @@ void *rrdpush_sender_thread(void *ptr) {
571 else {
572 debug(D_STREAM, "STREAM: send() returned %zd - closing the socket...", ret);
573 error("STREAM %s [send to %s]: failed to send metrics (send() returned %zd) - closing connection - we have sent %zu bytes on this connection.",
536 - host->hostname, connected_to, ret, sent_connection);
574 + host->hostname, connected_to, ret, sent_bytes_on_this_connection);
575 rrdpush_sender_thread_close_socket(host);
576 }
577
@@ -552,17 +590,17 @@ void *rrdpush_sender_thread(void *ptr) {
590
591 if(unlikely(ofd->revents & POLLERR)) {
592 debug(D_STREAM, "STREAM: Send failed (POLLERR) - closing socket...");
555 - 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);
593 + 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_bytes_on_this_connection);
594 rrdpush_sender_thread_close_socket(host);
595 }
596 else if(unlikely(ofd->revents & POLLHUP)) {
597 debug(D_STREAM, "STREAM: Send failed (POLLHUP) - closing socket...");
560 - 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);
598 + 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_bytes_on_this_connection);
599 rrdpush_sender_thread_close_socket(host);
600 }
601 else if(unlikely(ofd->revents & POLLNVAL)) {
602 debug(D_STREAM, "STREAM: Send failed (POLLNVAL) - closing socket...");
565 - 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);
603 + 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_bytes_on_this_connection);
604 rrdpush_sender_thread_close_socket(host);
605 }
606 }
@@ -574,7 +612,7 @@ void *rrdpush_sender_thread(void *ptr) {
612 if(buffer_strlen(host->rrdpush_sender_buffer) > max_size) {
613 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);
614 errno = 0;
577 - 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);
615 + 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_bytes_on_this_connection);
616 rrdpush_sender_thread_close_socket(host);
617 }
618 }