@cryptotaxi247 / netdata-1 / commits / 559e1e697

allow STREAM receiver to unblock the single threaded web server

Costa Tsaousis committed Feb 23, 2017 at 17:44 UTC 559e1e69738652813a432b48ea6814f6ea31c6fb
4 files changed +237 -176
src/main.c
+2 -2
@@ -55,7 +55,7 @@ struct netdata_static_thread static_threads[] = {
55 {"plugins.d", NULL, NULL, 1, NULL, NULL, pluginsd_main},
56 {"web", NULL, NULL, 1, NULL, NULL, socket_listen_main_multi_threaded},
57 {"web-single-threaded", NULL, NULL, 0, NULL, NULL, socket_listen_main_single_threaded},
58 - {"central-netdata-push",NULL, NULL, 0, NULL, NULL, central_netdata_push_thread},
58 + {"central-netdata-push",NULL, NULL, 0, NULL, NULL, rrdpush_sender_thread},
59 {NULL, NULL, NULL, 0, NULL, NULL, NULL}
60 };
61
@@ -82,7 +82,7 @@ void web_server_threading_selection(void) {
82 if(static_threads[i].start_routine == socket_listen_main_single_threaded)
83 static_threads[i].enabled = single_threaded;
84
85 - if(static_threads[i].start_routine == central_netdata_push_thread)
85 + if(static_threads[i].start_routine == rrdpush_sender_thread)
86 static_threads[i].enabled = rrdpush_thread;
87
88 if(static_threads[i].start_routine == backends_main)
src/rrdpush.c
+231 -1
@@ -179,7 +179,7 @@ int rrdpush_init() {
179 return rrdpush_enabled;
180 }
181
182 -void *central_netdata_push_thread(void *ptr) {
182 +void *rrdpush_sender_thread(void *ptr) {
183 struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
184
185 info("STREAM: central netdata push thread created with task id %d", gettid());
@@ -394,3 +394,233 @@ cleanup:
394 pthread_exit(NULL);
395 return NULL;
396 }
397 +
398 +
399 +// ----------------------------------------------------------------------------
400 +// STREAM receiver
401 +
402 +int rrdpush_receive(int fd, const char *key, const char *hostname, const char *machine_guid, const char *os, int update_every, char *client_ip, char *client_port) {
403 + RRDHOST *host;
404 + int history = default_rrd_history_entries;
405 + RRD_MEMORY_MODE mode = default_rrd_memory_mode;
406 + int health_enabled = default_health_enabled;
407 +
408 + update_every = (int)appconfig_get_number(&stream_config, machine_guid, "update every", update_every);
409 + if(update_every < 0) update_every = 1;
410 +
411 + history = (int)appconfig_get_number(&stream_config, key, "default history", history);
412 + history = (int)appconfig_get_number(&stream_config, machine_guid, "history", history);
413 + if(history < 5) history = 5;
414 +
415 + mode = rrd_memory_mode_id(appconfig_get(&stream_config, key, "default memory mode", rrd_memory_mode_name(mode)));
416 + mode = rrd_memory_mode_id(appconfig_get(&stream_config, machine_guid, "memory mode", rrd_memory_mode_name(mode)));
417 +
418 + health_enabled = appconfig_get_boolean_ondemand(&stream_config, key, "health enabled by default", health_enabled);
419 + health_enabled = appconfig_get_boolean_ondemand(&stream_config, machine_guid, "health enabled", health_enabled);
420 +
421 + if(!strcmp(machine_guid, "localhost"))
422 + host = localhost;
423 + else
424 + host = rrdhost_find_or_create(hostname, machine_guid, os, update_every, history, mode, health_enabled?1:0);
425 +
426 + info("STREAM request from client '%s:%s' for host '%s' with machine_guid '%s': update every = %d, history = %d, memory mode = %s, health %s",
427 + client_ip, client_port,
428 + hostname, machine_guid,
429 + update_every,
430 + history,
431 + rrd_memory_mode_name(mode),
432 + (health_enabled == CONFIG_BOOLEAN_NO)?"disabled":((health_enabled == CONFIG_BOOLEAN_YES)?"enabled":"auto")
433 + );
434 +
435 + struct plugind cd = {
436 + .enabled = 1,
437 + .update_every = default_rrd_update_every,
438 + .pid = 0,
439 + .serial_failures = 0,
440 + .successful_collections = 0,
441 + .obsolete = 0,
442 + .started_t = now_realtime_sec(),
443 + .next = NULL,
444 + };
445 +
446 + // put the client IP and port into the buffers used by plugins.d
447 + snprintfz(cd.id, CONFIG_MAX_NAME, "%s:%s", client_ip, client_port);
448 + snprintfz(cd.filename, FILENAME_MAX, "%s:%s", client_ip, client_port);
449 + snprintfz(cd.fullfilename, FILENAME_MAX, "%s:%s", client_ip, client_port);
450 + snprintfz(cd.cmd, PLUGINSD_CMD_MAX, "%s:%s", client_ip, client_port);
451 +
452 + info("STREAM [%s]:%s: sending STREAM to initiate streaming...", client_ip, client_port);
453 + if(send_timeout(fd, "STREAM", 6, 0, 60) != 6) {
454 + error("STREAM [%s]:%s: cannot send STREAM.", client_ip, client_port);
455 + return 0;
456 + }
457 +
458 + // remove the non-blocking flag from the socket
459 + if(fcntl(fd, F_SETFL, fcntl(fd, F_GETFL, 0) & ~O_NONBLOCK) == -1)
460 + error("STREAM [%s]:%s: cannot remove the non-blocking flag from socket %d", client_ip, client_port, fd);
461 +
462 + // convert the socket to a FILE *
463 + FILE *fp = fdopen(fd, "r");
464 + if(!fp) {
465 + error("STREAM [%s]:%s: failed to get a FILE for FD %d.", client_ip, client_port, fd);
466 + return 0;
467 + }
468 +
469 + rrdhost_wrlock(host);
470 + host->use_counter++;
471 + rrdhost_unlock(host);
472 +
473 + // call the plugins.d processor to receive the metrics
474 + info("STREAM [%s]:%s: connecting client to plugins.d on host '%s' with machine GUID '%s'.", client_ip, client_port, host->hostname, host->machine_guid);
475 + size_t count = pluginsd_process(host, &cd, fp, 1);
476 + error("STREAM [%s]:%s: client disconnected (host '%s', machine GUID '%s').", client_ip, client_port, host->hostname, host->machine_guid);
477 +
478 + rrdhost_wrlock(host);
479 + host->use_counter--;
480 + if(!host->use_counter && health_enabled == CONFIG_BOOLEAN_AUTO)
481 + host->health_enabled = 0;
482 + rrdhost_unlock(host);
483 +
484 + // cleanup
485 + fclose(fp);
486 +
487 + return (int)count;
488 +}
489 +
490 +struct rrdpush_thread {
491 + int fd;
492 + char *key;
493 + char *hostname;
494 + char *machine_guid;
495 + char *os;
496 + char *client_ip;
497 + char *client_port;
498 + int update_every;
499 +};
500 +
501 +void *rrdpush_receiver_thread(void *ptr) {
502 + struct rrdpush_thread *rpt = (struct rrdpush_thread *)ptr;
503 +
504 + info("STREAM: central netdata receive thread created with task id %d, for client [%s]:%s", gettid(), rpt->client_ip, rpt->client_port);
505 +
506 + if (pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
507 + error("STREAM: cannot set pthread cancel type to DEFERRED.");
508 +
509 + if (pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
510 + error("STREAM: cannot set pthread cancel state to ENABLE.");
511 +
512 +
513 + rrdpush_receive(rpt->fd, rpt->key, rpt->hostname, rpt->machine_guid, rpt->os, rpt->update_every, rpt->client_ip, rpt->client_port);
514 +
515 + close(rpt->fd);
516 + freez(rpt->key);
517 + freez(rpt->hostname);
518 + freez(rpt->machine_guid);
519 + freez(rpt->os);
520 + freez(rpt->client_ip);
521 + freez(rpt->client_port);
522 + freez(rpt);
523 +
524 + pthread_exit(NULL);
525 + return NULL;
526 +}
527 +
528 +static inline int rrdpush_receive_validate_api_key(const char *key) {
529 + return appconfig_get_boolean(&stream_config, key, "enabled", 0);
530 +}
531 +
532 +int rrdpush_receiver_thread_spawn(RRDHOST *host, struct web_client *w, char *url) {
533 + (void)host;
534 +
535 + info("STREAM [%s]:%s: client connection.", w->client_ip, w->client_port);
536 +
537 + char *key = NULL, *hostname = NULL, *machine_guid = NULL, *os = NULL;
538 + int update_every = default_rrd_update_every;
539 +
540 + while(url) {
541 + char *value = mystrsep(&url, "?&");
542 + if(!value || !*value) continue;
543 +
544 + char *name = mystrsep(&value, "=");
545 + if(!name || !*name) continue;
546 + if(!value || !*value) continue;
547 +
548 + if(!strcmp(name, "key"))
549 + key = value;
550 + else if(!strcmp(name, "hostname"))
551 + hostname = value;
552 + else if(!strcmp(name, "machine_guid"))
553 + machine_guid = value;
554 + else if(!strcmp(name, "update_every"))
555 + update_every = (int)strtoul(value, NULL, 0);
556 + else if(!strcmp(name, "os"))
557 + os = value;
558 + }
559 +
560 + if(!key || !*key) {
561 + error("STREAM [%s]:%s: request without an API key. Forbidding access.", w->client_ip, w->client_port);
562 + buffer_flush(w->response.data);
563 + buffer_sprintf(w->response.data, "You need an API key for this request.");
564 + return 401;
565 + }
566 +
567 + if(!hostname || !*hostname) {
568 + error("STREAM [%s]:%s: request without a hostname. Forbidding access.", w->client_ip, w->client_port);
569 + buffer_flush(w->response.data);
570 + buffer_sprintf(w->response.data, "You need to send a hostname too.");
571 + return 400;
572 + }
573 +
574 + if(!machine_guid || !*machine_guid) {
575 + error("STREAM [%s]:%s: request without a machine GUID. Forbidding access.", w->client_ip, w->client_port);
576 + buffer_flush(w->response.data);
577 + buffer_sprintf(w->response.data, "You need to send a machine GUID too.");
578 + return 400;
579 + }
580 +
581 + if(!rrdpush_receive_validate_api_key(key)) {
582 + error("STREAM [%s]:%s: API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
583 + buffer_flush(w->response.data);
584 + buffer_sprintf(w->response.data, "Your API key is not permitted access.");
585 + return 401;
586 + }
587 +
588 + if(!appconfig_get_boolean(&stream_config, machine_guid, "enabled", 1)) {
589 + error("STREAM [%s]:%s: machine GUID '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, machine_guid);
590 + buffer_flush(w->response.data);
591 + buffer_sprintf(w->response.data, "Your machine guide is not permitted access.");
592 + return 404;
593 + }
594 +
595 + struct rrdpush_thread *rpt = mallocz(sizeof(struct rrdpush_thread));
596 + rpt->fd = w->ifd;
597 + rpt->key = strdupz(key);
598 + rpt->hostname = strdupz(hostname);
599 + rpt->machine_guid = strdupz(machine_guid);
600 + rpt->os = strdupz(os);
601 + rpt->client_ip = strdupz(w->client_ip);
602 + rpt->client_port = strdupz(w->client_port);
603 + rpt->update_every = update_every;
604 +
605 + pthread_t *thread = mallocz(sizeof(pthread_t));
606 + pthread_attr_t attr;
607 +
608 + debug(D_SYSTEM, "Starting STREAM thread for client [%s]:%s.", w->client_ip, w->client_port);
609 +
610 + if(pthread_create(thread, &attr, rrdpush_receiver_thread, rpt))
611 + error("failed to create new STREAM thread for client [%s]:%s.", w->client_ip, w->client_port);
612 +
613 + else if(pthread_detach(*thread))
614 + error("Cannot request detach newly created thread for client [%s]:%s.", w->client_ip, w->client_port);
615 +
616 + rrdpush_receive(w->ifd, key, hostname, machine_guid, os, update_every, w->client_ip, w->client_port);
617 +
618 + // prevent the caller from closing the streaming socket
619 + if(w->ifd == w->ofd)
620 + w->ifd = w->ofd = -1;
621 + else
622 + w->ifd = -1;
623 +
624 + buffer_flush(w->response.data);
625 + return 200;
626 +}
src/rrdpush.h
+3 -1
@@ -6,6 +6,8 @@ extern int rrdpush_exclusive;
6
7 extern int rrdpush_init();
8 extern void rrdset_done_push(RRDSET *st);
9 -extern void *central_netdata_push_thread(void *ptr);
9 +extern void *rrdpush_sender_thread(void *ptr);
10 +
11 +extern int rrdpush_receiver_thread_spawn(RRDHOST *host, struct web_client *w, char *url);
12
13 #endif //NETDATA_RRDPUSH_H
src/web_client.c
+1 -172
@@ -1665,177 +1665,6 @@ int web_client_api_old_data_request(RRDHOST *host, struct web_client *w, char *u
1665 return 200;
1666 }
1667
1668 -int validate_stream_api_key(const char *key) {
1669 - if(appconfig_get_boolean(&stream_config, key, "enabled", 0))
1670 - return 1;
1671 -
1672 - return 0;
1673 -}
1674 -
1675 -int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1676 - info("STREAM [%s]:%s: client connection.", w->client_ip, w->client_port);
1677 -
1678 - char *key = NULL, *hostname = NULL, *machine_guid = NULL, *os = NULL;
1679 - int update_every = default_rrd_update_every;
1680 - int history = default_rrd_history_entries;
1681 - RRD_MEMORY_MODE mode = default_rrd_memory_mode;
1682 - int health_enabled = default_health_enabled;
1683 -
1684 - while(url) {
1685 - char *value = mystrsep(&url, "?&");
1686 - if(!value || !*value) continue;
1687 -
1688 - char *name = mystrsep(&value, "=");
1689 - if(!name || !*name) continue;
1690 - if(!value || !*value) continue;
1691 -
1692 - if(!strcmp(name, "key"))
1693 - key = value;
1694 - else if(!strcmp(name, "hostname"))
1695 - hostname = value;
1696 - else if(!strcmp(name, "machine_guid"))
1697 - machine_guid = value;
1698 - else if(!strcmp(name, "update_every"))
1699 - update_every = (int)strtoul(value, NULL, 0);
1700 - else if(!strcmp(name, "os"))
1701 - os = value;
1702 - }
1703 -
1704 - if(!key || !*key) {
1705 - error("STREAM [%s]:%s: request without an API key. Forbidding access.", w->client_ip, w->client_port);
1706 - buffer_flush(w->response.data);
1707 - buffer_sprintf(w->response.data, "You need an API key for this request.");
1708 - return 401;
1709 - }
1710 -
1711 - if(!hostname || !*hostname) {
1712 - error("STREAM [%s]:%s: request without a hostname. Forbidding access.", w->client_ip, w->client_port);
1713 - buffer_flush(w->response.data);
1714 - buffer_sprintf(w->response.data, "You need to send a hostname too.");
1715 - return 400;
1716 - }
1717 -
1718 - if(!machine_guid || !*machine_guid) {
1719 - error("STREAM [%s]:%s: request without a machine GUID. Forbidding access.", w->client_ip, w->client_port);
1720 - buffer_flush(w->response.data);
1721 - buffer_sprintf(w->response.data, "You need to send a machine GUID too.");
1722 - return 400;
1723 - }
1724 -
1725 - if(!validate_stream_api_key(key)) {
1726 - error("STREAM [%s]:%s: API key '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, key);
1727 - buffer_flush(w->response.data);
1728 - buffer_sprintf(w->response.data, "Your API key is not permitted access.");
1729 - return 401;
1730 - }
1731 -
1732 - if(!appconfig_get_boolean(&stream_config, machine_guid, "enabled", 1)) {
1733 - error("STREAM [%s]:%s: machine GUID '%s' is not allowed. Forbidding access.", w->client_ip, w->client_port, machine_guid);
1734 - buffer_flush(w->response.data);
1735 - buffer_sprintf(w->response.data, "Your machine guide is not permitted access.");
1736 - return 404;
1737 - }
1738 -
1739 - update_every = (int)appconfig_get_number(&stream_config, machine_guid, "update every", update_every);
1740 - if(update_every < 0) update_every = 1;
1741 -
1742 - history = (int)appconfig_get_number(&stream_config, key, "default history", history);
1743 - history = (int)appconfig_get_number(&stream_config, machine_guid, "history", history);
1744 - if(history < 5) history = 5;
1745 -
1746 - mode = rrd_memory_mode_id(appconfig_get(&stream_config, key, "default memory mode", rrd_memory_mode_name(mode)));
1747 - mode = rrd_memory_mode_id(appconfig_get(&stream_config, machine_guid, "memory mode", rrd_memory_mode_name(mode)));
1748 -
1749 - health_enabled = appconfig_get_boolean_ondemand(&stream_config, key, "health enabled by default", health_enabled);
1750 - health_enabled = appconfig_get_boolean_ondemand(&stream_config, machine_guid, "health enabled", health_enabled);
1751 -
1752 - if(!strcmp(machine_guid, "localhost"))
1753 - host = localhost;
1754 - else
1755 - host = rrdhost_find_or_create(hostname, machine_guid, os, update_every, history, mode, health_enabled?1:0);
1756 -
1757 - info("STREAM request from client '%s:%s' for host '%s' with machine_guid '%s': update every = %d, history = %d, memory mode = %s, health %s",
1758 - w->client_ip, w->client_port,
1759 - hostname, machine_guid,
1760 - update_every,
1761 - history,
1762 - rrd_memory_mode_name(mode),
1763 - (health_enabled == CONFIG_BOOLEAN_NO)?"disabled":((health_enabled == CONFIG_BOOLEAN_YES)?"enabled":"auto")
1764 - );
1765 -
1766 - struct plugind cd = {
1767 - .enabled = 1,
1768 - .update_every = default_rrd_update_every,
1769 - .pid = 0,
1770 - .serial_failures = 0,
1771 - .successful_collections = 0,
1772 - .obsolete = 0,
1773 - .started_t = now_realtime_sec(),
1774 - .next = NULL,
1775 - };
1776 -
1777 - // put the client IP and port into the buffers used by plugins.d
1778 - snprintfz(cd.id, CONFIG_MAX_NAME, "%s:%s", w->client_ip, w->client_port);
1779 - snprintfz(cd.filename, FILENAME_MAX, "%s:%s", w->client_ip, w->client_port);
1780 - snprintfz(cd.fullfilename, FILENAME_MAX, "%s:%s", w->client_ip, w->client_port);
1781 - snprintfz(cd.cmd, PLUGINSD_CMD_MAX, "%s:%s", w->client_ip, w->client_port);
1782 -
1783 - info("STREAM [%s]:%s: sending STREAM to initiate streaming...", w->client_ip, w->client_port);
1784 - if(send_timeout(w->ifd, "STREAM", 6, 0, 60) != 6) {
1785 - error("STREAM [%s]:%s: cannot send STREAM.", w->client_ip, w->client_port);
1786 - buffer_flush(w->response.data);
1787 - buffer_sprintf(w->response.data, "Failed to reply back with STREAM");
1788 - return 400;
1789 - }
1790 -
1791 - // remove the non-blocking flag from the socket
1792 - if(fcntl(w->ifd, F_SETFL, fcntl(w->ifd, F_GETFL, 0) & ~O_NONBLOCK) == -1)
1793 - error("STREAM [%s]:%s: cannot remove the non-blocking flag from socket %d", w->client_ip, w->client_port, w->ifd);
1794 -
1795 - /*
1796 - char buffer[1000 + 1];
1797 - ssize_t len;
1798 - while((len = read(w->ifd, buffer, 1000)) != -1) {
1799 - buffer[len] = '\0';
1800 - fprintf(stderr, "BEGIN READ %zu bytes\n%s\nEND READ\n", (size_t)len, buffer);
1801 - }
1802 - */
1803 -
1804 - // convert the socket to a FILE *
1805 - FILE *fp = fdopen(w->ifd, "r");
1806 - if(!fp) {
1807 - error("STREAM [%s]:%s: failed to get a FILE for FD %d.", w->client_ip, w->client_port, w->ifd);
1808 - buffer_flush(w->response.data);
1809 - buffer_sprintf(w->response.data, "Failed to get a FILE for an FD.");
1810 - return 500;
1811 - }
1812 -
1813 - rrdhost_wrlock(host);
1814 - host->use_counter++;
1815 - rrdhost_unlock(host);
1816 -
1817 - // call the plugins.d processor to receive the metrics
1818 - info("STREAM [%s]:%s: connecting client to plugins.d on host '%s' with machine GUID '%s'.", w->client_ip, w->client_port, host->hostname, host->machine_guid);
1819 - size_t count = pluginsd_process(host, &cd, fp, 1);
1820 - error("STREAM [%s]:%s: client disconnected (host '%s', machine GUID '%s').", w->client_ip, w->client_port, host->hostname, host->machine_guid);
1821 -
1822 - rrdhost_wrlock(host);
1823 - host->use_counter--;
1824 - if(!host->use_counter && health_enabled == CONFIG_BOOLEAN_AUTO)
1825 - host->health_enabled = 0;
1826 - rrdhost_unlock(host);
1827 -
1828 - // cleanup
1829 - fclose(fp);
1830 - w->ifd = -1;
1831 -
1832 - // this will not send anything
1833 - // the socket is closed
1834 - buffer_flush(w->response.data);
1835 - if(count) return 200;
1836 - return 400;
1837 -}
1838 -
1668 const char *web_content_type_to_string(uint8_t contenttype) {
1669 switch(contenttype) {
1670 case CT_TEXT_HTML:
@@ -2485,7 +2314,7 @@ void web_client_process_request(struct web_client *w) {
2314 case HTTP_VALIDATION_OK:
2315 switch(w->mode) {
2316 case WEB_CLIENT_MODE_STREAM:
2488 - w->response.code = web_client_stream_request(localhost, w, w->decoded_url);
2317 + w->response.code = rrdpush_receiver_thread_spawn(localhost, w, w->decoded_url);
2318 return;
2319
2320 case WEB_CLIENT_MODE_OPTIONS: