@cryptotaxi247 / netdata-1 / commits / df18e07f2

provide socktype to poll_events() add/del events

Costa Tsaousis (ktsaou) committed Sep 24, 2017 at 21:09 UTC df18e07f2b8a441c1abe77258789b4ad28936097
4 files changed +155 -12
src/socket.c
+10 -8
@@ -952,8 +952,8 @@ struct poll {
952 struct pollinfo *inf;
953 struct pollinfo *first_free;
954
955 - void *(*add_callback)(int fd, short int *events);
956 - void (*del_callback)(int fd, void *data);
955 + void *(*add_callback)(int fd, int socktype, short int *events);
956 + void (*del_callback)(int fd, int socktype, void *data);
957 int (*rcv_callback)(int fd, int socktype, void *data, short int *events);
958 int (*snd_callback)(int fd, int socktype, void *data, short int *events);
959 };
@@ -1008,7 +1008,7 @@ static inline struct pollinfo *poll_add_fd(struct poll *p, int fd, int socktype,
1008 p->max = pi->slot;
1009
1010 if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1011 - pi->data = p->add_callback(fd, &pf->events);
1011 + pi->data = p->add_callback(fd, pi->socktype, &pf->events);
1012 }
1013
1014 if(pi->flags & POLLINFO_FLAG_SERVER_SOCKET) {
@@ -1027,7 +1027,7 @@ static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1027 if(unlikely(pf->fd == -1)) return;
1028
1029 if(pi->flags & POLLINFO_FLAG_CLIENT_SOCKET) {
1030 - p->del_callback(pf->fd, pi->data);
1030 + p->del_callback(pf->fd, pi->socktype, pi->data);
1031 }
1032
1033 close(pf->fd);
@@ -1060,14 +1060,16 @@ static inline void poll_close_fd(struct poll *p, struct pollinfo *pi) {
1060 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);
1061 }
1062
1063 -static void *add_callback_default(int fd, short int *events) {
1063 +static void *add_callback_default(int fd, int socktype, short int *events) {
1064 (void)fd;
1065 + (void)socktype;
1066 (void)events;
1067
1068 return NULL;
1069 }
1069 -static void del_callback_default(int fd, void *data) {
1070 +static void del_callback_default(int fd, int socktype, void *data) {
1071 (void)fd;
1072 + (void)socktype;
1073 (void)data;
1074
1075 if(data)
@@ -1124,8 +1126,8 @@ void poll_events_cleanup(void *data) {
1126 }
1127
1128 void poll_events(LISTEN_SOCKETS *sockets
1127 - , void *(*add_callback)(int fd, short int *events)
1128 - , void (*del_callback)(int fd, void *data)
1129 + , void *(*add_callback)(int fd, int socktype, short int *events)
1130 + , void (*del_callback)(int fd, int socktype, void *data)
1131 , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
1132 , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
1133 , SIMPLE_PATTERN *access_list
src/socket.h
+2 -2
@@ -52,8 +52,8 @@ extern int accept4(int sock, struct sockaddr *addr, socklen_t *addrlen, int flag
52
53
54 extern void poll_events(LISTEN_SOCKETS *sockets
55 - , void *(*add_callback)(int fd, short int *events)
56 - , void (*del_callback)(int fd, void *data)
55 + , void *(*add_callback)(int fd, int socktype, short int *events)
56 + , void (*del_callback)(int fd, int socktype, void *data)
57 , int (*rcv_callback)(int fd, int socktype, void *data, short int *events)
58 , int (*snd_callback)(int fd, int socktype, void *data, short int *events)
59 , SIMPLE_PATTERN *access_list
src/statsd.c
+4 -2
@@ -685,8 +685,9 @@ struct statsd_udp {
685 #endif
686
687 // new TCP client connected
688 -static void *statsd_add_callback(int fd, short int *events) {
688 +static void *statsd_add_callback(int fd, int socktype, short int *events) {
689 (void)fd;
690 + (void)socktype;
691 *events = POLLIN;
692
693 struct statsd_tcp *data = (struct statsd_tcp *)callocz(sizeof(struct statsd_tcp) + STATSD_TCP_BUFFER_SIZE, 1);
@@ -697,8 +698,9 @@ static void *statsd_add_callback(int fd, short int *events) {
698 }
699
700 // TCP client disconnected
700 -static void statsd_del_callback(int fd, void *data) {
701 +static void statsd_del_callback(int fd, int socktype, void *data) {
702 (void)fd;
703 + (void)socktype;
704
705 if(data) {
706 struct statsd_tcp *t = data;
src/web_server.c
+139
@@ -377,3 +377,142 @@ void *socket_listen_main_single_threaded(void *ptr) {
377 pthread_exit(NULL);
378 return NULL;
379 }
380 +
381 +
382 +#if 0
383 +// new TCP client connected
384 +static void *web_server_add_callback(int fd, int socktype, short int *events) {
385 + (void)fd;
386 + (void)socktype;
387 +
388 + *events = POLLIN;
389 +
390 + debug(D_WEB_CLIENT_ACCESS, "LISTENER on %d: new connection.", fd);
391 + struct web_client *w = web_client_create(fd);
392 +
393 + if(unlikely(socktype == AF_UNIX))
394 + web_client_set_unix(w);
395 + else
396 + web_client_set_tcp(w);
397 +
398 + return (void *)w;
399 +}
400 +
401 +// TCP client disconnected
402 +static void web_server_del_callback(int fd, int socktype, void *data) {
403 + (void)fd;
404 + (void)socktype;
405 +
406 + struct web_client *w = (struct web_client *)data;
407 +
408 + if(w) {
409 + if(w->ofd == -1 || fd == w->ofd) {
410 + // we free the client, only if the closing fd
411 + // is the client socket
412 + web_client_free(w);
413 + }
414 + }
415 +
416 + return;
417 +}
418 +
419 +// Receive data
420 +static int web_server_rcv_callback(int fd, int socktype, void *data, short int *events) {
421 + (void)fd;
422 + (void)socktype;
423 +
424 + *events = 0;
425 +
426 + struct web_client *w = (struct web_client *)data;
427 +
428 + if(unlikely(!web_client_has_wait_receive(w)))
429 + return -1;
430 +
431 + if(unlikely(web_client_receive(w) < 0))
432 + return -1;
433 +
434 + if(unlikely(w->mode == WEB_CLIENT_MODE_FILECOPY)) {
435 + if(unlikely(w->ifd != -1 && w->ifd != fd)) {
436 + // FIXME: we switched input fd
437 + // add a new socket to poll_events, with the same
438 + }
439 + else if(unlikely(w->ifd == -1)) {
440 + // FIXME: we closed input fd
441 + // instruct poll_events() to close fd
442 + return -1;
443 + }
444 + }
445 + else {
446 + debug(D_WEB_CLIENT, "%llu: Processing received data.", w->id);
447 + web_client_process_request(w);
448 + }
449 +
450 + if(unlikely(w->ifd == fd && web_client_has_wait_receive(w)))
451 + *events |= POLLIN;
452 +
453 + if(unlikely(w->ofd == fd && web_client_has_wait_send(w)))
454 + *events |= POLLOUT;
455 +
456 + if(unlikely(*events == 0))
457 + return -1;
458 +
459 + return 0;
460 +}
461 +
462 +static int web_server_snd_callback(int fd, int socktype, void *data, short int *events) {
463 + (void)fd;
464 + (void)socktype;
465 +
466 + struct web_client *w = (struct web_client *)data;
467 +
468 + if(unlikely(!web_client_has_wait_send(w)))
469 + return -1;
470 +
471 + if(unlikely(web_client_send(w) < 0))
472 + return -1;
473 +
474 + if(unlikely(w->ifd == fd && web_client_has_wait_receive(w)))
475 + *events |= POLLIN;
476 +
477 + if(unlikely(w->ofd == fd && web_client_has_wait_send(w)))
478 + *events |= POLLOUT;
479 +
480 + if(unlikely(*events == 0))
481 + return -1;
482 +
483 + return 0;
484 +}
485 +
486 +void *socket_listen_main_single_threaded(void *ptr) {
487 + struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
488 +
489 + web_server_mode = WEB_SERVER_MODE_SINGLE_THREADED;
490 +
491 + info("Single-threaded WEB SERVER thread created with task id %d", gettid());
492 +
493 + if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
494 + error("Cannot set pthread cancel type to DEFERRED.");
495 +
496 + if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
497 + error("Cannot set pthread cancel state to ENABLE.");
498 +
499 + if(!api_sockets.opened)
500 + fatal("LISTENER: no listen sockets available.");
501 +
502 + poll_events(&api_sockets
503 + , web_server_add_callback
504 + , web_server_del_callback
505 + , web_server_rcv_callback
506 + , web_server_snd_callback
507 + , web_allow_connections_from
508 + , NULL
509 + );
510 +
511 + debug(D_WEB_CLIENT, "LISTENER: exit!");
512 + listen_sockets_close(&api_sockets);
513 +
514 + static_thread->enabled = 0;
515 + pthread_exit(NULL);
516 + return NULL;
517 +}
518 +#endif