@cryptotaxi247 / netdata-1 / commits / cf3a8218e

working static threaded web servers - with issues - not ready yet

Costa Tsaousis (ktsaou) committed Jan 8, 2018 at 03:25 UTC cf3a8218e9f0953f983c4b910a72dc29898c4c18
5 files changed +138 -21
src/socket.c
+46 -10
@@ -961,7 +961,18 @@ int accept_socket(int fd, int flags, char *client_ip, size_t ipsize, char *clien
961
962 #define POLL_FDS_INCREASE_STEP 10
963
964 -static inline POLLINFO *poll_add_fd(POLLJOB *p, int fd, int socktype, short int events, uint32_t flags, const char *client_ip, const char *client_port) {
964 +inline POLLINFO *poll_add_fd(POLLJOB *p
965 + , int fd
966 + , int socktype
967 + , uint32_t flags
968 + , const char *client_ip
969 + , const char *client_port
970 + , void *(*add_callback)(POLLINFO *pi, short int *events, void *data)
971 + , void (*del_callback)(POLLINFO *pi)
972 + , int (*rcv_callback)(POLLINFO *pi, short int *events)
973 + , int (*snd_callback)(POLLINFO *pi, short int *events)
974 + , void *data
975 +) {
976 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);
977
978 if(unlikely(fd < 0)) return NULL;
@@ -1008,7 +1019,7 @@ static inline POLLINFO *poll_add_fd(POLLJOB *p, int fd, int socktype, short int
1019
1020 struct pollfd *pf = &p->fds[pi->slot];
1021 pf->fd = fd;
1011 - pf->events = events;
1022 + pf->events = POLLIN;
1023 pf->revents = 0;
1024
1025 pi->fd = fd;
@@ -1019,16 +1030,16 @@ static inline POLLINFO *poll_add_fd(POLLJOB *p, int fd, int socktype, short int
1030 pi->client_ip = strdupz(client_ip);
1031 pi->client_port = strdupz(client_port);
1032
1022 - pi->del_callback = p->del_callback;
1023 - pi->rcv_callback = p->rcv_callback;
1024 - pi->snd_callback = p->snd_callback;
1033 + pi->del_callback = del_callback;
1034 + pi->rcv_callback = rcv_callback;
1035 + pi->snd_callback = snd_callback;
1036
1037 p->used++;
1038 if(unlikely(pi->slot > p->max))
1039 p->max = pi->slot;
1040
1041 if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1031 - pi->data = p->add_callback(pi, &pf->events);
1042 + pi->data = add_callback(pi, &pf->events, data);
1043 }
1044
1045 if(pi->flags & POLLINFO_FLAG_SERVER_SOCKET) {
@@ -1091,9 +1102,10 @@ static inline void poll_close_fd(POLLINFO *pi) {
1102 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);
1103 }
1104
1094 -void *poll_default_add_callback(POLLINFO *pi, short int *events) {
1105 +void *poll_default_add_callback(POLLINFO *pi, short int *events, void *data) {
1106 (void)pi;
1107 (void)events;
1108 + (void)data;
1109
1110 return NULL;
1111 }
@@ -1191,7 +1203,18 @@ static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, sho
1203 else {
1204 // accept ok
1205 // info("POLLFD: LISTENER: client '[%s]:%s' connected to '%s' on fd %d", client_ip, client_port, sockets->fds_names[i], nfd);
1194 - poll_add_fd(p, nfd, SOCK_STREAM, POLLIN, POLLINFO_FLAG_CLIENT_SOCKET, client_ip, client_port);
1206 + poll_add_fd(p
1207 + , nfd
1208 + , SOCK_STREAM
1209 + , POLLINFO_FLAG_CLIENT_SOCKET
1210 + , client_ip
1211 + , client_port
1212 + , p->add_callback
1213 + , p->del_callback
1214 + , p->rcv_callback
1215 + , p->snd_callback
1216 + , NULL
1217 + );
1218
1219 // it may have reallocated them, so refresh our pointers
1220 pf = &p->fds[i];
@@ -1267,7 +1290,7 @@ static void poll_events_process(POLLJOB *p, POLLINFO *pi, struct pollfd *pf, sho
1290 }
1291
1292 void poll_events(LISTEN_SOCKETS *sockets
1270 - , void *(*add_callback)(POLLINFO *pi, short int *events)
1293 + , void *(*add_callback)(POLLINFO *pi, short int *events, void *data)
1294 , void (*del_callback)(POLLINFO *pi)
1295 , int (*rcv_callback)(POLLINFO *pi, short int *events)
1296 , int (*snd_callback)(POLLINFO *pi, short int *events)
@@ -1294,7 +1317,20 @@ void poll_events(LISTEN_SOCKETS *sockets
1317
1318 size_t i;
1319 for(i = 0; i < sockets->opened ;i++) {
1297 - 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", "");
1320 +
1321 + POLLINFO *pi = poll_add_fd(&p
1322 + , sockets->fds[i]
1323 + , sockets->fds_types[i]
1324 + , POLLINFO_FLAG_SERVER_SOCKET
1325 + , (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN"
1326 + , ""
1327 + , poll_default_add_callback
1328 + , poll_default_del_callback
1329 + , poll_default_rcv_callback
1330 + , poll_default_snd_callback
1331 + , NULL
1332 + );
1333 +
1334 pi->data = data;
1335 info("POLLFD: LISTENER: listening on '%s'", (sockets->fds_names[i])?sockets->fds_names[i]:"UNKNOWN");
1336 }
src/socket.h
+16 -3
@@ -96,7 +96,7 @@ typedef struct poll {
96
97 SIMPLE_PATTERN *access_list;
98
99 - void *(*add_callback)(POLLINFO *pi, short int *events);
99 + void *(*add_callback)(POLLINFO *pi, short int *events, void *data);
100 void (*del_callback)(POLLINFO *pi);
101 int (*rcv_callback)(POLLINFO *pi, short int *events);
102 int (*snd_callback)(POLLINFO *pi, short int *events);
@@ -105,10 +105,23 @@ typedef struct poll {
105 extern int poll_default_snd_callback(POLLINFO *pi, short int *events);
106 extern int poll_default_rcv_callback(POLLINFO *pi, short int *events);
107 extern void poll_default_del_callback(POLLINFO *pi);
108 -extern void *poll_default_add_callback(POLLINFO *pi, short int *events);
108 +extern void *poll_default_add_callback(POLLINFO *pi, short int *events, void *data);
109 +
110 +extern POLLINFO *poll_add_fd(POLLJOB *p
111 + , int fd
112 + , int socktype
113 + , uint32_t flags
114 + , const char *client_ip
115 + , const char *client_port
116 + , void *(*add_callback)(POLLINFO *pi, short int *events, void *data)
117 + , void (*del_callback)(POLLINFO *pi)
118 + , int (*rcv_callback)(POLLINFO *pi, short int *events)
119 + , int (*snd_callback)(POLLINFO *pi, short int *events)
120 + , void *data
121 +);
122
123 extern void poll_events(LISTEN_SOCKETS *sockets
111 - , void *(*add_callback)(POLLINFO *pi, short int *events)
124 + , void *(*add_callback)(POLLINFO *pi, short int *events, void *data)
125 , void (*del_callback)(POLLINFO *pi)
126 , int (*rcv_callback)(POLLINFO *pi, short int *events)
127 , int (*snd_callback)(POLLINFO *pi, short int *events)
src/statsd.c
+2 -1
@@ -697,8 +697,9 @@ struct statsd_udp {
697 #endif
698
699 // new TCP client connected
700 -static void *statsd_add_callback(POLLINFO *pi, short int *events) {
700 +static void *statsd_add_callback(POLLINFO *pi, short int *events, void *data) {
701 (void)pi;
702 + (void)data;
703
704 *events = POLLIN;
705
src/web_client.c
+4 -1
@@ -1778,7 +1778,10 @@ ssize_t web_client_receive(struct web_client *w)
1778 web_client_disable_wait_receive(w);
1779
1780 debug(D_WEB_CLIENT, "%llu: Read the whole file.", w->id);
1781 - if(w->ifd != w->ofd) close(w->ifd);
1781 +
1782 + if(w->ifd != w->ofd && web_server_mode != WEB_SERVER_MODE_STATIC_THREADED)
1783 + close(w->ifd);
1784 +
1785 w->ifd = w->ofd;
1786 }
1787 else {
src/web_server.c
+70 -6
@@ -422,9 +422,11 @@ static struct web_server_static_threaded_worker *static_workers_private_data = N
422 static __thread struct web_server_static_threaded_worker *worker_private = NULL;
423
424 // new TCP client connected
425 -static void *web_server_add_callback(POLLINFO *pi, short int *events) {
425 +static void *web_server_add_callback(POLLINFO *pi, short int *events, void *data) {
426 + (void)data;
427
428 worker_private->connected++;
429 +
430 size_t concurrent = worker_private->connected - worker_private->disconnected;
431 if(unlikely(concurrent > worker_private->max_concurrent))
432 worker_private->max_concurrent = concurrent;
@@ -462,6 +464,48 @@ static void web_server_del_callback(POLLINFO *pi) {
464 }
465 }
466
467 +// ----------------------------------------------------------------------------
468 +// web server files
469 +
470 +struct web_file_pollinfo {
471 + struct web_client *w;
472 + POLLINFO *pi;
473 +};
474 +
475 +static void *web_server_file_add_callback(POLLINFO *pi, short int *events, void *data) {
476 + (void)pi;
477 +
478 + info("ADDING FILE ON FD %d", pi->fd);
479 + *events = POLLIN;
480 + return data;
481 +}
482 +
483 +static void web_werver_file_del_callback(POLLINFO *pi) {
484 + (void)pi;
485 + info("DELETE FILE ON FD %d", pi->fd);
486 + freez(pi->data);
487 +}
488 +
489 +static int web_server_file_rcv_callback(POLLINFO *pi, short int *events) {
490 + *events = POLLIN;
491 +
492 + info("READING FILE ON FD %d", pi->fd);
493 +
494 + struct web_file_pollinfo *wfpi = (struct web_file_pollinfo *)pi->data;
495 + if(unlikely(web_client_receive(wfpi->w) < 0))
496 + return -1;
497 +
498 + wfpi->pi->p->fds[wfpi->pi->slot].events |= POLLOUT;
499 +
500 + if(unlikely(wfpi->w->ifd == wfpi->w->ofd))
501 + return -1;
502 +
503 + return 0;
504 +}
505 +
506 +// ----------------------------------------------------------------------------
507 +//
508 +
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;
@@ -482,21 +526,41 @@ static int web_server_rcv_callback(POLLINFO *pi, short int *events) {
526 if(unlikely(web_client_receive(w) < 0))
527 return -1;
528
529 + debug(D_WEB_CLIENT, "%llu: Processing received data.", w->id);
530 + web_client_process_request(w);
531 +
532 if(unlikely(w->mode == WEB_CLIENT_MODE_FILECOPY)) {
486 - if(unlikely(w->ifd != -1 && w->ifd != fd)) {
533 + info("FILECOPY %d", pi->fd);
534 +
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);
539 +
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
545 + , w->ifd
546 + , 0
547 + , POLLINFO_FLAG_CLIENT_SOCKET
548 + , "FILENAME"
549 + , ""
550 + , web_server_file_add_callback
551 + , 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;
562 }
563 }
496 - else {
497 - debug(D_WEB_CLIENT, "%llu: Processing received data.", w->id);
498 - web_client_process_request(w);
499 - }
564
565 if(unlikely(w->ifd == fd && web_client_has_wait_receive(w)))
566 *events |= POLLIN;