@cryptotaxi247 / netdata-1 / commits / b33b7ee4f

static-threaded web server completed

Costa Tsaousis (ktsaou) committed Jan 9, 2018 at 00:00 UTC b33b7ee4fdbc3c46480f828a155c42db2285036d
9 files changed +305 -245
src/global_statistics.c
+13 -8
@@ -7,7 +7,8 @@ volatile struct global_statistics global_statistics = {
7 .bytes_received = 0,
8 .bytes_sent = 0,
9 .content_size = 0,
10 - .compressed_content_size = 0
10 + .compressed_content_size = 0,
11 + .web_client_count = 1
12 };
13
14 netdata_mutex_t global_statistics_mutex = NETDATA_MUTEX_INITIALIZER;
@@ -38,7 +39,7 @@ void finished_web_request_statistics(uint64_t dt,
39 __atomic_fetch_add(&global_statistics.compressed_content_size, compressed_content_size, __ATOMIC_SEQ_CST);
40 #else
41 #warning NOT using atomic operations - using locks for global statistics
41 - if (web_server_mode == WEB_SERVER_MODE_MULTI_THREADED)
42 + if (web_server_is_multithreaded)
43 global_statistics_lock();
44
45 if (dt > global_statistics.web_usec_max)
@@ -51,35 +52,39 @@ void finished_web_request_statistics(uint64_t dt,
52 global_statistics.content_size += content_size;
53 global_statistics.compressed_content_size += compressed_content_size;
54
54 - if (web_server_mode == WEB_SERVER_MODE_MULTI_THREADED)
55 + if (web_server_is_multithreaded)
56 global_statistics_unlock();
57 #endif
58 }
59
59 -void web_client_connected(void) {
60 +uint64_t web_client_connected(void) {
61 #if defined(HAVE_C___ATOMIC) && !defined(NETDATA_NO_ATOMIC_INSTRUCTIONS)
62 __atomic_fetch_add(&global_statistics.connected_clients, 1, __ATOMIC_SEQ_CST);
63 + uint64_t id = __atomic_fetch_add(&global_statistics.web_client_count, 1, __ATOMIC_SEQ_CST);
64 #else
63 - if (web_server_mode == WEB_SERVER_MODE_MULTI_THREADED)
65 + if (web_server_is_multithreaded)
66 global_statistics_lock();
67
68 global_statistics.connected_clients++;
69 + uint64_t id = global_statistics.web_client_count++;
70
68 - if (web_server_mode == WEB_SERVER_MODE_MULTI_THREADED)
71 + if (web_server_is_multithreaded)
72 global_statistics_unlock();
73 #endif
74 +
75 + return id;
76 }
77
78 void web_client_disconnected(void) {
79 #if defined(HAVE_C___ATOMIC) && !defined(NETDATA_NO_ATOMIC_INSTRUCTIONS)
80 __atomic_fetch_sub(&global_statistics.connected_clients, 1, __ATOMIC_SEQ_CST);
81 #else
77 - if (web_server_mode == WEB_SERVER_MODE_MULTI_THREADED)
82 + if (web_server_is_multithreaded)
83 global_statistics_lock();
84
85 global_statistics.connected_clients--;
86
82 - if (web_server_mode == WEB_SERVER_MODE_MULTI_THREADED)
87 + if (web_server_is_multithreaded)
88 global_statistics_unlock();
89 #endif
90 }
src/global_statistics.h
+3 -1
@@ -14,6 +14,8 @@ struct global_statistics {
14 volatile uint64_t bytes_sent;
15 volatile uint64_t content_size;
16 volatile uint64_t compressed_content_size;
17 +
18 + volatile uint64_t web_client_count;
19 };
20
21 extern volatile struct global_statistics global_statistics;
@@ -26,7 +28,7 @@ extern void finished_web_request_statistics(uint64_t dt,
28 uint64_t content_size,
29 uint64_t compressed_content_size);
30
29 -extern void web_client_connected(void);
31 +extern uint64_t web_client_connected(void);
32 extern void web_client_disconnected(void);
33
34 #define GLOBAL_STATS_RESET_WEB_USEC_MAX 0x01
src/socket.c
+1 -1
@@ -1051,7 +1051,7 @@ inline POLLINFO *poll_add_fd(POLLJOB *p
1051 return pi;
1052 }
1053
1054 -static inline void poll_close_fd(POLLINFO *pi) {
1054 +inline void poll_close_fd(POLLINFO *pi) {
1055 POLLJOB *p = pi->p;
1056
1057 struct pollfd *pf = &p->fds[pi->slot];
src/socket.h
+1
@@ -119,6 +119,7 @@ extern POLLINFO *poll_add_fd(POLLJOB *p
119 , int (*snd_callback)(POLLINFO *pi, short int *events)
120 , void *data
121 );
122 +extern void poll_close_fd(POLLINFO *pi);
123
124 extern void poll_events(LISTEN_SOCKETS *sockets
125 , void *(*add_callback)(POLLINFO *pi, short int *events, void *data)
src/web_buffer.h
+5 -2
@@ -47,8 +47,6 @@ typedef struct web_buffer {
47 #define buffer_strlen(wb) ((wb)->len)
48 extern const char *buffer_tostring(BUFFER *wb);
49
50 -#define buffer_need_bytes(buffer, needed_free_size) do { if(unlikely((buffer)->size - (buffer)->len < (size_t)(needed_free_size))) buffer_increase((buffer), (size_t)(needed_free_size)); } while(0)
51 -
50 #define buffer_flush(wb) wb->buffer[(wb)->len = 0] = '\0'
51 extern void buffer_reset(BUFFER *wb);
52
@@ -75,4 +73,9 @@ extern char *print_number_llu_r_smart(char *str, unsigned long long uvalue);
73
74 extern void buffer_print_llu(BUFFER *wb, unsigned long long uvalue);
75
76 +static inline void buffer_need_bytes(BUFFER *buffer, size_t needed_free_size) {
77 + if(unlikely(buffer->size - buffer->len < needed_free_size))
78 + buffer_increase(buffer, needed_free_size);
79 +}
80 +
81 #endif /* NETDATA_WEB_BUFFER_H */
src/web_client.c
+107 -120
@@ -21,9 +21,6 @@ SIMPLE_PATTERN *web_allow_badges_from = NULL;
21 int web_enable_gzip = 1, web_gzip_level = 3, web_gzip_strategy = Z_DEFAULT_STRATEGY;
22 #endif /* NETDATA_WITH_ZLIB */
23
24 -struct web_client *web_clients = NULL;
25 -unsigned long long web_clients_count = 0;
26 -
24 static inline int web_client_crock_socket(struct web_client *w) {
25 #ifdef TCP_CORK
26 if(likely(web_client_is_corkable(w) && !w->tcp_cork && w->ofd != -1)) {
@@ -84,33 +81,14 @@ static void web_client_update_acl_matches(struct web_client *w) {
81 w->acl |= WEB_CLIENT_ACL_BADGE;
82 }
83
87 -struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const char *client_port) {
88 - struct web_client *w;
89 -
90 - w = callocz(1, sizeof(struct web_client));
91 - w->id = ++web_clients_count;
92 - w->mode = WEB_CLIENT_MODE_NORMAL;
93 - w->ifd = fd;
84 +static void web_client_fix_socket(struct web_client *w) {
85 + int flag = 1;
86 + if(setsockopt(w->ifd, IPPROTO_TCP, TCP_NODELAY, (char *) &flag, sizeof(int)) != 0)
87 + error("%llu: failed to enable TCP_NODELAY on socket fd %d.", w->id, w->ofd);
88
95 - strncpyz(w->client_ip, client_ip, sizeof(w->client_ip));
96 - strncpyz(w->client_port, client_port, sizeof(w->client_port));
97 -
98 - if(unlikely(!*w->client_ip)) strcpy(w->client_ip, "-");
99 - if(unlikely(!*w->client_port)) strcpy(w->client_port, "-");
100 -
101 - log_connection(w, "CONNECTED");
102 -
103 - w->ofd = w->ifd;
104 -
105 - {
106 - int flag = 1;
107 - if(setsockopt(w->ofd, IPPROTO_TCP, TCP_NODELAY, (char *) &flag, sizeof(int)) != 0)
108 - error("%llu: failed to enable TCP_NODELAY on socket fd %d.", w->id, w->ofd);
109 -
110 - flag = 1;
111 - if(setsockopt(w->ifd, SOL_SOCKET, SO_KEEPALIVE, (char *) &flag, sizeof(int)) != 0)
112 - error("%llu: Cannot set SO_KEEPALIVE on socket fd %d.", w->id, w->ifd);
113 - }
89 + flag = 1;
90 + if(setsockopt(w->ifd, SOL_SOCKET, SO_KEEPALIVE, (char *) &flag, sizeof(int)) != 0)
91 + error("%llu: Cannot set SO_KEEPALIVE on socket fd %d.", w->id, w->ifd);
92
93 web_client_update_acl_matches(w);
94
@@ -121,61 +99,53 @@ struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const
99 web_client_enable_wait_receive(w);
100
101 web_client_connected();
124 - return(w);
102 + log_connection(w, "CONNECTED");
103 }
104
127 -struct web_client *web_client_create_on_listenfd(int listener) {
105 +struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const char *client_port) {
106 struct web_client *w;
107
108 w = callocz(1, sizeof(struct web_client));
131 - w->id = ++web_clients_count;
109 + w->id = web_client_connected();
110 w->mode = WEB_CLIENT_MODE_NORMAL;
111 + w->ifd = w->ofd = fd;
112
134 - {
135 - w->ifd = accept_socket(listener, SOCK_NONBLOCK, w->client_ip, sizeof(w->client_ip), w->client_port, sizeof(w->client_port), web_allow_connections_from);
136 -
137 - if(unlikely(!*w->client_ip)) strcpy(w->client_ip, "-");
138 - if(unlikely(!*w->client_port)) strcpy(w->client_port, "-");
113 + strncpyz(w->client_ip, client_ip, sizeof(w->client_ip) - 1);
114 + strncpyz(w->client_port, client_port, sizeof(w->client_port) - 1);
115
140 - if (w->ifd == -1) {
141 - if(errno == EPERM)
142 - log_connection(w, "ACCESS DENIED");
143 - else {
144 - log_connection(w, "CONNECTION FAILED");
145 - error("%llu: Failed to accept new incoming connection.", w->id);
146 - }
147 -
148 - freez(w);
149 - return NULL;
150 - }
151 - else
152 - log_connection(w, "CONNECTED");
116 + if(unlikely(!*w->client_ip)) strcpy(w->client_ip, "-");
117 + if(unlikely(!*w->client_port)) strcpy(w->client_port, "-");
118
154 - w->ofd = w->ifd;
119 + web_client_fix_socket(w);
120 + return(w);
121 +}
122
156 - int flag = 1;
157 - if(setsockopt(w->ofd, IPPROTO_TCP, TCP_NODELAY, (char *) &flag, sizeof(int)) != 0)
158 - error("%llu: failed to enable TCP_NODELAY on socket fd %d.", w->id, w->ofd);
123 +struct web_client *web_client_create_on_listenfd(int listener) {
124 + struct web_client *w;
125
160 - flag = 1;
161 - if(setsockopt(w->ifd, SOL_SOCKET, SO_KEEPALIVE, (char *) &flag, sizeof(int)) != 0)
162 - error("%llu: Cannot set SO_KEEPALIVE on socket fd %d.", w->id, w->ifd);
163 - }
126 + w = callocz(1, sizeof(struct web_client));
127 + w->id = web_client_connected();
128 + w->mode = WEB_CLIENT_MODE_NORMAL;
129
165 - web_client_update_acl_matches(w);
130 + w->ifd = w->ofd = accept_socket(listener, SOCK_NONBLOCK, w->client_ip, sizeof(w->client_ip), w->client_port, sizeof(w->client_port), web_allow_connections_from);
131
167 - w->response.data = buffer_create(INITIAL_WEB_DATA_LENGTH);
168 - w->response.header = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
169 - w->response.header_output = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
170 - w->origin[0] = '*';
171 - web_client_enable_wait_receive(w);
132 + if(unlikely(!*w->client_ip)) strcpy(w->client_ip, "-");
133 + if(unlikely(!*w->client_port)) strcpy(w->client_port, "-");
134
173 - if(web_clients) web_clients->prev = w;
174 - w->next = web_clients;
175 - web_clients = w;
135 + if (w->ifd == -1) {
136 + if(errno == EPERM)
137 + log_connection(w, "ACCESS DENIED");
138 + else {
139 + log_connection(w, "CONNECTION FAILED");
140 + error("%llu: Failed to accept new incoming connection.", w->id);
141 + }
142
177 - web_client_connected();
143 + freez(w);
144 + web_client_disconnected();
145 + return NULL;
146 + }
147
148 + web_client_fix_socket(w);
149 return(w);
150 }
151
@@ -253,7 +223,11 @@ void web_client_reset(struct web_client *w) {
223 if(unlikely(w->mode == WEB_CLIENT_MODE_FILECOPY)) {
224 if(w->ifd != w->ofd) {
225 debug(D_WEB_CLIENT, "%llu: Closing filecopy input file descriptor %d.", w->id, w->ifd);
256 - if(w->ifd != -1) close(w->ifd);
226 +
227 + if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
228 + if (w->ifd != -1) close(w->ifd);
229 + }
230 +
231 w->ifd = w->ofd;
232 }
233 }
@@ -300,30 +274,27 @@ void web_client_reset(struct web_client *w) {
274 #endif // NETDATA_WITH_ZLIB
275 }
276
303 -struct web_client *web_client_free(struct web_client *w) {
277 +void web_client_free(struct web_client *w) {
278 debug(D_WEB_CLIENT_ACCESS, "%llu: Closing web client from %s port %s.", w->id, w->client_ip, w->client_port);
279
280 web_client_reset(w);
281
308 - struct web_client *n = NULL;
309 - if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
310 - struct web_client *n = w->next;
311 - if (w == web_clients) web_clients = n;
312 -
313 - if(w->prev) w->prev->next = w->next;
314 - if(w->next) w->next->prev = w->prev;
315 - }
316 -
282 buffer_free(w->response.header_output);
283 + w->response.header_output = NULL;
284 +
285 buffer_free(w->response.header);
286 + w->response.header = NULL;
287 +
288 buffer_free(w->response.data);
320 - if(w->ifd != -1) close(w->ifd);
321 - if(w->ofd != -1 && w->ofd != w->ifd) close(w->ofd);
322 - freez(w);
289 + w->response.data = NULL;
290
324 - web_client_disconnected();
291 + if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
292 + if (w->ifd != -1) close(w->ifd);
293 + if (w->ofd != -1 && w->ofd != w->ifd) close(w->ofd);
294 + }
295
326 - return(n);
296 + freez(w);
297 + web_client_disconnected();
298 }
299
300 uid_t web_files_uid(void) {
@@ -569,6 +540,7 @@ int mysendfile(struct web_client *w, char *filename) {
540 web_client_enable_wait_receive(w);
541 web_client_disable_wait_send(w);
542 buffer_flush(w->response.data);
543 + buffer_need_bytes(w->response.data, (size_t)statbuf.st_size);
544 w->response.rlen = (size_t)statbuf.st_size;
545 #ifdef __APPLE__
546 w->response.data->date = statbuf.st_mtimespec.tv_sec;
@@ -1734,72 +1706,87 @@ ssize_t web_client_send(struct web_client *w) {
1706 return(bytes);
1707 }
1708
1737 -ssize_t web_client_receive(struct web_client *w)
1709 +ssize_t web_client_read_file(struct web_client *w)
1710 {
1739 - // do we have any space for more data?
1740 - buffer_need_bytes(w->response.data, WEB_REQUEST_LENGTH);
1711 + if(unlikely(w->response.rlen > w->response.data->size))
1712 + buffer_need_bytes(w->response.data, w->response.rlen - w->response.data->size);
1713
1742 - ssize_t left = w->response.data->size - w->response.data->len;
1743 - ssize_t bytes;
1744 -
1745 - if(unlikely(w->mode == WEB_CLIENT_MODE_FILECOPY))
1746 - bytes = read(w->ifd, &w->response.data->buffer[w->response.data->len], (size_t) (left - 1));
1747 - else
1748 - bytes = recv(w->ifd, &w->response.data->buffer[w->response.data->len], (size_t) (left - 1), MSG_DONTWAIT);
1714 + if(unlikely(w->response.rlen <= w->response.data->len))
1715 + return 0;
1716
1717 + ssize_t left = w->response.rlen - w->response.data->len;
1718 + ssize_t bytes = read(w->ifd, &w->response.data->buffer[w->response.data->len], (size_t) (left - 1));
1719 if(likely(bytes > 0)) {
1751 - if(w->mode != WEB_CLIENT_MODE_FILECOPY)
1752 - w->stats_received_bytes += bytes;
1753 -
1720 size_t old = w->response.data->len;
1721 w->response.data->len += bytes;
1722 w->response.data->buffer[w->response.data->len] = '\0';
1723
1758 - debug(D_WEB_CLIENT, "%llu: Received %zd bytes.", w->id, bytes);
1759 - debug(D_WEB_DATA, "%llu: Received data: '%s'.", w->id, &w->response.data->buffer[old]);
1724 + debug(D_WEB_CLIENT, "%llu: Read %zd bytes.", w->id, bytes);
1725 + debug(D_WEB_DATA, "%llu: Read data: '%s'.", w->id, &w->response.data->buffer[old]);
1726
1761 - if(w->mode == WEB_CLIENT_MODE_FILECOPY) {
1762 - web_client_enable_wait_send(w);
1727 + web_client_enable_wait_send(w);
1728
1764 - if(w->response.rlen && w->response.data->len >= w->response.rlen)
1765 - web_client_disable_wait_receive(w);
1766 - }
1729 + if(w->response.rlen && w->response.data->len >= w->response.rlen)
1730 + web_client_disable_wait_receive(w);
1731 }
1732 else if(likely(bytes == 0)) {
1769 - debug(D_WEB_CLIENT, "%llu: Out of input data.", w->id);
1733 + debug(D_WEB_CLIENT, "%llu: Out of input file data.", w->id);
1734
1735 // if we cannot read, it means we have an error on input.
1736 // if however, we are copying a file from ifd to ofd, we should not return an error.
1737 // in this case, the error should be generated when the file has been sent to the client.
1738
1775 - if(w->mode == WEB_CLIENT_MODE_FILECOPY) {
1776 - // we are copying data from ifd to ofd
1777 - // let it finish copying...
1778 - web_client_disable_wait_receive(w);
1779 -
1780 - debug(D_WEB_CLIENT, "%llu: Read the whole file.", w->id);
1739 + // we are copying data from ifd to ofd
1740 + // let it finish copying...
1741 + web_client_disable_wait_receive(w);
1742
1782 - if(w->ifd != w->ofd && web_server_mode != WEB_SERVER_MODE_STATIC_THREADED)
1783 - close(w->ifd);
1743 + debug(D_WEB_CLIENT, "%llu: Read the whole file.", w->id);
1744
1785 - w->ifd = w->ofd;
1786 - }
1787 - else {
1788 - debug(D_WEB_CLIENT, "%llu: failed to receive data.", w->id);
1789 - WEB_CLIENT_IS_DEAD(w);
1745 + if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
1746 + if (w->ifd != w->ofd) close(w->ifd);
1747 }
1748 +
1749 + w->ifd = w->ofd;
1750 }
1751 else {
1793 - debug(D_WEB_CLIENT, "%llu: receive data failed.", w->id);
1752 + debug(D_WEB_CLIENT, "%llu: read data failed.", w->id);
1753 WEB_CLIENT_IS_DEAD(w);
1754 }
1755
1756 return(bytes);
1757 }
1758
1759 +ssize_t web_client_receive(struct web_client *w)
1760 +{
1761 + if(unlikely(w->mode == WEB_CLIENT_MODE_FILECOPY))
1762 + return web_client_read_file(w);
1763 +
1764 + // do we have any space for more data?
1765 + buffer_need_bytes(w->response.data, WEB_REQUEST_LENGTH);
1766 +
1767 + ssize_t left = w->response.data->size - w->response.data->len;
1768 + ssize_t bytes = recv(w->ifd, &w->response.data->buffer[w->response.data->len], (size_t) (left - 1), MSG_DONTWAIT);
1769 +
1770 + if(likely(bytes > 0)) {
1771 + w->stats_received_bytes += bytes;
1772 +
1773 + size_t old = w->response.data->len;
1774 + w->response.data->len += bytes;
1775 + w->response.data->buffer[w->response.data->len] = '\0';
1776 +
1777 + debug(D_WEB_CLIENT, "%llu: Received %zd bytes.", w->id, bytes);
1778 + debug(D_WEB_DATA, "%llu: Received data: '%s'.", w->id, &w->response.data->buffer[old]);
1779 + }
1780 + else {
1781 + debug(D_WEB_CLIENT, "%llu: receive data failed.", w->id);
1782 + WEB_CLIENT_IS_DEAD(w);
1783 + }
1784 +
1785 + return(bytes);
1786 +}
1787
1788 // --------------------------------------------------------------------------------------
1802 -// the thread of a single client
1789 +// the thread of a single client - for the MULTI-THREADED web server
1790
1791 // 1. waits for input and output, using async I/O
1792 // 2. it processes HTTP requests
src/web_client.h
+15 -9
@@ -129,11 +129,11 @@ typedef enum web_client_acl {
129 struct web_client {
130 unsigned long long id;
131
132 - WEB_CLIENT_FLAGS flags; // status flags for the client
133 - WEB_CLIENT_MODE mode; // the operational mode of the client
134 - WEB_CLIENT_ACL acl; // the access list of the client
132 + WEB_CLIENT_FLAGS flags; // status flags for the client
133 + WEB_CLIENT_MODE mode; // the operational mode of the client
134 + WEB_CLIENT_ACL acl; // the access list of the client
135
136 - int tcp_cork; // 1 = we have a cork on the socket
136 + int tcp_cork; // 1 = we have a cork on the socket
137
138 int ifd;
139 int ofd;
@@ -155,13 +155,16 @@ struct web_client {
155 size_t stats_received_bytes;
156 size_t stats_sent_bytes;
157
158 - netdata_thread_t thread; // the thread servicing this client
158 + // STATIC-THREADED WEB SERVER MEMBERS
159 + size_t pollinfo_slot; // POLLINFO slot of the web client
160 + size_t pollinfo_filecopy_slot; // POLLINFO slot of the file read
161
160 - struct web_client *prev;
161 - struct web_client *next;
162 + // MULTI-THREADED WEB SERVER MEMBERS
163 + netdata_thread_t thread; // the thread servicing this client
164 + struct web_client *prev; // maintain a linked list of web clients
165 + struct web_client *next; // for the web servers that need it
166 };
167
164 -extern struct web_client *web_clients;
168 extern SIMPLE_PATTERN *web_allow_connections_from;
169 extern SIMPLE_PATTERN *web_allow_dashboard_from;
170 extern SIMPLE_PATTERN *web_allow_registry_from;
@@ -176,9 +179,12 @@ extern int web_client_permission_denied(struct web_client *w);
179
180 extern struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const char *client_port);
181 extern struct web_client *web_client_create_on_listenfd(int listener);
179 -extern struct web_client *web_client_free(struct web_client *w);
182 +extern void web_client_free(struct web_client *w);
183 +
184 extern ssize_t web_client_send(struct web_client *w);
185 extern ssize_t web_client_receive(struct web_client *w);
186 +extern ssize_t web_client_read_file(struct web_client *w);
187 +
188 extern void web_client_process_request(struct web_client *w);
189 extern void web_client_reset(struct web_client *w);
190
src/web_server.c
+159 -104
@@ -8,6 +8,9 @@ static LISTEN_SOCKETS api_sockets = {
8 };
9
10 WEB_SERVER_MODE web_server_mode = WEB_SERVER_MODE_MULTI_THREADED;
11 +int web_server_is_multithreaded = 1;
12 +
13 +// --------------------------------------------------------------------------------------
14
15 #ifdef NETDATA_INTERNAL_CHECKS
16 static void log_allocations(void)
@@ -19,21 +22,13 @@ static void log_allocations(void)
22
23 mi = mallinfo();
24 if(mi.uordblks > used) {
22 - int clients = 0;
23 -
24 - if(web_server_mode != WEB_SERVER_MODE_STATIC_THREADED) {
25 - struct web_client *w;
26 - for (w = web_clients; w; w = w->next) clients++;
27 - }
28 -
29 - info("Allocated memory: used %d KB (+%d B), mmap %d KB (+%d B), heap %d KB (+%d B). %d web clients connected.",
25 + info("Allocated memory: used %d KB (+%d B), mmap %d KB (+%d B), heap %d KB (+%d B).",
26 mi.uordblks / 1024,
27 mi.uordblks - used,
28 mi.hblkhd / 1024,
29 mi.hblkhd - mmap,
30 mi.arena / 1024,
35 - mi.arena - heap,
36 - clients);
31 + mi.arena - heap);
32
33 used = mi.uordblks;
34 heap = mi.arena;
@@ -93,13 +88,35 @@ int api_listen_sockets_setup(void) {
88 // --------------------------------------------------------------------------------------
89 // the main socket listener - MULTI-THREADED
90
91 +// maintain a linked list of web clients - for the web servers that need it
92 +static struct web_client *web_clients = NULL;
93 +
94 +static void web_client_link(struct web_client *w) {
95 + if (web_clients) web_clients->prev = w;
96 + w->next = web_clients;
97 + web_clients = w;
98 +}
99 +
100 +static struct web_client *web_client_unlink(struct web_client *w) {
101 + struct web_client *n = w->next;
102 + if (w == web_clients) web_clients = n;
103 +
104 + if(w->prev) w->prev->next = w->next;
105 + if(w->next) w->next->prev = w->prev;
106 +
107 + return n;
108 +}
109 +
110 static inline void multi_threaded_cleanup_web_clients(void) {
111 struct web_client *w;
112
113 for (w = web_clients; w;) {
114 if (web_client_check_obsolete(w)) {
115 debug(D_WEB_CLIENT, "%llu: Removing client.", w->id);
102 - w = web_client_free(w);
116 + struct web_client *t = web_client_unlink(w);
117 + web_client_free(w);
118 + w = t;
119 +
120 #ifdef NETDATA_INTERNAL_CHECKS
121 log_allocations();
122 #endif
@@ -147,6 +164,7 @@ void *socket_listen_main_multi_threaded(void *ptr) {
164 netdata_thread_cleanup_push(socket_listen_main_multi_threaded_cleanup, ptr);
165
166 web_server_mode = WEB_SERVER_MODE_MULTI_THREADED;
167 + web_server_is_multithreaded = 1;
168
169 struct web_client *w;
170 int retval, counter = 0;
@@ -195,6 +213,7 @@ void *socket_listen_main_multi_threaded(void *ptr) {
213 // no need for error log - web_client_create_on_listenfd already logged the error
214 continue;
215 }
216 + web_client_link(w);
217
218 if(api_sockets.fds_families[i] == AF_UNIX)
219 web_client_set_unix(w);
@@ -286,9 +305,9 @@ static void socket_listen_main_single_threaded_cleanup(void *data) {
305 void *socket_listen_main_single_threaded(void *ptr) {
306 netdata_thread_cleanup_push(socket_listen_main_single_threaded_cleanup, ptr);
307 web_server_mode = WEB_SERVER_MODE_SINGLE_THREADED;
308 + web_server_is_multithreaded = 0;
309
310 struct web_client *w;
291 - int retval;
311
312 if(!api_sockets.opened)
313 fatal("LISTENER: no listen sockets available.");
@@ -322,7 +341,7 @@ void *socket_listen_main_single_threaded(void *ptr) {
341 rifds = ifds;
342 rofds = ofds;
343 refds = efds;
325 - retval = select(fdmax+1, &rifds, &rofds, &refds, &tv);
344 + int retval = select(fdmax+1, &rifds, &rofds, &refds, &tv);
345
346 if(unlikely(retval == -1)) {
347 error("LISTENER: select() failed.");
@@ -335,6 +354,8 @@ void *socket_listen_main_single_threaded(void *ptr) {
354 if (FD_ISSET(api_sockets.fds[i], &rifds)) {
355 debug(D_WEB_CLIENT_ACCESS, "LISTENER: new connection.");
356 w = web_client_create_on_listenfd(api_sockets.fds[i]);
357 + if(unlikely(!w))
358 + continue;
359
360 if(api_sockets.fds_families[i] == AF_UNIX)
361 web_client_set_unix(w);
@@ -392,9 +413,6 @@ void *socket_listen_main_single_threaded(void *ptr) {
413 }
414 else {
415 debug(D_WEB_CLIENT_ACCESS, "LISTENER: single threaded web server timeout.");
395 -#ifdef NETDATA_INTERNAL_CHECKS
396 - log_allocations();
397 -#endif
416 }
417 }
418
@@ -415,133 +433,163 @@ struct web_server_static_threaded_worker {
433 volatile size_t receptions;
434 volatile size_t sends;
435 volatile size_t max_concurrent;
436 +
437 + volatile size_t files_read;
438 + volatile size_t file_reads;
439 };
440
441 static long long static_threaded_workers_count = 1;
442 static struct web_server_static_threaded_worker *static_workers_private_data = NULL;
443 static __thread struct web_server_static_threaded_worker *worker_private = NULL;
444
424 -// new TCP client connected
425 -static void *web_server_add_callback(POLLINFO *pi, short int *events, void *data) {
426 - (void)data;
445 +// ----------------------------------------------------------------------------
446
428 - worker_private->connected++;
447 +static inline int web_server_check_client_status(struct web_client *w) {
448 + if(unlikely(web_client_check_obsolete(w) || web_client_check_dead(w) || (!web_client_has_wait_receive(w) && !web_client_has_wait_send(w))))
449 + return -1;
450
430 - size_t concurrent = worker_private->connected - worker_private->disconnected;
431 - if(unlikely(concurrent > worker_private->max_concurrent))
432 - worker_private->max_concurrent = concurrent;
451 + return 0;
452 +}
453
434 - *events = POLLIN;
454 +// ----------------------------------------------------------------------------
455 +// web server files
456
436 - debug(D_WEB_CLIENT_ACCESS, "LISTENER on %d: new connection.", pi->fd);
437 - struct web_client *w = web_client_create_on_fd(pi->fd, pi->client_ip, pi->client_port);
457 +static void *web_server_file_add_callback(POLLINFO *pi, short int *events, void *data) {
458 + struct web_client *w = (struct web_client *)data;
459
439 - if(unlikely(pi->socktype == AF_UNIX))
440 - web_client_set_unix(w);
441 - else
442 - web_client_set_tcp(w);
460 + worker_private->files_read++;
461
444 - return (void *)w;
462 + debug(D_WEB_CLIENT, "%llu: ADDED FILE READ ON FD %d", w->id, pi->fd);
463 + *events = POLLIN;
464 + pi->data = w;
465 + return w;
466 }
467
447 -// TCP client disconnected
448 -static void web_server_del_callback(POLLINFO *pi) {
449 - worker_private->disconnected++;
450 -
468 +static void web_werver_file_del_callback(POLLINFO *pi) {
469 struct web_client *w = (struct web_client *)pi->data;
452 - int fd = pi->fd;
453 -
454 - if(likely(w)) {
455 - if(w->ofd == -1 || fd == w->ofd) {
456 - // we free the client, only if the closing fd
457 - // is the client socket
470 + debug(D_WEB_CLIENT, "%llu: RELEASE FILE READ ON FD %d", w->id, pi->fd);
471
459 - // prevent it from closing the file descriptors
460 - w->ifd = w->ofd = -1;
472 + w->pollinfo_filecopy_slot = 0;
473
462 - web_client_free(w);
463 - }
474 + if(unlikely(!w->pollinfo_slot)) {
475 + debug(D_WEB_CLIENT, "%llu: CROSS WEB CLIENT CLEANUP (iFD %d, oFD %d)", w->id, pi->fd, w->ofd);
476 + web_client_free(w);
477 }
478 }
479
467 -// ----------------------------------------------------------------------------
468 -// web server files
480 +static int web_server_file_read_callback(POLLINFO *pi, short int *events) {
481 + struct web_client *w = (struct web_client *)pi->data;
482
470 -struct web_file_pollinfo {
471 - struct web_client *w;
472 - POLLINFO *pi;
473 -};
483 + // if there is no POLLINFO linked to this, it means the client disconnected
484 + // stop the file reading too
485 + if(unlikely(!w->pollinfo_slot)) {
486 + debug(D_WEB_CLIENT, "%llu: PREVENTED ATTEMPT TO READ FILE ON FD %d, ON CLOSED WEB CLIENT", w->id, pi->fd);
487 + return -1;
488 + }
489
475 -static void *web_server_file_add_callback(POLLINFO *pi, short int *events, void *data) {
476 - (void)pi;
490 + if(unlikely(w->mode != WEB_CLIENT_MODE_FILECOPY || w->ifd == w->ofd)) {
491 + debug(D_WEB_CLIENT, "%llu: PREVENTED ATTEMPT TO READ FILE ON FD %d, ON NON-FILECOPY WEB CLIENT", w->id, pi->fd);
492 + return -1;
493 + }
494 +
495 + debug(D_WEB_CLIENT, "%llu: READING FILE ON FD %d", w->id, pi->fd);
496 +
497 + worker_private->file_reads++;
498 + ssize_t ret = unlikely(web_client_read_file(w));
499 +
500 + if(likely(web_client_has_wait_send(w))) {
501 + POLLJOB *p = pi->p; // our POLLJOB
502 + POLLINFO *wpi = &p->inf[w->pollinfo_slot]; // POLLINFO of the client socket
503 +
504 + debug(D_WEB_CLIENT, "%llu: SIGNALING W TO SEND (iFD %d, oFD %d)", w->id, pi->fd, wpi->fd);
505 + p->fds[wpi->slot].events |= POLLOUT;
506 + }
507 +
508 + if(unlikely(ret <= 0 || w->ifd == w->ofd)) {
509 + debug(D_WEB_CLIENT, "%llu: DONE READING FILE ON FD %d", w->id, pi->fd);
510 + return -1;
511 + }
512
478 - info("ADDING FILE ON FD %d", pi->fd);
513 *events = POLLIN;
480 - return data;
514 + return 0;
515 }
516
483 -static void web_werver_file_del_callback(POLLINFO *pi) {
517 +static int web_server_file_write_callback(POLLINFO *pi, short int *events) {
518 (void)pi;
485 - info("DELETE FILE ON FD %d", pi->fd);
486 - freez(pi->data);
519 + (void)events;
520 +
521 + error("Writing to web files is not supported!");
522 +
523 + return -1;
524 }
525
489 -static int web_server_file_rcv_callback(POLLINFO *pi, short int *events) {
490 - *events = POLLIN;
526 +// ----------------------------------------------------------------------------
527 +// web server clients
528
492 - info("READING FILE ON FD %d", pi->fd);
529 +static void *web_server_add_callback(POLLINFO *pi, short int *events, void *data) {
530 + (void)data;
531
494 - struct web_file_pollinfo *wfpi = (struct web_file_pollinfo *)pi->data;
495 - if(unlikely(web_client_receive(wfpi->w) < 0))
496 - return -1;
532 + worker_private->connected++;
533
498 - wfpi->pi->p->fds[wfpi->pi->slot].events |= POLLOUT;
534 + size_t concurrent = worker_private->connected - worker_private->disconnected;
535 + if(unlikely(concurrent > worker_private->max_concurrent))
536 + worker_private->max_concurrent = concurrent;
537
500 - if(unlikely(wfpi->w->ifd == wfpi->w->ofd))
501 - return -1;
538 + *events = POLLIN;
539
503 - return 0;
540 + debug(D_WEB_CLIENT_ACCESS, "LISTENER on %d: new connection.", pi->fd);
541 + struct web_client *w = web_client_create_on_fd(pi->fd, pi->client_ip, pi->client_port);
542 + w->pollinfo_slot = pi->slot;
543 +
544 + if(unlikely(pi->socktype == AF_UNIX))
545 + web_client_set_unix(w);
546 + else
547 + web_client_set_tcp(w);
548 +
549 + debug(D_WEB_CLIENT, "%llu: ADDED CLIENT FD %d", w->id, pi->fd);
550 + return w;
551 }
552
506 -// ----------------------------------------------------------------------------
507 -//
553 +// TCP client disconnected
554 +static void web_server_del_callback(POLLINFO *pi) {
555 + worker_private->disconnected++;
556
509 -static inline int web_server_check_client_status(struct web_client *w) {
510 - if(unlikely(web_client_check_obsolete(w) || web_client_check_dead(w) || (!web_client_has_wait_receive(w) && !web_client_has_wait_send(w))))
511 - return -1;
557 + struct web_client *w = (struct web_client *)pi->data;
558
513 - return 0;
559 + w->pollinfo_slot = 0;
560 + if(unlikely(w->pollinfo_filecopy_slot)) {
561 + POLLJOB *p = pi->p; // our POLLJOB
562 + POLLINFO *fpi = &p->inf[w->pollinfo_filecopy_slot]; // POLLINFO of the client socket
563 + debug(D_WEB_CLIENT, "%llu: THE CLIENT WILL BE FRED BY READING FILE JOB ON FD %d", w->id, fpi->fd);
564 + }
565 + else {
566 + debug(D_WEB_CLIENT, "%llu: CLOSING CLIENT FD %d", w->id, pi->fd);
567 + web_client_free(w);
568 + }
569 }
570
516 -// Receive data
571 static int web_server_rcv_callback(POLLINFO *pi, short int *events) {
572 worker_private->receptions++;
573
574 struct web_client *w = (struct web_client *)pi->data;
575 int fd = pi->fd;
576
523 - if(unlikely(!web_client_has_wait_receive(w)))
524 - return -1;
525 -
577 if(unlikely(web_client_receive(w) < 0))
578 return -1;
579
529 - debug(D_WEB_CLIENT, "%llu: Processing received data.", w->id);
580 + debug(D_WEB_CLIENT, "%llu: processing received data on fd %d.", w->id, fd);
581 web_client_process_request(w);
582
583 if(unlikely(w->mode == WEB_CLIENT_MODE_FILECOPY)) {
533 - info("FILECOPY %d", pi->fd);
584 + if(w->pollinfo_filecopy_slot == 0) {
585 + debug(D_WEB_CLIENT, "%llu: FILECOPY DETECTED ON FD %d", w->id, pi->fd);
586
535 - if(unlikely(w->ifd != -1 && w->ifd != w->ofd && w->ifd != fd)) {
536 - // FIXME: we switched input fd
537 - // add a new socket to poll_events, with the same
538 - info("DETECTED FILECOPY ON FD %d", pi->fd);
587 + if (unlikely(w->ifd != -1 && w->ifd != w->ofd && w->ifd != fd)) {
588 + // add a new socket to poll_events, with the same
589 + debug(D_WEB_CLIENT, "%llu: CREATING FILECOPY SLOT ON FD %d", w->id, pi->fd);
590
540 - struct web_file_pollinfo *wfpi = callocz(1, sizeof(struct web_file_pollinfo));
541 - wfpi->w = w;
542 - wfpi->pi = pi;
543 -
544 - poll_add_fd(pi->p
591 + POLLINFO *fpi = poll_add_fd(
592 + pi->p
593 , w->ifd
594 , 0
595 , POLLINFO_FLAG_CLIENT_SOCKET
@@ -549,21 +597,19 @@ static int web_server_rcv_callback(POLLINFO *pi, short int *events) {
597 , ""
598 , web_server_file_add_callback
599 , web_werver_file_del_callback
552 - , web_server_file_rcv_callback
553 - , poll_default_snd_callback
554 - , (void *)wfpi
555 - );
556 - }
557 - else if(unlikely(w->ifd == -1)) {
558 - // FIXME: we closed input fd
559 - // instruct poll_events() to close fd
560 - info("INPUT CLOSED ON FD %d", pi->fd);
561 - return -1;
600 + , web_server_file_read_callback
601 + , web_server_file_write_callback
602 + , (void *) w
603 + );
604 +
605 + w->pollinfo_filecopy_slot = fpi->slot;
606 + }
607 }
608 }
564 -
565 - if(unlikely(w->ifd == fd && web_client_has_wait_receive(w)))
566 - *events |= POLLIN;
609 + else {
610 + if(unlikely(w->ifd == fd && web_client_has_wait_receive(w)))
611 + *events |= POLLIN;
612 + }
613
614 if(unlikely(w->ofd == fd && web_client_has_wait_send(w)))
615 *events |= POLLOUT;
@@ -577,8 +623,7 @@ static int web_server_snd_callback(POLLINFO *pi, short int *events) {
623 struct web_client *w = (struct web_client *)pi->data;
624 int fd = pi->fd;
625
580 - if(unlikely(!web_client_has_wait_send(w)))
581 - return -1;
626 + debug(D_WEB_CLIENT, "%llu: sending data on fd %d.", w->id, fd);
627
628 if(unlikely(web_client_send(w) < 0))
629 return -1;
@@ -592,6 +637,10 @@ static int web_server_snd_callback(POLLINFO *pi, short int *events) {
637 return web_server_check_client_status(w);
638 }
639
640 +
641 +// ----------------------------------------------------------------------------
642 +// web server worker thread
643 +
644 static void socket_listen_main_static_threaded_worker_cleanup(void *ptr) {
645 worker_private = (struct web_server_static_threaded_worker *)ptr;
646 worker_private->running = 0;
@@ -624,6 +673,10 @@ void *socket_listen_main_static_threaded_worker(void *ptr) {
673 return NULL;
674 }
675
676 +
677 +// ----------------------------------------------------------------------------
678 +// web server main thread - also becomes a worker
679 +
680 static void socket_listen_main_static_threaded_cleanup(void *ptr) {
681 struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
682 if(static_thread->enabled) {
@@ -661,6 +714,8 @@ void *socket_listen_main_static_threaded(void *ptr) {
714
715 static_workers_private_data = callocz((size_t)static_threaded_workers_count, sizeof(struct web_server_static_threaded_worker));
716
717 + web_server_is_multithreaded = (static_threaded_workers_count > 1);
718 +
719 int i;
720 for(i = 1; i < static_threaded_workers_count; i++) {
721 char tag[50 + 1];
src/web_server.h
+1
@@ -22,6 +22,7 @@ typedef enum web_server_mode {
22 } WEB_SERVER_MODE;
23
24 extern WEB_SERVER_MODE web_server_mode;
25 +extern int web_server_is_multithreaded;
26
27 extern WEB_SERVER_MODE web_server_mode_id(const char *mode);
28 extern const char *web_server_mode_name(WEB_SERVER_MODE id);