@cryptotaxi247 / netdata-1 / commits / 6f07d0dd5

static-threaded web server

Costa Tsaousis (ktsaou) committed Jan 7, 2018 at 23:07 UTC 6f07d0dd5c3fc0707d2c42577153d0724ba01a46
11 files changed +199 -87
src/log.c
+2 -2
@@ -338,8 +338,8 @@ void error_int( const char *prefix, const char *file, const char *function, cons
338 log_lock();
339
340 va_start( args, fmt );
341 - if(debug_flags) fprintf(stderr, "%s: %s %-5.5s : %s: (%04lu@%-10.10s:%-15.15s): ", date, program_name, prefix, netdata_thread_tag(), line, file, function);
342 - else fprintf(stderr, "%s: %s %-5.5s : %s: ", date, program_name, prefix, netdata_thread_tag());
341 + if(debug_flags) fprintf(stderr, "%s: %s %-5.5s : %s : (%04lu@%-10.10s:%-15.15s): ", date, program_name, prefix, netdata_thread_tag(), line, file, function);
342 + else fprintf(stderr, "%s: %s %-5.5s : %s : ", date, program_name, prefix, netdata_thread_tag());
343 vfprintf( stderr, fmt, args );
344 va_end( args );
345
src/main.c
+5
@@ -76,6 +76,7 @@ struct netdata_static_thread static_threads[] = {
76 {"PLUGINSD", NULL, NULL, 1, NULL, NULL, pluginsd_main},
77 {"WEB_SERVER[multi]", NULL, NULL, 1, NULL, NULL, socket_listen_main_multi_threaded},
78 {"WEB_SERVER[single]", NULL, NULL, 0, NULL, NULL, socket_listen_main_single_threaded},
79 + {"WEB_SERVER[static]", NULL, NULL, 0, NULL, NULL, socket_listen_main_static_threaded},
80 {"STREAM", NULL, NULL, 0, NULL, NULL, rrdpush_sender_thread},
81 {"STATSD", NULL, NULL, 1, NULL, NULL, statsd_main},
82
@@ -87,6 +88,7 @@ void web_server_threading_selection(void) {
88
89 int multi_threaded = (web_server_mode == WEB_SERVER_MODE_MULTI_THREADED);
90 int single_threaded = (web_server_mode == WEB_SERVER_MODE_SINGLE_THREADED);
91 + int static_threaded = (web_server_mode == WEB_SERVER_MODE_STATIC_THREADED);
92
93 int i;
94 for (i = 0; static_threads[i].name; i++) {
@@ -95,6 +97,9 @@ void web_server_threading_selection(void) {
97
98 if (static_threads[i].start_routine == socket_listen_main_single_threaded)
99 static_threads[i].enabled = single_threaded;
100 +
101 + if (static_threads[i].start_routine == socket_listen_main_static_threaded)
102 + static_threads[i].enabled = static_threaded;
103 }
104 }
105
src/procfile.c
+9 -3
@@ -289,7 +289,7 @@ procfile *procfile_readall(procfile *ff) {
289 debug(D_PROCFILE, "Reading file '%s', from position %zd with length %zd", procfile_filename(ff), s, (ssize_t)(ff->size - s));
290 r = read(ff->fd, &ff->data[s], ff->size - s);
291 if(unlikely(r == -1)) {
292 - if(unlikely(!(ff->flags & PROCFILE_FLAG_NO_ERROR_ON_FILE_IO))) error(PF_PREFIX ": Cannot read from file '%s'", procfile_filename(ff));
292 + if(unlikely(!(ff->flags & PROCFILE_FLAG_NO_ERROR_ON_FILE_IO))) error(PF_PREFIX ": Cannot read from file '%s' on fd %d", procfile_filename(ff), ff->fd);
293 procfile_close(ff);
294 return NULL;
295 }
@@ -408,6 +408,8 @@ procfile *procfile_open(const char *filename, const char *separators, uint32_t f
408 return NULL;
409 }
410
411 + // info("PROCFILE: opened '%s' on fd %d", filename, fd);
412 +
413 size_t size = (unlikely(procfile_adaptive_initial_allocation)) ? procfile_max_allocation : PROCFILE_INCREMENT_BUFFER;
414 procfile *ff = mallocz(sizeof(procfile) + size);
415
@@ -431,7 +433,10 @@ procfile *procfile_open(const char *filename, const char *separators, uint32_t f
433 procfile *procfile_reopen(procfile *ff, const char *filename, const char *separators, uint32_t flags) {
434 if(unlikely(!ff)) return procfile_open(filename, separators, flags);
435
434 - if(likely(ff->fd != -1)) close(ff->fd);
436 + if(likely(ff->fd != -1)) {
437 + // info("PROCFILE: closing fd %d", ff->fd);
438 + close(ff->fd);
439 + }
440
441 ff->fd = open(filename, O_RDONLY, 0666);
442 if(unlikely(ff->fd == -1)) {
@@ -439,9 +444,10 @@ procfile *procfile_reopen(procfile *ff, const char *filename, const char *separa
444 return NULL;
445 }
446
447 + // info("PROCFILE: opened '%s' on fd %d", filename, ff->fd);
448 +
449 //strncpyz(ff->filename, filename, FILENAME_MAX);
450 ff->filename[0] = '\0';
444 -
451 ff->flags = flags;
452
453 // do not do the separators again if NULL is given
src/socket.c
+87 -47
@@ -78,6 +78,29 @@ int sock_enlarge_out(int fd) {
78 }
79
80
81 +// --------------------------------------------------------------------------------------------------------------------
82 +
83 +char *strdup_client_description(int family, const char *protocol, const char *ip, int port) {
84 + char buffer[100 + 1];
85 +
86 + switch(family) {
87 + case AF_INET:
88 + snprintfz(buffer, 100, "%s:%s:%d", protocol, ip, port);
89 + break;
90 +
91 + case AF_INET6:
92 + default:
93 + snprintfz(buffer, 100, "%s:[%s]:%d", protocol, ip, port);
94 + break;
95 +
96 + case AF_UNIX:
97 + snprintfz(buffer, 100, "%s:%s", protocol, ip);
98 + break;
99 + }
100 +
101 + return strdupz(buffer);
102 +}
103 +
104 // --------------------------------------------------------------------------------------------------------------------
105 // listening sockets
106
@@ -231,25 +254,7 @@ static inline int listen_sockets_add(LISTEN_SOCKETS *sockets, int fd, int family
254 sockets->fds[sockets->opened] = fd;
255 sockets->fds_types[sockets->opened] = socktype;
256 sockets->fds_families[sockets->opened] = family;
234 -
235 - char buffer[100 + 1];
236 -
237 - switch(family) {
238 - case AF_INET:
239 - snprintfz(buffer, 100, "%s:%s:%d", protocol, ip, port);
240 - break;
241 -
242 - case AF_INET6:
243 - default:
244 - snprintfz(buffer, 100, "%s:[%s]:%d", protocol, ip, port);
245 - break;
246 -
247 - case AF_UNIX:
248 - snprintfz(buffer, 100, "%s:%s", protocol, ip);
249 - break;
250 - }
251 -
252 - sockets->fds_names[sockets->opened] = strdupz(buffer);
257 + sockets->fds_names[sockets->opened] = strdup_client_description(family, protocol, ip, port);
258
259 sockets->opened++;
260 return 0;
@@ -842,7 +847,7 @@ ssize_t send_timeout(int sockfd, void *buf, size_t len, int flags, int timeout)
847 // --------------------------------------------------------------------------------------------------------------------
848 // accept4() replacement for systems that do not have one
849
845 -#ifndef HAVE_ACCEPT4
850 +//#ifndef HAVE_ACCEPT4
851 int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flags) {
852 int fd = accept(sock, addr, addrlen);
853 int newflags = 0;
@@ -877,7 +882,7 @@ int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flags) {
882
883 return fd;
884 }
880 -#endif
885 +//#endif
886
887
888 // --------------------------------------------------------------------------------------------------------------------
@@ -961,11 +966,16 @@ int accept_socket(int fd, int flags, char *client_ip, size_t ipsize, char *clien
966
967 struct pollinfo {
968 size_t slot;
964 - char *client;
969 + char *client_ip;
970 + char *client_port;
971 struct pollinfo *next;
972 uint32_t flags;
973 int socktype;
974
975 + void (*del_callback)(int fd, int socktype, void *data);
976 + int (*rcv_callback)(int fd, int socktype, void *data, short int *events);
977 + int (*snd_callback)(int fd, int socktype, void *data, short int *events);
978 +
979 void *data;
980 };
981
@@ -978,13 +988,13 @@ struct poll {
988 struct pollinfo *inf;
989 struct pollinfo *first_free;
990
981 - void *(*add_callback)(int fd, int socktype, short int *events);
991 + void *(*add_callback)(int fd, int socktype, short int *events, const char *client_ip, const char *client_port);
992 void (*del_callback)(int fd, int socktype, void *data);
993 int (*rcv_callback)(int fd, int socktype, void *data, short int *events);
994 int (*snd_callback)(int fd, int socktype, void *data, short int *events);
995 };
996
987 -static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype, short int events, uint32_t flags) {
997 +static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype, short int events, uint32_t flags, const char *client_ip, const char *client_port) {
998 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);
999
1000 if(unlikely(fd < 0)) return NULL;
@@ -996,6 +1006,7 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1006 p->fds = reallocz(p->fds, sizeof(struct pollfd) * new_slots);
1007 p->inf = reallocz(p->inf, sizeof(struct pollinfo) * new_slots);
1008
1009 + // reset all the newly added slots
1010 ssize_t i;
1011 for(i = new_slots - 1; i >= (ssize_t)p->slots ; i--) {
1012 debug(D_POLLFD, "POLLFD: ADD: resetting new slot %zd", i);
@@ -1006,8 +1017,15 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1017 p->inf[i].slot = (size_t)i;
1018 p->inf[i].flags = 0;
1019 p->inf[i].socktype = -1;
1009 - p->inf[i].client = NULL;
1020 + p->inf[i].client_ip = NULL;
1021 + p->inf[i].client_port = NULL;
1022 + p->inf[i].del_callback = p->del_callback;
1023 + p->inf[i].rcv_callback = p->rcv_callback;
1024 + p->inf[i].snd_callback = p->snd_callback;
1025 p->inf[i].data = NULL;
1026 +
1027 + // link them so that the first free will be earlier in the array
1028 + // (we loop decrementing i)
1029 p->inf[i].next = p->first_free;
1030 p->first_free = &p->inf[i];
1031 }
@@ -1028,13 +1046,19 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1046 pi->socktype = socktype;
1047 pi->flags = flags;
1048 pi->next = NULL;
1049 + pi->client_ip = strdupz(client_ip);
1050 + pi->client_port = strdupz(client_port);
1051 +
1052 + pi->del_callback = p->del_callback;
1053 + pi->rcv_callback = p->rcv_callback;
1054 + pi->snd_callback = p->snd_callback;
1055
1056 p->used++;
1057 if(unlikely(pi->slot > p->max))
1058 p->max = pi->slot;
1059
1060 if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1037 - pi->data = p->add_callback(fd, pi->socktype, &pf->events);
1061 + pi->data = p->add_callback(fd, pi->socktype, &pf->events, client_ip, client_port);
1062 }
1063
1064 if(pi->flags & POLLINFO_FLAG_SERVER_SOCKET) {
@@ -1053,9 +1077,10 @@ static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1077 if(unlikely(pf->fd == -1)) return;
1078
1079 if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1056 - p->del_callback(pf->fd, pi->socktype, pi->data);
1080 + pi->del_callback(pf->fd, pi->socktype, pi->data);
1081 }
1082
1083 + info("POLLFD: closing fd %d", pf->fd);
1084 close(pf->fd);
1085 pf->fd = -1;
1086 pf->events = 0;
@@ -1065,14 +1090,21 @@ static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1090 pi->flags = 0;
1091 pi->data = NULL;
1092
1068 - freez(pi->client);
1069 - pi->client = NULL;
1093 + pi->del_callback = NULL;
1094 + pi->rcv_callback = NULL;
1095 + pi->snd_callback = NULL;
1096 +
1097 + freez(pi->client_ip);
1098 + pi->client_ip = NULL;
1099 +
1100 + freez(pi->client_port);
1101 + pi->client_port = NULL;
1102
1103 pi->next = p->first_free;
1104 p->first_free = pi;
1105
1106 p->used--;
1075 - if(p->max == pi->slot) {
1107 + if(unlikely(p->max == pi->slot)) {
1108 p->max = p->min;
1109 ssize_t i;
1110 for(i = (ssize_t)pi->slot; i > (ssize_t)p->min ;i--) {
@@ -1086,10 +1118,12 @@ static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1118 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);
1119 }
1120
1089 -static void *add_callback_default(int fd, int socktype, short int *events) {
1121 +static void *add_callback_default(int fd, int socktype, short int *events, const char *client_ip, const char *client_port) {
1122 (void)fd;
1123 (void)socktype;
1124 (void)events;
1125 + (void)client_ip;
1126 + (void)client_port;
1127
1128 return NULL;
1129 }
@@ -1138,7 +1172,7 @@ static int snd_callback_default(int fd, int socktype, void *data, short int *eve
1172 return 0;
1173 }
1174
1141 -void poll_events_cleanup(void *data) {
1175 +static void poll_events_cleanup(void *data) {
1176 struct poll *p = (struct poll *)data;
1177
1178 size_t i;
@@ -1152,7 +1186,7 @@ void poll_events_cleanup(void *data) {
1186 }
1187
1188 void poll_events(LISTEN_SOCKETS *sockets
1155 - , void *(*add_callback)(int fd, int socktype, short int *events)
1189 + , void *(*add_callback)(int fd, int socktype, short int *events, const char *client_ip, const char *client_port)
1190 , void (*del_callback)(int fd, int socktype, void *data)
1191 , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
1192 , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
@@ -1177,7 +1211,7 @@ void poll_events(LISTEN_SOCKETS *sockets
1211
1212 size_t i;
1213 for(i = 0; i < sockets->opened ;i++) {
1180 - struct pollinfo *pi = poll_add_fd(&p, sockets->fds[i], sockets->fds_types[i], POLLIN, POLLINFO_FLAG_SERVER_SOCKET);
1214 + struct pollinfo *pi = poll_add_fd(&p, sockets->fds[i], sockets->fds_types[i], POLLIN, POLLINFO_FLAG_SERVER_SOCKET, (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN", "");
1215 pi->data = data;
1216 info("POLLFD: LISTENER: listening on '%s'", (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN");
1217 }
@@ -1207,7 +1241,7 @@ void poll_events(LISTEN_SOCKETS *sockets
1241 struct pollfd *pf = &p.fds[i];
1242 struct pollinfo *pi = &p.inf[i];
1243 int fd = pf->fd;
1210 - short int revents = pf->revents;
1244 + short int events = pf->events, revents = pf->revents;
1245 pf->revents = 0;
1246
1247 if(unlikely(fd == -1)) {
@@ -1215,7 +1249,7 @@ void poll_events(LISTEN_SOCKETS *sockets
1249 continue;
1250 }
1251
1218 - debug(D_POLLFD, "POLLFD: LISTENER: processing events for slot %zu (events = %d, revents = %d)", i, pf->events, revents);
1252 + debug(D_POLLFD, "POLLFD: LISTENER: processing events for slot %zu (events = %d, revents = %d)", i, events, revents);
1253
1254 if(revents & POLLIN || revents & POLLPRI) {
1255 // receiving data
@@ -1236,7 +1270,7 @@ void poll_events(LISTEN_SOCKETS *sockets
1270
1271 debug(D_POLLFD, "POLLFD: LISTENER: calling accept4() slot %zu (fd %d)", i, fd);
1272 nfd = accept_socket(fd, SOCK_NONBLOCK, client_ip, NI_MAXHOST + 1, client_port, NI_MAXSERV + 1, access_list);
1239 - if (nfd < 0) {
1273 + if (unlikely(nfd < 0)) {
1274 // accept failed
1275
1276 debug(D_POLLFD, "POLLFD: LISTENER: accept4() slot %zu (fd %d) failed.", i, fd);
@@ -1248,14 +1282,14 @@ void poll_events(LISTEN_SOCKETS *sockets
1282 }
1283 else {
1284 // accept ok
1251 - info("POLLFD: LISTENER: client '[%s]:%s' connected to '%s'", client_ip, client_port, sockets->fds_names[i]);
1252 - poll_add_fd(&p, nfd, SOCK_STREAM, POLLIN, POLLINFO_FLAG_CLIENT_SOCKET);
1285 + info("POLLFD: LISTENER: client '[%s]:%s' connected to '%s' on fd %d", client_ip, client_port, sockets->fds_names[i], nfd);
1286 + poll_add_fd(&p, nfd, SOCK_STREAM, POLLIN, POLLINFO_FLAG_CLIENT_SOCKET, client_ip, client_port);
1287
1254 - // it may have realloced them, so refresh our pointers
1288 + // it may have reallocated them, so refresh our pointers
1289 pf = &p.fds[i];
1290 pi = &p.inf[i];
1291 }
1258 - } while (nfd != -1);
1292 + } while (nfd >= 0);
1293 break;
1294 }
1295
@@ -1267,7 +1301,8 @@ void poll_events(LISTEN_SOCKETS *sockets
1301
1302 // FIXME: access_list is not applied to UDP
1303
1270 - p.rcv_callback(fd, pi->socktype, pi->data, &pf->events);
1304 + pf->events = 0;
1305 + pi->rcv_callback(fd, pi->socktype, pi->data, &pf->events);
1306 break;
1307 }
1308
@@ -1282,7 +1317,8 @@ void poll_events(LISTEN_SOCKETS *sockets
1317 // read data from client TCP socket
1318 debug(D_POLLFD, "POLLFD: LISTENER: reading data from TCP client slot %zu (fd %d)", i, fd);
1319
1285 - if (p.rcv_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
1320 + pf->events = 0;
1321 + if (pi->rcv_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
1322 poll_close_fd(&p, pi);
1323 continue;
1324 }
@@ -1293,26 +1329,30 @@ void poll_events(LISTEN_SOCKETS *sockets
1329 // sending data
1330 debug(D_POLLFD, "POLLFD: LISTENER: sending data to socket on slot %zu (fd %d)", i, fd);
1331
1296 - if (p.snd_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
1332 + pf->events = 0;
1333 + if (pi->snd_callback(fd, pi->socktype, pi->data, &pf->events) == -1) {
1334 poll_close_fd(&p, pi);
1335 continue;
1336 }
1337 }
1338
1339 if(unlikely(revents & POLLERR)) {
1303 - error("POLLFD: LISTENER: processing POLLERR events for slot %zu (events = %d, revents = %d)", i, pf->events, revents);
1340 + error("POLLFD: LISTENER: processing POLLERR events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1341 + pf->events = 0;
1342 poll_close_fd(&p, pi);
1343 continue;
1344 }
1345
1346 if(unlikely(revents & POLLHUP)) {
1309 - error("POLLFD: LISTENER: processing POLLHUP events for slot %zu (events = %d, revents = %d)", i, pf->events, pf->revents);
1347 + error("POLLFD: LISTENER: processing POLLHUP events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1348 + pf->events = 0;
1349 poll_close_fd(&p, pi);
1350 continue;
1351 }
1352
1353 if(unlikely(revents & POLLNVAL)) {
1315 - error("POLLFD: LISTENER: processing POLLNVAP events for slot %zu (events = %d, revents = %d)", i, pf->events, revents);
1354 + error("POLLFD: LISTENER: processing POLLNVAL events for slot %zu fd %d (events = %d, revents = %d)", i, events, revents, fd);
1355 + pf->events = 0;
1356 poll_close_fd(&p, pi);
1357 continue;
1358 }
src/socket.h
+3 -1
@@ -19,6 +19,8 @@ typedef struct listen_sockets {
19 int fds_families[MAX_LISTEN_FDS]; // the family of the open sockets (AF_UNIX, AF_INET, AF_INET6)
20 } LISTEN_SOCKETS;
21
22 +extern char *strdup_client_description(int family, const char *protocol, const char *ip, int port);
23 +
24 extern int listen_sockets_setup(LISTEN_SOCKETS *sockets);
25 extern void listen_sockets_close(LISTEN_SOCKETS *sockets);
26
@@ -52,7 +54,7 @@ extern int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flag
54
55
56 extern void poll_events(LISTEN_SOCKETS *sockets
55 - , void *(*add_callback)(int fd, int socktype, short int *events)
57 + , void *(*add_callback)(int fd, int socktype, short int *events, const char *client_ip, const char *client_port)
58 , void (*del_callback)(int fd, int socktype, void *data)
59 , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
60 , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
src/statsd.c
+4 -1
@@ -697,9 +697,12 @@ struct statsd_udp {
697 #endif
698
699 // new TCP client connected
700 -static void *statsd_add_callback(int fd, int socktype, short int *events) {
700 +static void *statsd_add_callback(int fd, int socktype, short int *events, const char *client_ip, const char *client_port) {
701 (void)fd;
702 (void)socktype;
703 + (void)client_ip;
704 + (void)client_port;
705 +
706 *events = POLLIN;
707
708 struct statsd_tcp *data = (struct statsd_tcp *)callocz(sizeof(struct statsd_tcp) + STATSD_TCP_BUFFER_SIZE, 1);
src/sys_fs_cgroup.c
+2 -2
@@ -532,14 +532,14 @@ static inline void cgroup_read_cpuacct_usage(struct cpuacct_usage *ca) {
532 }
533
534 static inline void cgroup_read_blkio(struct blkio *io) {
535 - static procfile *ff = NULL;
536 -
535 if(unlikely(io->enabled == CONFIG_BOOLEAN_AUTO && io->delay_counter > 0)) {
536 io->delay_counter--;
537 return;
538 }
539
540 if(likely(io->filename)) {
541 + static procfile *ff = NULL;
542 +
543 ff = procfile_reopen(ff, io->filename, NULL, PROCFILE_FLAG_DEFAULT);
544 if(unlikely(!ff)) {
545 io->updated = 0;
src/web_client.c
+47 -3
@@ -84,7 +84,51 @@ static void web_client_update_acl_matches(struct web_client *w) {
84 w->acl |= WEB_CLIENT_ACL_BADGE;
85 }
86
87 -struct web_client *web_client_create(int listener) {
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;
94 +
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 + }
114 +
115 + web_client_update_acl_matches(w);
116 +
117 + w->response.data = buffer_create(INITIAL_WEB_DATA_LENGTH);
118 + w->response.header = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
119 + w->response.header_output = buffer_create(HTTP_RESPONSE_HEADER_SIZE);
120 + w->origin[0] = '*';
121 + web_client_enable_wait_receive(w);
122 +
123 + if(web_clients) web_clients->prev = w;
124 + w->next = web_clients;
125 + web_clients = w;
126 +
127 + web_client_connected();
128 + return(w);
129 +}
130 +
131 +struct web_client *web_client_create_on_listenfd(int listener) {
132 struct web_client *w;
133
134 w = callocz(1, sizeof(struct web_client));
@@ -115,11 +159,11 @@ struct web_client *web_client_create(int listener) {
159
160 int flag = 1;
161 if(setsockopt(w->ofd, IPPROTO_TCP, TCP_NODELAY, (char *) &flag, sizeof(int)) != 0)
118 - error("%llu: failed to enable TCP_NODELAY on socket.", w->id);
162 + error("%llu: failed to enable TCP_NODELAY on socket fd %d.", w->id, w->ofd);
163
164 flag = 1;
165 if(setsockopt(w->ifd, SOL_SOCKET, SO_KEEPALIVE, (char *) &flag, sizeof(int)) != 0)
122 - error("%llu: Cannot set SO_KEEPALIVE on socket.", w->id);
166 + error("%llu: Cannot set SO_KEEPALIVE on socket fd %d.", w->id, w->ifd);
167 }
168
169 web_client_update_acl_matches(w);
src/web_client.h
+2 -1
@@ -174,7 +174,8 @@ extern uid_t web_files_gid(void);
174
175 extern int web_client_permission_denied(struct web_client *w);
176
177 -extern struct web_client *web_client_create(int listener);
177 +extern struct web_client *web_client_create_on_fd(int fd, const char *client_ip, const char *client_port);
178 +extern struct web_client *web_client_create_on_listenfd(int listener);
179 extern struct web_client *web_client_free(struct web_client *w);
180 extern ssize_t web_client_send(struct web_client *w);
181 extern ssize_t web_client_receive(struct web_client *w);
src/web_server.c
+36 -27
@@ -53,6 +53,8 @@ WEB_SERVER_MODE web_server_mode_id(const char *mode) {
53 return WEB_SERVER_MODE_NONE;
54 else if(!strcmp(mode, "single") || !strcmp(mode, "single-threaded"))
55 return WEB_SERVER_MODE_SINGLE_THREADED;
56 + else if(!strcmp(mode, "static") || !strcmp(mode, "static-threaded"))
57 + return WEB_SERVER_MODE_STATIC_THREADED;
58 else // if(!strcmp(mode, "multi") || !strcmp(mode, "multi-threaded"))
59 return WEB_SERVER_MODE_MULTI_THREADED;
60 }
@@ -65,6 +67,9 @@ const char *web_server_mode_name(WEB_SERVER_MODE id) {
67 case WEB_SERVER_MODE_SINGLE_THREADED:
68 return "single-threaded";
69
70 + case WEB_SERVER_MODE_STATIC_THREADED:
71 + return "static-threaded";
72 +
73 default:
74 case WEB_SERVER_MODE_MULTI_THREADED:
75 return "multi-threaded";
@@ -83,9 +88,9 @@ int api_listen_sockets_setup(void) {
88 }
89
90 // --------------------------------------------------------------------------------------
86 -// the main socket listener
91 +// the main socket listener - MULTI-THREADED
92
88 -static inline void cleanup_web_clients(void) {
93 +static inline void multi_threaded_cleanup_web_clients(void) {
94 struct web_client *w;
95
96 for (w = web_clients; w;) {
@@ -171,7 +176,7 @@ void *socket_listen_main_multi_threaded(void *ptr) {
176 else if(unlikely(!retval)) {
177 debug(D_WEB_CLIENT, "LISTENER: select() timeout.");
178 counter = 0;
174 - cleanup_web_clients();
179 + multi_threaded_cleanup_web_clients();
180 continue;
181 }
182
@@ -182,9 +187,9 @@ void *socket_listen_main_multi_threaded(void *ptr) {
187 if(revents & POLLIN || revents & POLLPRI) {
188 socket_listen_main_multi_threaded_fds[i].revents = 0;
189
185 - w = web_client_create(socket_listen_main_multi_threaded_fds[i].fd);
190 + w = web_client_create_on_listenfd(socket_listen_main_multi_threaded_fds[i].fd);
191 if(unlikely(!w)) {
187 - // no need for error log - web_client_create already logged the error
192 + // no need for error log - web_client_create_on_listenfd already logged the error
193 continue;
194 }
195
@@ -205,7 +210,7 @@ void *socket_listen_main_multi_threaded(void *ptr) {
210 counter++;
211 if(counter >= CLEANUP_EVERY_EVENTS) {
212 counter = 0;
208 - cleanup_web_clients();
213 + multi_threaded_cleanup_web_clients();
214 }
215 }
216
@@ -213,6 +218,9 @@ void *socket_listen_main_multi_threaded(void *ptr) {
218 return NULL;
219 }
220
221 +// --------------------------------------------------------------------------------------
222 +// the main socket listener - SINGLE-THREADED
223 +
224 struct web_client *single_threaded_clients[FD_SETSIZE];
225
226 static inline int single_threaded_link_client(struct web_client *w, fd_set *ifds, fd_set *ofds, fd_set *efds, int *max) {
@@ -323,7 +331,7 @@ void *socket_listen_main_single_threaded(void *ptr) {
331 for(i = 0; i < api_sockets.opened ; i++) {
332 if (FD_ISSET(api_sockets.fds[i], &rifds)) {
333 debug(D_WEB_CLIENT_ACCESS, "LISTENER: new connection.");
326 - w = web_client_create(api_sockets.fds[i]);
334 + w = web_client_create_on_listenfd(api_sockets.fds[i]);
335
336 if(api_sockets.fds_families[i] == AF_UNIX)
337 web_client_set_unix(w);
@@ -391,17 +399,18 @@ void *socket_listen_main_single_threaded(void *ptr) {
399 return NULL;
400 }
401
402 +// --------------------------------------------------------------------------------------
403 +// the main socket listener - STATIC-THREADED
404
395 -#if 0
405 // new TCP client connected
397 -static void *web_server_add_callback(int fd, int socktype, short int *events) {
406 +static void *web_server_add_callback(int fd, int socktype, short int *events, const char *client_ip, const char *client_port) {
407 (void)fd;
408 (void)socktype;
409
410 *events = POLLIN;
411
412 debug(D_WEB_CLIENT_ACCESS, "LISTENER on %d: new connection.", fd);
404 - struct web_client *w = web_client_create(fd);
413 + struct web_client *w = web_client_create_on_fd(fd, client_ip, client_port);
414
415 if(unlikely(socktype == AF_UNIX))
416 web_client_set_unix(w);
@@ -418,15 +427,24 @@ static void web_server_del_callback(int fd, int socktype, void *data) {
427
428 struct web_client *w = (struct web_client *)data;
429
421 - if(w) {
430 + if(likely(w)) {
431 if(w->ofd == -1 || fd == w->ofd) {
432 // we free the client, only if the closing fd
433 // is the client socket
434 +
435 + // prevent it from closing the file descriptors
436 + w->ifd = w->ofd = -1;
437 +
438 web_client_free(w);
439 }
440 }
441 +}
442
429 - return;
443 +static inline int web_server_check_client_status(struct web_client *w) {
444 + if(unlikely(web_client_check_obsolete(w) || web_client_check_dead(w) || (!web_client_has_wait_receive(w) && !web_client_has_wait_send(w))))
445 + return -1;
446 +
447 + return 0;
448 }
449
450 // Receive data
@@ -434,8 +452,6 @@ static int web_server_rcv_callback(int fd, int socktype, void *data, short int *
452 (void)fd;
453 (void)socktype;
454
437 - *events = 0;
438 -
455 struct web_client *w = (struct web_client *)data;
456
457 if(unlikely(!web_client_has_wait_receive(w)))
@@ -466,10 +482,7 @@ static int web_server_rcv_callback(int fd, int socktype, void *data, short int *
482 if(unlikely(w->ofd == fd && web_client_has_wait_send(w)))
483 *events |= POLLOUT;
484
469 - if(unlikely(*events == 0))
470 - return -1;
471 -
472 - return 0;
485 + return web_server_check_client_status(w);
486 }
487
488 static int web_server_snd_callback(int fd, int socktype, void *data, short int *events) {
@@ -490,13 +503,10 @@ static int web_server_snd_callback(int fd, int socktype, void *data, short int *
503 if(unlikely(w->ofd == fd && web_client_has_wait_send(w)))
504 *events |= POLLOUT;
505
493 - if(unlikely(*events == 0))
494 - return -1;
495 -
496 - return 0;
506 + return web_server_check_client_status(w);
507 }
508
499 -static void socket_listen_main_single_threaded_cleanup(void *ptr) {
509 +static void socket_listen_main_static_threaded_cleanup(void *ptr) {
510 struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
511 if(static_thread->enabled) {
512 static_thread->enabled = 0;
@@ -506,9 +516,9 @@ static void socket_listen_main_single_threaded_cleanup(void *ptr) {
516 }
517 }
518
509 -void *socket_listen_main_single_threaded(void *ptr) {
510 - netdata_thread_cleanup_push(socket_listen_main_single_threaded_cleanup, ptr);
511 - web_server_mode = WEB_SERVER_MODE_SINGLE_THREADED;
519 +void *socket_listen_main_static_threaded(void *ptr) {
520 + netdata_thread_cleanup_push(socket_listen_main_static_threaded_cleanup, ptr);
521 + web_server_mode = WEB_SERVER_MODE_STATIC_THREADED;
522
523 if(!api_sockets.opened)
524 fatal("LISTENER: no listen sockets available.");
@@ -525,4 +535,3 @@ void *socket_listen_main_single_threaded(void *ptr) {
535 netdata_thread_cleanup_pop(1);
536 return NULL;
537 }
528 -#endif
src/web_server.h
+2
@@ -16,6 +16,7 @@
16
17 typedef enum web_server_mode {
18 WEB_SERVER_MODE_SINGLE_THREADED,
19 + WEB_SERVER_MODE_STATIC_THREADED,
20 WEB_SERVER_MODE_MULTI_THREADED,
21 WEB_SERVER_MODE_NONE
22 } WEB_SERVER_MODE;
@@ -27,6 +28,7 @@ extern const char *web_server_mode_name(WEB_SERVER_MODE id);
28
29 extern void *socket_listen_main_multi_threaded(void *ptr);
30 extern void *socket_listen_main_single_threaded(void *ptr);
31 +extern void *socket_listen_main_static_threaded(void *ptr);
32 extern int api_listen_sockets_setup(void);
33
34 #endif /* NETDATA_WEB_SERVER_H */