@cryptotaxi247 / netdata-1 / commits / 32afa0289

skeleton for a statsd plugin for netdata - not finished yet

Costa Tsaousis (ktsaou) committed Apr 23, 2017 at 17:28 UTC 32afa028926a89a8479d3bcb54e100d8473736bc
14 files changed +589 -68
CMakeLists.txt
+1 -1
@@ -151,7 +151,7 @@ set(NETDATA_SOURCE_FILES
151 src/web_client.h
152 src/web_server.c
153 src/web_server.h
154 - src/locks.h)
154 + src/locks.h src/statsd.c src/statsd.h)
155
156 set(APPS_PLUGIN_SOURCE_FILES
157 src/appconfig.c
src/Makefile.am
+2
@@ -119,6 +119,8 @@ netdata_SOURCES = \
119 simple_pattern.h \
120 socket.c \
121 socket.h \
122 + statsd.c \
123 + statsd.h \
124 storage_number.c \
125 storage_number.h \
126 sys_devices_system_edac_mc.c \
src/appconfig.c
+1
@@ -532,6 +532,7 @@ void appconfig_generate(struct config *root, BUFFER *wb, int only_changed)
532 for(co = root->sections; co ; co = co->next) {
533 if(!strcmp(co->name, CONFIG_SECTION_GLOBAL)
534 || !strcmp(co->name, CONFIG_SECTION_WEB)
535 + || !strcmp(co->name, CONFIG_SECTION_STATSD)
536 || !strcmp(co->name, CONFIG_SECTION_PLUGINS)
537 || !strcmp(co->name, CONFIG_SECTION_REGISTRY)
538 || !strcmp(co->name, CONFIG_SECTION_HEALTH)
src/appconfig.h
+1
@@ -5,6 +5,7 @@
5
6 #define CONFIG_SECTION_GLOBAL "global"
7 #define CONFIG_SECTION_WEB "web"
8 +#define CONFIG_SECTION_STATSD "statsd"
9 #define CONFIG_SECTION_PLUGINS "plugins"
10 #define CONFIG_SECTION_REGISTRY "registry"
11 #define CONFIG_SECTION_HEALTH "health"
src/common.h
+2
@@ -40,6 +40,7 @@
40 #include <strings.h>
41 #include <arpa/inet.h>
42 #include <netinet/tcp.h>
43 +#include <sys/ioctl.h>
44
45 #ifdef HAVE_NETINET_IN_H
46 #include <netinet/in.h>
@@ -207,6 +208,7 @@
208 #include "rrd.h"
209 #include "plugin_tc.h"
210 #include "plugins_d.h"
211 +#include "statsd.h"
212 #include "rrd2json.h"
213 #include "rrd2json_api_old.h"
214 #include "web_client.h"
src/log.h
+2
@@ -29,6 +29,8 @@
29 #define D_RRDHOST 0x0000000002000000
30 #define D_LOCKS 0x0000000004000000
31 #define D_BACKEND 0x0000000008000000
32 +#define D_STATSD 0x0000000010000000
33 +#define D_POLLFD 0x0000000020000000
34 #define D_SYSTEM 0x8000000000000000
35
36 //#define DEBUG (D_WEB_CLIENT_ACCESS|D_LISTENER|D_RRD_STATS)
src/main.c
+1
@@ -56,6 +56,7 @@ struct netdata_static_thread static_threads[] = {
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 {"push-metrics", NULL, NULL, 0, NULL, NULL, rrdpush_sender_thread},
59 + {"statsd", NULL, NULL, 1, NULL, NULL, statsd_main},
60 {NULL, NULL, NULL, 0, NULL, NULL, NULL}
61 };
62
src/socket.c
+351 -29
@@ -6,18 +6,23 @@
6 int create_listen_socket4(int socktype, const char *ip, int port, int listen_backlog) {
7 int sock;
8 int sockopt = 1;
9 + int nonblock = 1;
10
10 - debug(D_LISTENER, "IPv4 creating new listening socket on ip '%s' port %d", ip, port);
11 + debug(D_LISTENER, "LISTENER: IPv4 creating new listening socket on ip '%s' port %d, socktype %d", ip, port, socktype);
12
13 sock = socket(AF_INET, socktype, 0);
14 if(sock < 0) {
14 - error("IPv4 socket() on ip '%s' port %d failed.", ip, port);
15 + error("LISTENER: IPv4 socket() on ip '%s' port %d, socktype %d failed.", ip, port, socktype);
16 return -1;
17 }
18
19 /* avoid "address already in use" */
20 if(setsockopt(sock, SOL_SOCKET, SO_REUSEADDR, (void*)&sockopt, sizeof(sockopt)) != 0)
20 - error("Cannot set SO_REUSEADDR on ip '%s' port's %d.", ip, port);
21 + error("LISTENER: Cannot set SO_REUSEADDR on ip '%s' port %d, socktype %d.", ip, port, socktype);
22 +
23 + /* set non-blocking mode */
24 + if(ioctl(sock, FIONBIO, (char *)&nonblock))
25 + error("LISTENER: Cannot set FIONBIO on ip '%s' port %d, socktype %d.", ip, port, socktype);
26
27 struct sockaddr_in name;
28 memset(&name, 0, sizeof(struct sockaddr_in));
@@ -26,24 +31,24 @@ int create_listen_socket4(int socktype, const char *ip, int port, int listen_bac
31
32 int ret = inet_pton(AF_INET, ip, (void *)&name.sin_addr.s_addr);
33 if(ret != 1) {
29 - error("Failed to convert IP '%s' to a valid IPv4 address.", ip);
34 + error("LISTENER: Failed to convert IP '%s' to a valid IPv4 address.", ip);
35 close(sock);
36 return -1;
37 }
38
39 if(bind (sock, (struct sockaddr *) &name, sizeof (name)) < 0) {
40 close(sock);
36 - error("IPv4 bind() on ip '%s' port %d failed.", ip, port);
41 + error("LISTENER: IPv4 bind() on ip '%s' port %d, socktype %d failed.", ip, port, socktype);
42 return -1;
43 }
44
40 - if(listen(sock, listen_backlog) < 0) {
45 + if(socktype == SOCK_STREAM && listen(sock, listen_backlog) < 0) {
46 close(sock);
42 - error("IPv4 listen() on ip '%s' port %d failed.", ip, port);
47 + error("LISTENER: IPv4 listen() on ip '%s' port %d, socktype %d failed.", ip, port, socktype);
48 return -1;
49 }
50
46 - debug(D_LISTENER, "Listening on IPv4 ip '%s' port %d", ip, port);
51 + debug(D_LISTENER, "LISTENER: Listening on IPv4 ip '%s' port %d, socktype %d", ip, port, socktype);
52 return sock;
53 }
54
@@ -51,22 +56,27 @@ int create_listen_socket6(int socktype, uint32_t scope_id, const char *ip, int p
56 int sock = -1;
57 int sockopt = 1;
58 int ipv6only = 1;
59 + int nonblock = 1;
60
55 - debug(D_LISTENER, "IPv6 creating new listening socket on ip '%s' port %d", ip, port);
61 + debug(D_LISTENER, "LISTENER: IPv6 creating new listening socket on ip '%s' port %d, socktype %d", ip, port, socktype);
62
63 sock = socket(AF_INET6, socktype, 0);
64 if (sock < 0) {
59 - error("IPv6 socket() on ip '%s' port %d failed.", ip, port);
65 + error("LISTENER: IPv6 socket() on ip '%s' port %d, socktype %d, failed.", ip, port, socktype);
66 return -1;
67 }
68
69 /* avoid "address already in use" */
70 if(setsockopt(sock, SOL_SOCKET, SO_REUSEADDR, (void*)&sockopt, sizeof(sockopt)) != 0)
65 - error("Cannot set SO_REUSEADDR on ip '%s' port's %d.", ip, port);
71 + error("LISTENER: Cannot set SO_REUSEADDR on ip '%s' port %d, socktype %d.", ip, port, socktype);
72 +
73 + /* set non-blocking mode */
74 + if(ioctl(sock, FIONBIO, (char *)&nonblock))
75 + error("LISTENER: Cannot set FIONBIO on ip '%s' port %d, socktype %d.", ip, port, socktype);
76
77 /* IPv6 only */
78 if(setsockopt(sock, IPPROTO_IPV6, IPV6_V6ONLY, (void*)&ipv6only, sizeof(ipv6only)) != 0)
69 - error("Cannot set IPV6_V6ONLY on ip '%s' port's %d.", ip, port);
79 + error("LISTENER: Cannot set IPV6_V6ONLY on ip '%s' port %d, socktype %d.", ip, port, socktype);
80
81 struct sockaddr_in6 name;
82 memset(&name, 0, sizeof(struct sockaddr_in6));
@@ -76,7 +86,7 @@ int create_listen_socket6(int socktype, uint32_t scope_id, const char *ip, int p
86
87 int ret = inet_pton(AF_INET6, ip, (void *)&name.sin6_addr.s6_addr);
88 if(ret != 1) {
79 - error("Failed to convert IP '%s' to a valid IPv6 address.", ip);
89 + error("LISTENER: Failed to convert IP '%s' to a valid IPv6 address.", ip);
90 close(sock);
91 return -1;
92 }
@@ -85,23 +95,23 @@ int create_listen_socket6(int socktype, uint32_t scope_id, const char *ip, int p
95
96 if (bind (sock, (struct sockaddr *) &name, sizeof (name)) < 0) {
97 close(sock);
88 - error("IPv6 bind() on ip '%s' port %d failed.", ip, port);
98 + error("LISTENER: IPv6 bind() on ip '%s' port %d, socktype %d failed.", ip, port, socktype);
99 return -1;
100 }
101
92 - if (listen(sock, listen_backlog) < 0) {
102 + if (socktype == SOCK_STREAM && listen(sock, listen_backlog) < 0) {
103 close(sock);
94 - error("IPv6 listen() on ip '%s' port %d failed.", ip, port);
104 + error("LISTENER: IPv6 listen() on ip '%s' port %d, socktype %d failed.", ip, port, socktype);
105 return -1;
106 }
107
98 - debug(D_LISTENER, "Listening on IPv6 ip '%s' port %d", ip, port);
108 + debug(D_LISTENER, "LISTENER: Listening on IPv6 ip '%s' port %d, socktype %d", ip, port, socktype);
109 return sock;
110 }
111
102 -static inline int listen_sockets_add(LISTEN_SOCKETS *sockets, int fd, const char *protocol, const char *ip, int port) {
112 +static inline int listen_sockets_add(LISTEN_SOCKETS *sockets, int fd, int socktype, const char *protocol, const char *ip, int port) {
113 if(sockets->opened >= MAX_LISTEN_FDS) {
104 - error("Too many listening sockets. Failed to add listening %s socket at ip '%s' port %d", protocol, ip, port);
114 + error("LISTENER: Too many listening sockets. Failed to add listening %s socket at ip '%s' port %d, protocol %s, socktype %d", protocol, ip, port, protocol, socktype);
115 close(fd);
116 return -1;
117 }
@@ -111,6 +121,7 @@ static inline int listen_sockets_add(LISTEN_SOCKETS *sockets, int fd, const char
121 char buffer[100 + 1];
122 snprintfz(buffer, 100, "%s:[%s]:%d", protocol, ip, port);
123 sockets->fds_names[sockets->opened] = strdupz(buffer);
124 + sockets->fds_types[sockets->opened] = socktype;
125
126 sockets->opened++;
127 return 0;
@@ -129,6 +140,7 @@ static inline void listen_sockets_init(LISTEN_SOCKETS *sockets) {
140 for(i = 0; i < MAX_LISTEN_FDS ;i++) {
141 sockets->fds[i] = -1;
142 sockets->fds_names[i] = NULL;
143 + sockets->fds_types[i] = -1;
144 }
145
146 sockets->opened = 0;
@@ -143,6 +155,8 @@ void listen_sockets_close(LISTEN_SOCKETS *sockets) {
155
156 freez(sockets->fds_names[i]);
157 sockets->fds_names[i] = NULL;
158 +
159 + sockets->fds_types[i] = -1;
160 }
161
162 sockets->opened = 0;
@@ -207,7 +221,7 @@ static inline int bind_to_one(LISTEN_SOCKETS *sockets, const char *definition, i
221 if(*interface) {
222 scope_id = if_nametoindex(interface);
223 if(!scope_id)
210 - error("Cannot find a network interface named '%s'. Continuing with limiting the network interface", interface);
224 + error("LISTENER: Cannot find a network interface named '%s'. Continuing with limiting the network interface", interface);
225 }
226
227 if(!*ip || *ip == '*' || !strcmp(ip, "any") || !strcmp(ip, "all"))
@@ -227,7 +241,7 @@ static inline int bind_to_one(LISTEN_SOCKETS *sockets, const char *definition, i
241
242 int r = getaddrinfo(ip, port, &hints, &result);
243 if (r != 0) {
230 - error("getaddrinfo('%s', '%s'): %s\n", ip, port, gai_strerror(r));
244 + error("LISTENER: getaddrinfo('%s', '%s'): %s\n", ip, port, gai_strerror(r));
245 return -1;
246 }
247
@@ -255,16 +269,16 @@ static inline int bind_to_one(LISTEN_SOCKETS *sockets, const char *definition, i
269 }
270
271 default:
258 - debug(D_LISTENER, "Unknown socket family %d", rp->ai_addr->sa_family);
272 + debug(D_LISTENER, "LISTENER: Unknown socket family %d", rp->ai_addr->sa_family);
273 break;
274 }
275
276 if (fd == -1) {
263 - error("Cannot bind to ip '%s', port %d", rip, rport);
277 + error("LISTENER: Cannot bind to ip '%s', port %d", rip, rport);
278 sockets->failed++;
279 }
280 else {
267 - listen_sockets_add(sockets, fd, protocol_str, rip, rport);
281 + listen_sockets_add(sockets, fd, socktype, protocol_str, rip, rport);
282 added++;
283 }
284 }
@@ -282,12 +296,12 @@ int listen_sockets_setup(LISTEN_SOCKETS *sockets) {
296 int old_port = sockets->default_port;
297 sockets->default_port = (int) config_get_number(sockets->config_section, "default port", sockets->default_port);
298 if(sockets->default_port < 1 || sockets->default_port > 65535) {
285 - error("Invalid listen port %d given. Defaulting to %d.", sockets->default_port, old_port);
299 + error("LISTENER: Invalid listen port %d given. Defaulting to %d.", sockets->default_port, old_port);
300 sockets->default_port = (int) config_set_number(sockets->config_section, "default port", old_port);
301 }
288 - debug(D_OPTIONS, "Default listen port set to %d.", sockets->default_port);
302 + debug(D_OPTIONS, "LISTENER: Default listen port set to %d.", sockets->default_port);
303
290 - char *s = config_get(sockets->config_section, "bind to", "*");
304 + char *s = config_get(sockets->config_section, "bind to", sockets->default_bind_to);
305 while(*s) {
306 char *e = s;
307
@@ -308,12 +322,12 @@ int listen_sockets_setup(LISTEN_SOCKETS *sockets) {
322 }
323
324 if(!sockets->opened)
311 - fatal("Cannot listen on any socket. Exiting...");
325 + fatal("LISTENER: Cannot listen on any socket. Exiting...");
326
327 else if(sockets->failed) {
328 size_t i;
329 for(i = 0; i < sockets->opened ;i++)
316 - info("Listen socket %s opened successfully.", sockets->fds_names[i]);
330 + info("LISTENER: Listen socket %s opened successfully.", sockets->fds_names[i]);
331 }
332
333 return (int)sockets->opened;
@@ -658,3 +672,311 @@ int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flags) {
672 }
673 #endif
674
675 +
676 +// --------------------------------------------------------------------------------------------------------------------
677 +// accept_socket() - accept a socket and store client IP and port
678 +
679 +int accept_socket(int fd, int flags, char *client_ip, size_t ipsize, char *client_port, size_t portsize) {
680 + struct sockaddr_storage sadr;
681 + socklen_t addrlen = sizeof(sadr);
682 +
683 + int nfd = accept4(fd, (struct sockaddr *)&sadr, &addrlen, flags);
684 + if (nfd >= 0) {
685 + if (getnameinfo((struct sockaddr *)&sadr, addrlen, client_ip, (socklen_t)ipsize, client_port, (socklen_t)portsize, NI_NUMERICHOST | NI_NUMERICSERV) != 0) {
686 + error("LISTENER: cannot getnameinfo() on received client connection.");
687 + strncpyz(client_ip, "UNKNOWN", ipsize - 1);
688 + strncpyz(client_port, "UNKNOWN", portsize - 1);
689 + }
690 +
691 + client_ip[ipsize - 1] = '\0';
692 + client_port[portsize - 1] = '\0';
693 +
694 + switch (((struct sockaddr *)&sadr)->sa_family) {
695 + case AF_INET:
696 + debug(D_LISTENER, "New IPv4 web client from %s port %s on socket %d.", client_ip, client_port, fd);
697 + break;
698 +
699 + case AF_INET6:
700 + if (strncmp(client_ip, "::ffff:", 7) == 0) {
701 + memmove(client_ip, &client_ip[7], strlen(&client_ip[7]) + 1);
702 + debug(D_LISTENER, "New IPv4 web client from %s port %s on socket %d.", client_ip, client_port, fd);
703 + } else
704 + debug(D_LISTENER, "New IPv6 web client from %s port %s on socket %d.", client_ip, client_port, fd);
705 + break;
706 +
707 + default:
708 + debug(D_LISTENER, "New UNKNOWN web client from %s port %s on socket %d.", client_ip, client_port, fd);
709 + break;
710 + }
711 + }
712 +
713 + return nfd;
714 +}
715 +
716 +
717 +// --------------------------------------------------------------------------------------------------------------------
718 +// poll() based listener
719 +// this should be the fastest possible listener for up to 100 sockets
720 +// above 100, an epoll() interface is needed on Linux
721 +
722 +#define POLL_FDS_INCREASE_STEP 10
723 +
724 +#define POLLINFO_FLAG_SERVER_SOCKET 0x00000001
725 +#define POLLINFO_FLAG_CLIENT_SOCKET 0x00000002
726 +
727 +struct pollinfo {
728 + size_t slot;
729 + char *client;
730 + struct pollinfo *next;
731 + uint32_t flags;
732 + int socktype;
733 +
734 + void *data;
735 +};
736 +
737 +struct poll {
738 + size_t slots;
739 + size_t used;
740 + size_t min;
741 + size_t max;
742 + struct pollfd *fds;
743 + struct pollinfo *inf;
744 + struct pollinfo *first_free;
745 +
746 + void *(*add_callback)(int fd, short int *events);
747 + void (*del_callback)(int fd, void *data);
748 + int (*rcv_callback)(int fd, int socktype, void *data, short int *events);
749 +};
750 +
751 +static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype, short int events, uint32_t flags) {
752 + debug(D_POLLFD, "POLLFD: ADD: request to add fd %d, slots = %zu, used = %zu, min = %zu, max = %zu, next free = %zd", fd, p->slots, p->used, p->min, p->max, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
753 +
754 + if(unlikely(fd < 0)) return NULL;
755 +
756 + if(unlikely(!p->first_free)) {
757 + size_t new_slots = p->slots + POLL_FDS_INCREASE_STEP;
758 + debug(D_POLLFD, "POLLFD: ADD: increasing size (current = %zu, new = %zu, used = %zu, min = %zu, max = %zu)", p->slots, new_slots, p->used, p->min, p->max);
759 +
760 + p->fds = reallocz(p->fds, sizeof(struct pollfd) * new_slots);
761 + p->inf = reallocz(p->inf, sizeof(struct pollinfo) * new_slots);
762 +
763 + ssize_t i;
764 + for(i = new_slots - 1; i >= (ssize_t)p->slots ; i--) {
765 + debug(D_POLLFD, "POLLFD: ADD: reseting new slot %zd", i);
766 + p->fds[i].fd = -1;
767 + p->fds[i].events = 0;
768 + p->fds[i].revents = 0;
769 +
770 + p->inf[i].slot = (size_t)i;
771 + p->inf[i].flags = 0;
772 + p->inf[i].socktype = -1;
773 + p->inf[i].client = NULL;
774 + p->inf[i].data = NULL;
775 + p->inf[i].next = p->first_free;
776 + p->first_free = &p->inf[i];
777 + }
778 +
779 + p->slots = new_slots;
780 + }
781 +
782 + struct pollinfo *pi = p->first_free;
783 + p->first_free = p->first_free->next;
784 +
785 + debug(D_POLLFD, "POLLFD: ADD: selected slot %zu, next free is %zd", pi->slot, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
786 +
787 + struct pollfd *pf = &p->fds[pi->slot];
788 + pf->fd = fd;
789 + pf->events = events;
790 + pf->revents = 0;
791 +
792 + pi->socktype = socktype;
793 + pi->flags = flags;
794 + pi->next = NULL;
795 +
796 + p->used++;
797 + if(unlikely(pi->slot > p->max))
798 + p->max = pi->slot;
799 +
800 + if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
801 + pi->data = p->add_callback(fd, &pf->events);
802 + }
803 +
804 + if(pi->flags & POLLINFO_FLAG_SERVER_SOCKET) {
805 + p->min = pi->slot;
806 + }
807 +
808 + debug(D_POLLFD, "POLLFD: ADD: completed, slots = %zu, used = %zu, min = %zu, max = %zu, next free = %zd", p->slots, p->used, p->min, p->max, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
809 +
810 + return pi;
811 +}
812 +
813 +static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
814 + struct pollfd *pf = &p->fds[pi->slot];
815 + debug(D_POLLFD, "POLLFD: DEL: request to clear slot %zu (fd %d), old next free was %zd", pi->slot, pf->fd, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
816 +
817 + if(unlikely(pf->fd == -1)) return;
818 +
819 + if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
820 + p->del_callback(pf->fd, pi->data);
821 + }
822 +
823 + close(pf->fd);
824 + pf->fd = -1;
825 + pf->events = 0;
826 + pf->revents = 0;
827 +
828 + pi->socktype = -1;
829 + pi->flags = 0;
830 + pi->data = NULL;
831 +
832 + freez(pi->client);
833 + pi->client = NULL;
834 +
835 + pi->next = p->first_free;
836 + p->first_free = pi;
837 +
838 + p->used--;
839 + if(p->max == pi->slot) {
840 + p->max = p->min;
841 + ssize_t i;
842 + for(i = (ssize_t)pi->slot; i > (ssize_t)p->min ;i--) {
843 + if (unlikely(p->fds[i].fd != -1)) {
844 + p->max = (size_t)i;
845 + break;
846 + }
847 + }
848 + }
849 +
850 + debug(D_POLLFD, "POLLFD: DEL: completed, slots = %zu, used = %zu, min = %zu, max = %zu, next free = %zd", p->slots, p->used, p->min, p->max, p->first_free?(ssize_t)p->first_free->slot:(ssize_t)-1);
851 +}
852 +
853 +void poll_events(LISTEN_SOCKETS *sockets
854 + , void *(*add_callback)(int fd, short int *events)
855 + , void (*del_callback)(int fd, void *data)
856 + , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
857 +) {
858 + int retval;
859 +
860 + struct poll p = {
861 + .slots = 0,
862 + .used = 0,
863 + .max = 0,
864 + .fds = NULL,
865 + .inf = NULL,
866 + .first_free = NULL,
867 +
868 + .add_callback = add_callback,
869 + .del_callback = del_callback,
870 + .rcv_callback = rcv_callback
871 + };
872 +
873 + size_t i;
874 + for(i = 0; i < sockets->opened ;i++) {
875 + poll_add_fd(&p, sockets->fds[i], sockets->fds_types[i], POLLIN, POLLINFO_FLAG_SERVER_SOCKET);
876 + info("POLLFD: LISTENER: listening on '%s'", (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN");
877 + }
878 +
879 + int timeout = 10 * 1000;
880 +
881 + for(;;) {
882 + if(unlikely(netdata_exit)) break;
883 +
884 + debug(D_POLLFD, "POLLFD: LISTENER: Waiting on %zu sockets...", p.max + 1);
885 + retval = poll(p.fds, p.max + 1, timeout);
886 +
887 + if(unlikely(retval == -1)) {
888 + error("POLLFD: LISTENER: poll() failed.");
889 + continue;
890 + }
891 + else if(unlikely(!retval)) {
892 + debug(D_POLLFD, "POLLFD: LISTENER: poll() timeout.");
893 + continue;
894 + }
895 +
896 + if(unlikely(netdata_exit)) break;
897 +
898 + for(i = 0 ; i <= p.max ; i++) {
899 + struct pollfd *pf = &p.fds[i];
900 + struct pollinfo *pi = &p.inf[i];
901 + int fd = pf->fd;
902 +
903 + if(unlikely(fd == -1)) {
904 + debug(D_POLLFD, "POLLFD: LISTENER: ignoring slot %zu, it does not have an fd", i);
905 + continue;
906 + }
907 +
908 + // check for new incoming connections
909 + if(pf->revents & POLLIN || pf->revents & POLLPRI) {
910 + debug(D_POLLFD, "POLLFD: LISTENER: processing events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
911 +
912 + pf->revents = 0;
913 +
914 + if(likely(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET)) {
915 + // read data from client TCP socket
916 +
917 + debug(D_POLLFD, "POLLFD: LISTENER: reading data from TCP client slot %zu (fd %d)", i, fd);
918 +
919 + if (p.rcv_callback(fd, pi->socktype, pi->data, &pf->events) == -1)
920 + poll_close_fd(&p, pi);
921 + }
922 +
923 + if(likely(pi->flags & POLLINFO_FLAG_SERVER_SOCKET)) {
924 + // new connection
925 +
926 + debug(D_POLLFD, "POLLFD: LISTENER: accepting connections from slot %zu (fd %d)", i, fd);
927 +
928 + switch(pi->socktype) {
929 + case SOCK_STREAM: {
930 + // a TCP socket
931 + // we accept the connection
932 +
933 + int nfd;
934 + do {
935 + char client_ip[NI_MAXHOST + 1];
936 + char client_port[NI_MAXSERV + 1];
937 +
938 + debug(D_POLLFD, "POLLFD: LISTENER: calling accept4() slot %zu (fd %d)", i, fd);
939 + nfd = accept_socket(fd, SOCK_NONBLOCK, client_ip, NI_MAXHOST + 1, client_port, NI_MAXSERV + 1);
940 + if (nfd < 0) {
941 + // accept failed
942 +
943 + debug(D_POLLFD, "POLLFD: LISTENER: accept4() slot %zu (fd %d) failed.", i, fd);
944 +
945 + if(errno != EWOULDBLOCK && errno != EAGAIN)
946 + error("POLLFD: LISTENER: accept() failed.");
947 +
948 + break;
949 + }
950 + else {
951 + // accept ok
952 + info("POLLFD: LISTENER: client '[%s]:%s' connected to '%s'", client_ip, client_port, sockets->fds_names[i]);
953 + poll_add_fd(&p, nfd, SOCK_STREAM, POLLIN, POLLINFO_FLAG_CLIENT_SOCKET);
954 + }
955 + } while (nfd != -1);
956 + break;
957 + }
958 +
959 + case SOCK_DGRAM: {
960 + // a UDP socket
961 + // we read data from the server socket
962 +
963 + debug(D_POLLFD, "POLLFD: LISTENER: reading data from UDP slot %zu (fd %d)", i, fd);
964 +
965 + p.rcv_callback(fd, pi->socktype, pi->data, &pf->events);
966 + break;
967 + }
968 +
969 + default: {
970 + error("POLLFD: LISTENER: Unknown socktype %d on slot %zu", pi->socktype, pi->slot);
971 + break;
972 + }
973 + }
974 + }
975 + }
976 + else {
977 + debug(D_POLLFD, "POLLFD: LISTENER: no events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
978 + }
979 + }
980 + }
981 +}
982 +
src/socket.h
+11
@@ -7,6 +7,7 @@
7
8 typedef struct listen_sockets {
9 const char *config_section; // the netdata configuration section to read settings from
10 + const char *default_bind_to; // the default bind to configuration string
11 int default_port; // the default port to use
12 int backlog; // the default listen backlog to use
13
@@ -14,6 +15,7 @@ typedef struct listen_sockets {
15 size_t failed; // the number of sockets attempted to open, but failed
16 int fds[MAX_LISTEN_FDS]; // the open sockets
17 char *fds_names[MAX_LISTEN_FDS]; // descriptions for the open sockets
18 + int fds_types[MAX_LISTEN_FDS]; // the socktype for the open sockets
19 } LISTEN_SOCKETS;
20
21 extern int listen_sockets_setup(LISTEN_SOCKETS *sockets);
@@ -25,6 +27,8 @@ extern int connect_to_one_of(const char *destination, int default_port, struct t
27 extern ssize_t recv_timeout(int sockfd, void *buf, size_t len, int flags, int timeout);
28 extern ssize_t send_timeout(int sockfd, void *buf, size_t len, int flags, int timeout);
29
30 +extern int accept_socket(int fd, int flags, char *client_ip, size_t ipsize, char *client_port, size_t portsize);
31 +
32 #ifndef HAVE_ACCEPT4
33 extern int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flags);
34
@@ -38,4 +42,11 @@ extern int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flag
42
43 #endif /* #ifndef HAVE_ACCEPT4 */
44
45 +
46 +extern void poll_events(LISTEN_SOCKETS *sockets
47 + , void *(*add_callback)(int fd, short int *events)
48 + , void (*del_callback)(int fd, void *data)
49 + , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
50 +);
51 +
52 #endif //NETDATA_SOCKET_H
src/statsd.c new
+203
@@ -0,0 +1,203 @@
1 +#include "common.h"
2 +
3 +static int statsd_threads = 0;
4 +
5 +// --------------------------------------------------------------------------------------
6 +
7 +static LISTEN_SOCKETS statsd_sockets = {
8 + .config_section = CONFIG_SECTION_STATSD,
9 + .default_bind_to = "udp:* tcp:*",
10 + .default_port = STATSD_LISTEN_PORT,
11 + .backlog = STATSD_LISTEN_BACKLOG
12 +};
13 +
14 +int statsd_listen_sockets_setup(void) {
15 + return listen_sockets_setup(&statsd_sockets);
16 +}
17 +
18 +// --------------------------------------------------------------------------------------
19 +
20 +typedef enum statsd_metric_type {
21 + STATSD_METRIC_TYPE_GAUGE = 'g',
22 + STATSD_METRIC_TYPE_COUNTER = 'c',
23 + STATSD_METRIC_TYPE_TIMER = 't',
24 + STATSD_METRIC_TYPE_HISTOGRAM = 'h',
25 + STATSD_METRIC_TYPE_METER = 'm'
26 +} STATSD_METRIC_TYPE;
27 +
28 +typedef struct statsd_metric {
29 + avl avl; // indexing
30 +
31 + const char *key; // "type|name" for indexing
32 + uint32_t hash; // hash of the key
33 +
34 + const char *name;
35 + STATSD_METRIC_TYPE type;
36 +
37 + usec_t last_collected_ut; // the last time this metric was updated
38 + usec_t last_exposed_ut; // the last time this metric was sent to netdata
39 + size_t events; // the number of times this metrics has been collected
40 +
41 + size_t count; // number of events since the last exposure to netdata
42 + calculated_number last; // the last value collected
43 + calculated_number total; // the sum of all values collected since the last exposure to netdata
44 + calculated_number min; // the min value collected since the last exposure to netdata
45 + calculated_number max; // the max value collected since the last exposure to netdata
46 +
47 + RRDSET *st;
48 + RRDDIM *rd_min;
49 + RRDDIM *rd_max;
50 + RRDDIM *rd_avg;
51 +
52 + netdata_mutex_t mutex;
53 +
54 + struct statsd_metric *next;
55 +} STATSD_METRIC;
56 +
57 +static inline void statsd_collected_value(STATSD_METRIC *mt, calculated_number value, const char *options) {
58 + (void)options;
59 +
60 + int lock = 0;
61 + if(unlikely(statsd_threads > 1)) {
62 + netdata_mutex_lock(&mt->mutex);
63 + lock = 1;
64 + }
65 +
66 + mt->last_collected_ut = now_realtime_usec();
67 + mt->events++;
68 + mt->count++;
69 +
70 + switch(mt->type) {
71 + case STATSD_METRIC_TYPE_HISTOGRAM:
72 + // FIXME: not implemented yet
73 +
74 + case STATSD_METRIC_TYPE_GAUGE:
75 + case STATSD_METRIC_TYPE_TIMER:
76 + case STATSD_METRIC_TYPE_METER:
77 + mt->last = value;
78 + mt->total += value;
79 + if(value < mt->min)
80 + mt->min = value;
81 + if(value > mt->max)
82 + mt->max = value;
83 + break;
84 +
85 + case STATSD_METRIC_TYPE_COUNTER:
86 + mt->total = mt->last = value;
87 + break;
88 + }
89 +
90 + if(unlikely(lock))
91 + netdata_mutex_unlock(&mt->mutex);
92 +}
93 +
94 +static void statsd_process(char *buffer, size_t size) {
95 + buffer[size] = '\0';
96 + debug(D_STATSD, "RECEIVED: '%s'", buffer);
97 +}
98 +
99 +static char statsd_read_buffer[65536];
100 +
101 +// new TCP client connected
102 +static void *statsd_add_callback(int fd, short int *events) {
103 + (void)fd;
104 + *events = POLLIN;
105 +
106 + return NULL;
107 +}
108 +
109 +
110 +// TCP client disconnected
111 +static void statsd_del_callback(int fd, void *data) {
112 + (void)fd;
113 + (void)data;
114 +
115 + return;
116 +}
117 +
118 +// Receive data
119 +static int statsd_rcv_callback(int fd, int socktype, void *data, short int *events) {
120 + (void)data;
121 +
122 + switch(socktype) {
123 + case SOCK_STREAM: {
124 + ssize_t rc;
125 + do {
126 + rc = recv(fd, statsd_read_buffer, sizeof(statsd_read_buffer), MSG_DONTWAIT);
127 + if (rc < 0) {
128 + // read failed
129 + if (errno != EWOULDBLOCK && errno != EAGAIN) {
130 + error("STATSD: recv() failed.");
131 + return -1;
132 + }
133 + } else if (!rc) {
134 + // connection closed
135 + error("STATSD: client disconnected.");
136 + return -1;
137 + } else {
138 + // data received
139 + statsd_process(statsd_read_buffer, (size_t) rc);
140 + }
141 + } while (rc != -1);
142 + break;
143 + }
144 +
145 + case SOCK_DGRAM: {
146 + ssize_t rc;
147 + do {
148 + // FIXME: collect sender information
149 + rc = recvfrom(fd, statsd_read_buffer, sizeof(statsd_read_buffer), MSG_DONTWAIT, NULL, NULL);
150 + if (rc < 0) {
151 + // read failed
152 + if (errno != EWOULDBLOCK && errno != EAGAIN) {
153 + error("STATSD: recvfrom() failed.");
154 + return -1;
155 + }
156 + } else if (rc) {
157 + // data received
158 + statsd_process(statsd_read_buffer, (size_t) rc);
159 + }
160 + } while (rc != -1);
161 + break;
162 + }
163 +
164 + default: {
165 + error("STATSD: unknown socktype %d on socket %d", socktype, fd);
166 + return -1;
167 + }
168 + }
169 +
170 + *events = POLLIN;
171 + return 0;
172 +}
173 +
174 +void *statsd_main(void *ptr) {
175 + struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
176 +
177 + info("STATSD thread created with task id %d", gettid());
178 +
179 + if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
180 + error("Cannot set pthread cancel type to DEFERRED.");
181 +
182 + if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
183 + error("Cannot set pthread cancel state to ENABLE.");
184 +
185 + statsd_listen_sockets_setup();
186 + if(!statsd_sockets.opened) {
187 + error("STATSD: No statsd sockets to listen to.");
188 + goto cleanup;
189 + }
190 +
191 + poll_events(&statsd_sockets
192 + , statsd_add_callback
193 + , statsd_del_callback
194 + , statsd_rcv_callback
195 + );
196 +
197 +cleanup:
198 + debug(D_WEB_CLIENT, "STATSD: exit!");
199 + listen_sockets_close(&statsd_sockets);
200 +
201 + pthread_exit(NULL);
202 + return NULL;
203 +}
src/statsd.h new
+9
@@ -0,0 +1,9 @@
1 +#ifndef NETDATA_STATSD_H
2 +#define NETDATA_STATSD_H
3 +
4 +#define STATSD_LISTEN_PORT 8125
5 +#define STATSD_LISTEN_BACKLOG 4096
6 +
7 +extern void *statsd_main(void *ptr);
8 +
9 +#endif //NETDATA_STATSD_H
src/web_client.c
+1 -34
@@ -57,13 +57,7 @@ struct web_client *web_client_create(int listener) {
57 w->mode = WEB_CLIENT_MODE_NORMAL;
58
59 {
60 - struct sockaddr *sadr;
61 - socklen_t addrlen;
62 -
63 - sadr = (struct sockaddr*) &w->clientaddr;
64 - addrlen = sizeof(w->clientaddr);
65 -
66 - w->ifd = accept4(listener, sadr, &addrlen, SOCK_NONBLOCK);
60 + w->ifd = accept_socket(listener, SOCK_NONBLOCK, w->client_ip, sizeof(w->client_ip), w->client_port, sizeof(w->client_port));
61 if (w->ifd == -1) {
62 error("%llu: Cannot accept new incoming connection.", w->id);
63 freez(w);
@@ -71,33 +65,6 @@ struct web_client *web_client_create(int listener) {
65 }
66 w->ofd = w->ifd;
67
74 - if(getnameinfo(sadr, addrlen, w->client_ip, NI_MAXHOST, w->client_port, NI_MAXSERV, NI_NUMERICHOST | NI_NUMERICSERV) != 0) {
75 - error("Cannot getnameinfo() on received client connection.");
76 - strncpyz(w->client_ip, "UNKNOWN", NI_MAXHOST);
77 - strncpyz(w->client_port, "UNKNOWN", NI_MAXSERV);
78 - }
79 - w->client_ip[NI_MAXHOST] = '\0';
80 - w->client_port[NI_MAXSERV] = '\0';
81 -
82 - switch(sadr->sa_family) {
83 - case AF_INET:
84 - debug(D_WEB_CLIENT_ACCESS, "%llu: New IPv4 web client from %s port %s on socket %d.", w->id, w->client_ip, w->client_port, w->ifd);
85 - break;
86 -
87 - case AF_INET6:
88 - if(strncmp(w->client_ip, "::ffff:", 7) == 0) {
89 - memmove(w->client_ip, &w->client_ip[7], strlen(&w->client_ip[7]) + 1);
90 - debug(D_WEB_CLIENT_ACCESS, "%llu: New IPv4 web client from %s port %s on socket %d.", w->id, w->client_ip, w->client_port, w->ifd);
91 - }
92 - else
93 - debug(D_WEB_CLIENT_ACCESS, "%llu: New IPv6 web client from %s port %s on socket %d.", w->id, w->client_ip, w->client_port, w->ifd);
94 - break;
95 -
96 - default:
97 - debug(D_WEB_CLIENT_ACCESS, "%llu: New UNKNOWN web client from %s port %s on socket %d.", w->id, w->client_ip, w->client_port, w->ifd);
98 - break;
99 - }
100 -
68 int flag = 1;
69 if(setsockopt(w->ofd, IPPROTO_TCP, TCP_NODELAY, (char *) &flag, sizeof(int)) != 0)
70 error("%llu: failed to enable TCP_NODELAY on socket.", w->id);
src/web_client.h
-1
@@ -83,7 +83,6 @@ struct web_client {
83 char cookie2[COOKIE_MAX+1];
84 char origin[ORIGIN_MAX+1];
85
86 - struct sockaddr_storage clientaddr;
86 struct response response;
87
88 size_t stats_received_bytes;
src/web_server.c
+4 -3
@@ -1,9 +1,10 @@
1 #include "common.h"
2
3 static LISTEN_SOCKETS api_sockets = {
4 - .config_section = CONFIG_SECTION_WEB,
5 - .default_port = API_LISTEN_PORT,
6 - .backlog = API_LISTEN_BACKLOG
4 + .config_section = CONFIG_SECTION_WEB,
5 + .default_bind_to = "*",
6 + .default_port = API_LISTEN_PORT,
7 + .backlog = API_LISTEN_BACKLOG
8 };
9
10 WEB_SERVER_MODE web_server_mode = WEB_SERVER_MODE_MULTI_THREADED;