@cryptotaxi247 / netdata-1 / commits / 2595eadef

attempt to fix crashing proxies under load

Costa Tsaousis (ktsaou) committed Sep 23, 2017 at 02:34 UTC 2595eadef02ba12e7cbd85b61f0fe6348b272cb4
2 files changed +253 -224
src/rrd.h
+3 -2
@@ -425,11 +425,12 @@ struct rrdhost {
425 // streaming of data to remote hosts - rrdpush
426
427 int rrdpush_enabled:1; // 1 when this host sends metrics to another netdata
428 - char *rrdpush_destination; // where to send metrics to
429 - char *rrdpush_api_key; // the api key at the receiving 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
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
src/rrdpush.c
+250 -222
@@ -170,7 +170,7 @@ void rrdset_done_push(RRDSET *st) {
170 send_chart_metrics(st);
171
172 // signal the sender there are more data
173 - if(write(host->rrdpush_pipe[PIPE_WRITE], " ", 1) == -1)
173 + if(host->rrdpush_pipe[PIPE_WRITE] != -1 && write(host->rrdpush_pipe[PIPE_WRITE], " ", 1) == -1)
174 error("STREAM %s [send]: cannot write to internal pipe", host->hostname);
175
176 rrdpush_unlock(host);
@@ -236,34 +236,76 @@ 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)
240 + pthread_detach(pthread_self());
241 +
242 host->rrdpush_spawn = 0;
243 +
244 + pthread_exit(NULL);
245 }
246
247 void rrdpush_sender_thread_stop(RRDHOST *host) {
248 rrdpush_lock(host);
249 rrdhost_wrlock(host);
250
251 + pthread_t thr = 0;
252 +
253 if(host->rrdpush_spawn) {
254 info("STREAM %s [send]: stopping sending thread...", host->hostname);
248 - pthread_cancel(host->rrdpush_thread);
249 - rrdpush_sender_thread_cleanup_locked_all(host);
255 +
256 + // signal the thread that we want to join it
257 + host->rrdpush_sender_join = 1;
258 +
259 + // copy the thread id, so that we will be waiting for the right one
260 + // even if a new one has been spawn
261 + thr = host->rrdpush_thread;
262 +
263 + // signal it to cancel
264 + int ret = pthread_cancel(host->rrdpush_thread);
265 + if(ret != 0)
266 + error("STREAM %s [send]: pthread_cancel() returned error.", host->hostname);
267 }
268
269 rrdhost_unlock(host);
270 rrdpush_unlock(host);
271 +
272 + if(thr != 0) {
273 + info("STREAM %s [send]: waiting for sending thread to stop...", host->hostname);
274 + void *result;
275 + int ret = pthread_join(thr, &result);
276 + if(ret != 0)
277 + error("STREAM %s [send]: pthread_join() returned error.", host->hostname);
278 + }
279 +}
280 +
281 +static void rrdpush_sender_thread_cleanup_callback(void *ptr) {
282 + RRDHOST *host = (RRDHOST *)ptr;
283 +
284 + rrdpush_lock(host);
285 + rrdhost_wrlock(host);
286 + info("STREAM %s [send]: sending thread self-exits.", host->hostname);
287 + rrdpush_sender_thread_cleanup_locked_all(host);
288 + rrdhost_unlock(host);
289 + rrdpush_unlock(host);
290 }
291
292 void *rrdpush_sender_thread(void *ptr) {
293 RRDHOST *host = (RRDHOST *)ptr;
294
259 - info("STREAM %s [send]: thread created (task id %d)", host->hostname, gettid());
295 + if(!host->rrdpush_enabled || !host->rrdpush_destination || !*host->rrdpush_destination || !host->rrdpush_api_key || !*host->rrdpush_api_key) {
296 + error("STREAM %s [send]: thread created (task id %d), but host has streaming disabled.", host->hostname, gettid());
297 + pthread_exit(NULL);
298 + return NULL;
299 + }
300
261 - if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
262 - error("STREAM %s [send]: cannot set pthread cancel type to DEFERRED.", host->hostname);
301 + info("STREAM %s [send]: thread created (task id %d)", host->hostname, gettid());
302
303 if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
304 error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
305
306 + if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
307 + error("STREAM %s [send]: cannot set pthread cancel type to DEFERRED.", host->hostname);
308 +
309 int timeout = (int)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "timeout seconds", 60);
310 int default_port = (int)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "default port", 19999);
311 size_t max_size = (size_t)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "buffer size bytes", 1024 * 1024);
@@ -271,9 +313,6 @@ void *rrdpush_sender_thread(void *ptr) {
313 remote_clock_resync_iterations = (unsigned int)appconfig_get_number(&stream_config, CONFIG_SECTION_STREAM, "initial clock resync iterations", remote_clock_resync_iterations);
314 char connected_to[CONNECTED_TO_SIZE + 1] = "";
315
274 - if(!host->rrdpush_enabled || !host->rrdpush_destination || !*host->rrdpush_destination || !host->rrdpush_api_key || !*host->rrdpush_api_key)
275 - goto cleanup;
276 -
316 // initialize rrdpush globals
317 host->rrdpush_buffer = buffer_create(1);
318 host->rrdpush_connected = 0;
@@ -297,260 +336,252 @@ void *rrdpush_sender_thread(void *ptr) {
336 ifd = &fds[0];
337 ofd = &fds[1];
338
300 - for(; host->rrdpush_enabled && !netdata_exit ;) {
301 - debug(D_STREAM, "STREAM: Checking if we need to timeout the connection...");
302 - if(host->rrdpush_socket != -1 && now_monotonic_sec() - last_sent_t > timeout) {
303 - 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);
304 - close(host->rrdpush_socket);
305 - host->rrdpush_socket = -1;
306 - }
307 -
308 - if(unlikely(host->rrdpush_socket == -1)) {
309 - debug(D_STREAM, "STREAM: Attempting to connect...");
310 -
311 - // stop appending data into rrdpush_buffer
312 - // they will be lost, so there is no point to do it
313 - host->rrdpush_connected = 0;
314 -
315 - info("STREAM %s [send to %s]: connecting...", host->hostname, host->rrdpush_destination);
316 - host->rrdpush_socket = connect_to_one_of(host->rrdpush_destination, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
339 + pthread_cleanup_push(rrdpush_sender_thread_cleanup_callback, host);
340
318 - if(unlikely(host->rrdpush_socket == -1)) {
319 - error("STREAM %s [send to %s]: failed to connect", host->hostname, host->rrdpush_destination);
320 - sleep(reconnect_delay);
321 - continue;
322 - }
323 -
324 - info("STREAM %s [send to %s]: initializing communication...", host->hostname, connected_to);
325 -
326 - #define HTTP_HEADER_SIZE 8192
327 - char http[HTTP_HEADER_SIZE + 1];
328 - snprintfz(http, HTTP_HEADER_SIZE,
329 - "STREAM key=%s&hostname=%s&registry_hostname=%s&machine_guid=%s&update_every=%d&os=%s&tags=%s HTTP/1.1\r\n"
330 - "User-Agent: netdata-push-service/%s\r\n"
331 - "Accept: */*\r\n\r\n"
332 - , host->rrdpush_api_key
333 - , host->hostname
334 - , host->registry_hostname
335 - , host->machine_guid
336 - , default_rrd_update_every
337 - , host->os
338 - , (host->tags)?host->tags:""
339 - , program_version
340 - );
341 + for(; host->rrdpush_enabled && !netdata_exit ;) {
342 + // check for outstanding cancellation requests
343 + pthread_testcancel();
344
342 - if(send_timeout(host->rrdpush_socket, http, strlen(http), 0, timeout) == -1) {
345 + debug(D_STREAM, "STREAM: Checking if we need to timeout the connection...");
346 + if(host->rrdpush_socket != -1 && now_monotonic_sec() - last_sent_t > timeout) {
347 + 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);
348 close(host->rrdpush_socket);
349 host->rrdpush_socket = -1;
345 - error("STREAM %s [send to %s]: failed to send http header to netdata", host->hostname, connected_to);
346 - sleep(reconnect_delay);
347 - continue;
350 }
351
350 - info("STREAM %s [send to %s]: waiting response from remote netdata...", host->hostname, connected_to);
352 + if(unlikely(host->rrdpush_socket == -1)) {
353 + debug(D_STREAM, "STREAM: Attempting to connect...");
354
352 - if(recv_timeout(host->rrdpush_socket, http, HTTP_HEADER_SIZE, 0, timeout) == -1) {
353 - close(host->rrdpush_socket);
354 - host->rrdpush_socket = -1;
355 - error("STREAM %s [send to %s]: failed to initialize communication", host->hostname, connected_to);
356 - sleep(reconnect_delay);
357 - continue;
358 - }
355 + // stop appending data into rrdpush_buffer
356 + // they will be lost, so there is no point to do it
357 + host->rrdpush_connected = 0;
358
360 - if(strncmp(http, START_STREAMING_PROMPT, strlen(START_STREAMING_PROMPT))) {
361 - close(host->rrdpush_socket);
362 - host->rrdpush_socket = -1;
363 - error("STREAM %s [send to %s]: server is not replying properly.", host->hostname, connected_to);
364 - sleep(reconnect_delay);
365 - continue;
366 - }
359 + info("STREAM %s [send to %s]: connecting...", host->hostname, host->rrdpush_destination);
360 + host->rrdpush_socket = connect_to_one_of(host->rrdpush_destination, default_port, &tv, &reconnects_counter, connected_to, CONNECTED_TO_SIZE);
361
368 - info("STREAM %s [send to %s]: established communication - sending metrics...", host->hostname, connected_to);
369 - last_sent_t = now_monotonic_sec();
362 + if(unlikely(host->rrdpush_socket == -1)) {
363 + error("STREAM %s [send to %s]: failed to connect", host->hostname, host->rrdpush_destination);
364 + sleep(reconnect_delay);
365 + continue;
366 + }
367
371 - if(sock_setnonblock(host->rrdpush_socket) < 0)
372 - error("STREAM %s [send to %s]: cannot set non-blocking mode for socket.", host->hostname, connected_to);
368 + info("STREAM %s [send to %s]: initializing communication...", host->hostname, connected_to);
369 +
370 + #define HTTP_HEADER_SIZE 8192
371 + char http[HTTP_HEADER_SIZE + 1];
372 + snprintfz(http, HTTP_HEADER_SIZE,
373 + "STREAM key=%s&hostname=%s&registry_hostname=%s&machine_guid=%s&update_every=%d&os=%s&tags=%s HTTP/1.1\r\n"
374 + "User-Agent: netdata-push-service/%s\r\n"
375 + "Accept: */*\r\n\r\n"
376 + , host->rrdpush_api_key
377 + , host->hostname
378 + , host->registry_hostname
379 + , host->machine_guid
380 + , default_rrd_update_every
381 + , host->os
382 + , (host->tags)?host->tags:""
383 + , program_version
384 + );
385 +
386 + if(send_timeout(host->rrdpush_socket, http, strlen(http), 0, timeout) == -1) {
387 + close(host->rrdpush_socket);
388 + host->rrdpush_socket = -1;
389 + error("STREAM %s [send to %s]: failed to send http header to netdata", host->hostname, connected_to);
390 + sleep(reconnect_delay);
391 + continue;
392 + }
393
374 - if(sock_enlarge_out(host->rrdpush_socket) < 0)
375 - error("STREAM %s [send to %s]: cannot enlarge the socket buffer.", host->hostname, connected_to);
394 + info("STREAM %s [send to %s]: waiting response from remote netdata...", host->hostname, connected_to);
395
377 - rrdpush_sender_thread_data_flush(host);
378 - sent_connection = 0;
396 + if(recv_timeout(host->rrdpush_socket, http, HTTP_HEADER_SIZE, 0, timeout) == -1) {
397 + close(host->rrdpush_socket);
398 + host->rrdpush_socket = -1;
399 + error("STREAM %s [send to %s]: failed to initialize communication", host->hostname, connected_to);
400 + sleep(reconnect_delay);
401 + continue;
402 + }
403
380 - // allow appending data into rrdpush_buffer
381 - host->rrdpush_connected = 1;
404 + if(strncmp(http, START_STREAMING_PROMPT, strlen(START_STREAMING_PROMPT))) {
405 + close(host->rrdpush_socket);
406 + host->rrdpush_socket = -1;
407 + error("STREAM %s [send to %s]: server is not replying properly.", host->hostname, connected_to);
408 + sleep(reconnect_delay);
409 + continue;
410 + }
411
383 - debug(D_STREAM, "STREAM: Connected on fd %d...", host->rrdpush_socket);
384 - }
412 + info("STREAM %s [send to %s]: established communication - sending metrics...", host->hostname, connected_to);
413 + last_sent_t = now_monotonic_sec();
414
386 - ifd->fd = host->rrdpush_pipe[PIPE_READ];
387 - ifd->events = POLLIN;
388 - ifd->revents = 0;
415 + if(sock_setnonblock(host->rrdpush_socket) < 0)
416 + error("STREAM %s [send to %s]: cannot set non-blocking mode for socket.", host->hostname, connected_to);
417
390 - ofd->fd = host->rrdpush_socket;
391 - ofd->revents = 0;
392 - if(ofd->fd != -1 && begin < buffer_strlen(host->rrdpush_buffer)) {
393 - debug(D_STREAM, "STREAM: Requesting data output on streaming socket %d...", ofd->fd);
394 - ofd->events = POLLOUT;
395 - fdmax = 2;
396 - }
397 - else {
398 - debug(D_STREAM, "STREAM: Not requesting data output on streaming socket %d (nothing to send now)...", ofd->fd);
399 - ofd->events = 0;
400 - fdmax = 1;
401 - }
418 + if(sock_enlarge_out(host->rrdpush_socket) < 0)
419 + error("STREAM %s [send to %s]: cannot enlarge the socket buffer.", host->hostname, connected_to);
420
403 - debug(D_STREAM, "STREAM: Waiting for poll() events (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
404 - if(netdata_exit) break;
405 - int retval = poll(fds, fdmax, 1000);
406 - if(netdata_exit) break;
421 + rrdpush_sender_thread_data_flush(host);
422 + sent_connection = 0;
423
408 - if(unlikely(retval == -1)) {
409 - debug(D_STREAM, "STREAM: poll() failed (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
424 + // allow appending data into rrdpush_buffer
425 + host->rrdpush_connected = 1;
426
411 - if(errno == EAGAIN || errno == EINTR) {
412 - debug(D_STREAM, "STREAM: poll() failed with EAGAIN or EINTR...");
413 - continue;
427 + debug(D_STREAM, "STREAM: Connected on fd %d...", host->rrdpush_socket);
428 }
429
416 - error("STREAM %s [send to %s]: failed to poll().", host->hostname, connected_to);
417 - close(host->rrdpush_socket);
418 - host->rrdpush_socket = -1;
419 - break;
420 - }
421 - else if(likely(retval)) {
422 - if (ifd->revents & POLLIN || ifd->revents & POLLPRI) {
423 - debug(D_STREAM, "STREAM: Data added to send buffer (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
430 + ifd->fd = host->rrdpush_pipe[PIPE_READ];
431 + ifd->events = POLLIN;
432 + ifd->revents = 0;
433
425 - char buffer[1000 + 1];
426 - if (read(host->rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
427 - error("STREAM %s [send to %s]: cannot read from internal pipe.", host->hostname, connected_to);
434 + ofd->fd = host->rrdpush_socket;
435 + ofd->revents = 0;
436 + if(ofd->fd != -1 && begin < buffer_strlen(host->rrdpush_buffer)) {
437 + debug(D_STREAM, "STREAM: Requesting data output on streaming socket %d...", ofd->fd);
438 + ofd->events = POLLOUT;
439 + fdmax = 2;
440 + }
441 + else {
442 + debug(D_STREAM, "STREAM: Not requesting data output on streaming socket %d (nothing to send now)...", ofd->fd);
443 + ofd->events = 0;
444 + fdmax = 1;
445 }
446
430 - if (ofd->revents & POLLOUT) {
431 - if (begin < buffer_strlen(host->rrdpush_buffer)) {
432 - debug(D_STREAM, "STREAM: Sending data (current buffer length %zu bytes, begin = %zu)...", buffer_strlen(host->rrdpush_buffer), begin);
447 + debug(D_STREAM, "STREAM: Waiting for poll() events (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
448 + if(netdata_exit) break;
449 + int retval = poll(fds, fdmax, 1000);
450 + if(netdata_exit) break;
451
434 - // BEGIN RRDPUSH LOCKED SESSION
452 + if(unlikely(retval == -1)) {
453 + debug(D_STREAM, "STREAM: poll() failed (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
454
436 - // during this session, data collectors
437 - // will not be able to append data to our buffer
438 - // but the socket is in non-blocking mode
439 - // so, we will not block at send()
455 + if(errno == EAGAIN || errno == EINTR) {
456 + debug(D_STREAM, "STREAM: poll() failed with EAGAIN or EINTR...");
457 + continue;
458 + }
459
441 - if (pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL) != 0)
442 - error("STREAM %s [send]: cannot set pthread cancel state to DISABLE.", host->hostname);
460 + error("STREAM %s [send to %s]: failed to poll().", host->hostname, connected_to);
461 + close(host->rrdpush_socket);
462 + host->rrdpush_socket = -1;
463 + break;
464 + }
465 + else if(likely(retval)) {
466 + if (ifd->revents & POLLIN || ifd->revents & POLLPRI) {
467 + debug(D_STREAM, "STREAM: Data added to send buffer (current buffer length %zu bytes)...", buffer_strlen(host->rrdpush_buffer));
468
444 - debug(D_STREAM, "STREAM: Getting exclusive lock on host...");
445 - rrdpush_lock(host);
469 + char buffer[1000 + 1];
470 + if (read(host->rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
471 + error("STREAM %s [send to %s]: cannot read from internal pipe.", host->hostname, connected_to);
472 + }
473
447 - debug(D_STREAM, "STREAM: Sending data, starting from %zu, size %zu...", begin, buffer_strlen(host->rrdpush_buffer));
448 - ssize_t ret = send(host->rrdpush_socket, &host->rrdpush_buffer->buffer[begin], buffer_strlen(host->rrdpush_buffer) - begin, MSG_DONTWAIT);
449 - if (unlikely(ret == -1)) {
450 - if (errno != EAGAIN && errno != EINTR && errno != EWOULDBLOCK) {
451 - debug(D_STREAM, "STREAM: Send failed - closing socket...");
452 - 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);
453 - close(host->rrdpush_socket);
454 - host->rrdpush_socket = -1;
474 + if (ofd->revents & POLLOUT) {
475 + if (begin < buffer_strlen(host->rrdpush_buffer)) {
476 + debug(D_STREAM, "STREAM: Sending data (current buffer length %zu bytes, begin = %zu)...", buffer_strlen(host->rrdpush_buffer), begin);
477 +
478 + // BEGIN RRDPUSH LOCKED SESSION
479 +
480 + // during this session, data collectors
481 + // will not be able to append data to our buffer
482 + // but the socket is in non-blocking mode
483 + // so, we will not block at send()
484 +
485 + if (pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL) != 0)
486 + error("STREAM %s [send]: cannot set pthread cancel state to DISABLE.", host->hostname);
487 +
488 + debug(D_STREAM, "STREAM: Getting exclusive lock on host...");
489 + rrdpush_lock(host);
490 +
491 + debug(D_STREAM, "STREAM: Sending data, starting from %zu, size %zu...", begin, buffer_strlen(host->rrdpush_buffer));
492 + ssize_t ret = send(host->rrdpush_socket, &host->rrdpush_buffer->buffer[begin], buffer_strlen(host->rrdpush_buffer) - begin, MSG_DONTWAIT);
493 + if (unlikely(ret == -1)) {
494 + if (errno != EAGAIN && errno != EINTR && errno != EWOULDBLOCK) {
495 + debug(D_STREAM, "STREAM: Send failed - closing socket...");
496 + 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);
497 + close(host->rrdpush_socket);
498 + host->rrdpush_socket = -1;
499 + }
500 + else {
501 + debug(D_STREAM, "STREAM: Send failed - will retry...");
502 + }
503 }
456 - else {
457 - debug(D_STREAM, "STREAM: Send failed - will retry...");
458 - }
459 - }
460 - else if (likely(ret > 0)) {
461 - // DEBUG - dump the scring to see it
462 - //char c = host->rrdpush_buffer->buffer[begin + ret];
463 - //host->rrdpush_buffer->buffer[begin + ret] = '\0';
464 - //debug(D_STREAM, "STREAM: sent from %zu to %zd:\n%s\n", begin, ret, &host->rrdpush_buffer->buffer[begin]);
465 - //host->rrdpush_buffer->buffer[begin + ret] = c;
466 -
467 - sent_connection += ret;
468 - sent_bytes += ret;
469 - begin += ret;
470 -
471 - if (begin == buffer_strlen(host->rrdpush_buffer)) {
472 - // we send it all
473 -
474 - debug(D_STREAM, "STREAM: Sent %zd bytes (the whole buffer)...", ret);
475 - buffer_flush(host->rrdpush_buffer);
476 - begin = 0;
504 + else if (likely(ret > 0)) {
505 + // DEBUG - dump the scring to see it
506 + //char c = host->rrdpush_buffer->buffer[begin + ret];
507 + //host->rrdpush_buffer->buffer[begin + ret] = '\0';
508 + //debug(D_STREAM, "STREAM: sent from %zu to %zd:\n%s\n", begin, ret, &host->rrdpush_buffer->buffer[begin]);
509 + //host->rrdpush_buffer->buffer[begin + ret] = c;
510 +
511 + sent_connection += ret;
512 + sent_bytes += ret;
513 + begin += ret;
514 +
515 + if (begin == buffer_strlen(host->rrdpush_buffer)) {
516 + // we send it all
517 +
518 + debug(D_STREAM, "STREAM: Sent %zd bytes (the whole buffer)...", ret);
519 + buffer_flush(host->rrdpush_buffer);
520 + begin = 0;
521 + }
522 + else {
523 + debug(D_STREAM, "STREAM: Sent %zd bytes (part of the data buffer)...", ret);
524 + }
525 +
526 + last_sent_t = now_monotonic_sec();
527 }
528 else {
479 - debug(D_STREAM, "STREAM: Sent %zd bytes (part of the data buffer)...", ret);
529 + debug(D_STREAM, "STREAM: send() returned %zd - closing the socket...", ret);
530 + error("STREAM %s [send to %s]: failed to send metrics (send() returned %zd) - closing connection - we have sent %zu bytes on this connection.",
531 + host->hostname, connected_to, ret, sent_connection);
532 + close(host->rrdpush_socket);
533 + host->rrdpush_socket = -1;
534 }
535
482 - last_sent_t = now_monotonic_sec();
536 + debug(D_STREAM, "STREAM: Releasing exclusive lock on host...");
537 + rrdpush_unlock(host);
538 +
539 + if (pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
540 + error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
541 +
542 + // END RRDPUSH LOCKED SESSION
543 }
544 else {
485 - debug(D_STREAM, "STREAM: send() returned %zd - closing the socket...", ret);
486 - error("STREAM %s [send to %s]: failed to send metrics (send() returned %zd) - closing connection - we have sent %zu bytes on this connection.",
487 - host->hostname, connected_to, ret, sent_connection);
488 - close(host->rrdpush_socket);
489 - host->rrdpush_socket = -1;
545 + debug(D_STREAM, "STREAM: we have sent the entire buffer, but we received POLLOUT...");
546 }
547 + }
548
492 - debug(D_STREAM, "STREAM: Releasing exclusive lock on host...");
493 - rrdpush_unlock(host);
494 -
495 - if (pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
496 - error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
497 -
498 - // END RRDPUSH LOCKED SESSION
549 + if(unlikely(ofd->revents & POLLERR)) {
550 + debug(D_STREAM, "STREAM: Send failed (POLLERR) - closing socket...");
551 + 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);
552 + close(host->rrdpush_socket);
553 + host->rrdpush_socket = -1;
554 }
500 - else {
501 - debug(D_STREAM, "STREAM: we have sent the entire buffer, but we received POLLOUT...");
555 + else if(unlikely(ofd->revents & POLLHUP)) {
556 + debug(D_STREAM, "STREAM: Send failed (POLLHUP) - closing socket...");
557 + 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);
558 + close(host->rrdpush_socket);
559 + host->rrdpush_socket = -1;
560 + }
561 + else if(unlikely(ofd->revents & POLLNVAL)) {
562 + debug(D_STREAM, "STREAM: Send failed (POLLNVAL) - closing socket...");
563 + 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);
564 + close(host->rrdpush_socket);
565 + host->rrdpush_socket = -1;
566 }
567 }
504 -
505 - if(unlikely(ofd->revents & POLLERR)) {
506 - debug(D_STREAM, "STREAM: Send failed (POLLERR) - closing socket...");
507 - 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);
508 - close(host->rrdpush_socket);
509 - host->rrdpush_socket = -1;
510 - }
511 - else if(unlikely(ofd->revents & POLLHUP)) {
512 - debug(D_STREAM, "STREAM: Send failed (POLLHUP) - closing socket...");
513 - 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);
514 - close(host->rrdpush_socket);
515 - host->rrdpush_socket = -1;
516 - }
517 - else if(unlikely(ofd->revents & POLLNVAL)) {
518 - debug(D_STREAM, "STREAM: Send failed (POLLNVAL) - closing socket...");
519 - 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);
520 - close(host->rrdpush_socket);
521 - host->rrdpush_socket = -1;
568 + else {
569 + debug(D_STREAM, "STREAM: poll() timed out.");
570 }
523 - }
524 - else {
525 - debug(D_STREAM, "STREAM: poll() timed out.");
526 - }
571
528 - // protection from overflow
529 - if(buffer_strlen(host->rrdpush_buffer) > max_size) {
530 - debug(D_STREAM, "STREAM: Buffer is too big (%zu bytes), bigger than the max (%zu) - flushing it...", buffer_strlen(host->rrdpush_buffer), max_size);
531 - errno = 0;
532 - 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);
533 - if(host->rrdpush_socket != -1) {
534 - close(host->rrdpush_socket);
535 - host->rrdpush_socket = -1;
572 + // protection from overflow
573 + if(buffer_strlen(host->rrdpush_buffer) > max_size) {
574 + debug(D_STREAM, "STREAM: Buffer is too big (%zu bytes), bigger than the max (%zu) - flushing it...", buffer_strlen(host->rrdpush_buffer), max_size);
575 + errno = 0;
576 + 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);
577 + if(host->rrdpush_socket != -1) {
578 + close(host->rrdpush_socket);
579 + host->rrdpush_socket = -1;
580 + }
581 }
582 }
538 - }
539 -
540 -cleanup:
541 - debug(D_WEB_CLIENT, "STREAM %s [send]: sending thread exits.", host->hostname);
583
543 - if(pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL) != 0)
544 - error("STREAM %s [send]: cannot set pthread cancel state to DISABLE.", host->hostname);
545 -
546 - rrdpush_lock(host);
547 - rrdhost_wrlock(host);
548 - rrdpush_sender_thread_cleanup_locked_all(host);
549 - rrdhost_unlock(host);
550 - rrdpush_unlock(host);
551 -
552 - if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
553 - error("STREAM %s [send]: cannot set pthread cancel state to ENABLE.", host->hostname);
584 + pthread_cleanup_pop(1);
585
586 pthread_exit(NULL);
587 return NULL;
@@ -772,11 +803,8 @@ static void rrdpush_sender_thread_spawn(RRDHOST *host) {
803 if(!host->rrdpush_spawn) {
804 if(pthread_create(&host->rrdpush_thread, NULL, rrdpush_sender_thread, (void *) host))
805 error("STREAM %s [send]: failed to create new thread for client.", host->hostname);
775 -
776 - else if(pthread_detach(host->rrdpush_thread))
777 - error("STREAM %s [send]: cannot request detach newly created thread.", host->hostname);
778 -
779 - host->rrdpush_spawn = 1;
806 + else
807 + host->rrdpush_spawn = 1;
808 }
809
810 rrdhost_unlock(host);