diff --git a/async.c b/async.c index fa787da3..4bb35b0c 100644 --- a/async.c +++ b/async.c @@ -176,23 +176,17 @@ int async_del(int queuefd, int fd, int etype) { #else int async_queue_init(int serverfd) { - int kfd; - struct kevent kev; - kfd = kqueue(); + int eqfd = event_queue_init(); + if (eqfd < 0) { + exit(1); + } - if (kfd < 0) { - uwsgi_error("kqueue()"); - return -1; + if (event_queue_add_fd_read(eqfd, serverfd) < 0) { + exit(1); } - EV_SET(&kev, serverfd, EVFILT_READ, EV_ADD, 0, 0, 0); - if (kevent(kfd, &kev, 1, NULL, 0, NULL) < 0) { - uwsgi_error("kevent()"); - return -1; - } - - return kfd; + return eqfd ; } int async_wait(int queuefd, void *events, int nevents, int block, int timeout) { diff --git a/buildconf/default.ini b/buildconf/default.ini index 9639494a..e73ffe6d 100644 --- a/buildconf/default.ini +++ b/buildconf/default.ini @@ -27,6 +27,11 @@ bin_name = uwsgi plugin_dir = . embedded_plugins = python, ping, nagios +locking = auto +event = auto +timer = auto +filemonitor = auto + [python] paste = true web3 = true diff --git a/event.c b/event.c new file mode 100644 index 00000000..2bde5010 --- /dev/null +++ b/event.c @@ -0,0 +1,99 @@ +#include "uwsgi.h" + +#ifdef UWSGI_EVENT_USE_KQUEUE +int event_queue_init() { + + int kfd = kqueue(); + + if (kfd < 0) { + uwsgi_error("kqueue()"); + return -1; + } + + return kfd; +} + +int event_queue_add_fd_read(int eq, int fd) { + + struct kevent kev; + + EV_SET(&kev, fd, EVFILT_READ, EV_ADD, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} + +int event_queue_add_fd_write(int eq, int fd) { + + struct kevent kev; + + EV_SET(&kev, fd, EVFILT_WRITE, EV_ADD, 0, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} + +int event_queue_wait(int eq, int timeout, int *interesting_fd) { + + int ret; + struct timespec ts; + struct kevent ev; + + if (timeout <= 0) { + ret = kevent(eq, NULL, 0, &ev, 1, NULL); + } + else { + memset(&ts, 0, sizeof(struct timespec)); + ts.tv_sec = timeout; + ret = kevent(eq, NULL, 0, &ev, 1, &ts); + } + + if (ret < 0) { + uwsgi_error("kevent()"); + } + + if (ret > 0) { + *interesting_fd = ev.ident; + uwsgi_log("FFLAGS: %d\n", ev.fflags); + } + + return ret; + +} +#endif + +#ifdef UWSGI_EVENT_FILEMONITOR_USE_KQUEUE +int event_queue_add_file_monitor(int eq, int fd) { + + struct kevent kev; + + EV_SET(&kev, fd, EVFILT_VNODE, EV_ADD|EV_CLEAR, NOTE_WRITE|NOTE_DELETE|NOTE_EXTEND|NOTE_ATTRIB|NOTE_RENAME|NOTE_REVOKE, 0, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} +#endif + +#ifdef UWSGI_EVENT_TIMER_USE_KQUEUE +int event_queue_add_timer(int eq, int id, int sec) { + + struct kevent kev; + + EV_SET(&kev, id, EVFILT_TIMER, EV_ADD, 0, sec*1000, 0); + if (kevent(eq, &kev, 1, NULL, 0, NULL) < 0) { + uwsgi_error("kevent()"); + return -1; + } + + return 0; +} +#endif diff --git a/master.c b/master.c index 52c63d1e..70c48921 100644 --- a/master.c +++ b/master.c @@ -61,7 +61,6 @@ void master_loop(char **argv, char **environ) { int master_has_children = 0; - struct pollfd *uwsgi_signal_poll; char uwsgi_signal; #ifdef UWSGI_UDP @@ -99,17 +98,11 @@ void master_loop(char **argv, char **environ) { signal(SIGUSR1, (void *) &stats); - uwsgi_signal_poll = malloc(sizeof(struct pollfd) * uwsgi.numproc); - if (!uwsgi_signal_poll) { - uwsgi_error("malloc()"); - exit(1); - } - memset(uwsgi_signal_poll, 0, sizeof(struct pollfd) * uwsgi.numproc); + int master_queue = event_queue_init(); for(i=1;i<=uwsgi.numproc;i++) { uwsgi_log("adding %d to signal poll\n", uwsgi.workers[i].pipe[0]); - uwsgi_signal_poll[i-1].fd = uwsgi.workers[i].pipe[0]; - uwsgi_signal_poll[i-1].events = POLLIN; + event_queue_add_fd_read(master_queue, uwsgi.workers[i].pipe[0]); } uwsgi.wsgi_req->buffer = uwsgi.async_buf[0]; @@ -144,11 +137,7 @@ void master_loop(char **argv, char **environ) { } //uwsgi_log("cluster opts size: %d\n", cluster_opt_size); - cluster_opt_buf = malloc(cluster_opt_size); - if (!cluster_opt_buf) { - uwsgi_error("malloc()"); - exit(1); - } + cluster_opt_buf = uwsgi_malloc(cluster_opt_size); uh = (struct uwsgi_header *) cluster_opt_buf; @@ -206,6 +195,17 @@ void master_loop(char **argv, char **environ) { } #endif + // add a fake timer + event_queue_add_timer(master_queue, 0xFFFF, 5); + + int fd_mon = open("/tmp/topolino", O_RDONLY); + if (fd_mon < 0) { + uwsgi_error("open()"); + exit(1); + } + + event_queue_add_file_monitor(master_queue, fd_mon); + for (;;) { //uwsgi_log("ready_to_reload %d %d\n", ready_to_reload, uwsgi.numproc); if (ready_to_die >= uwsgi.numproc && uwsgi.to_hell) { @@ -309,145 +309,123 @@ void master_loop(char **argv, char **environ) { if (!check_interval) check_interval = 1; -#ifdef UWSGI_UDP -#ifdef UWSGI_MULTICAST - if ( (uwsgi.udp_socket && udp_fd >= 0) || (uwsgi.cluster && uwsgi.cluster_fd >= 0)) { -#else - if ((uwsgi.udp_socket && udp_fd >= 0)) { -#endif - rlen = poll(uwsgi_poll, uwsgi_poll_size, check_interval * 1000); + int interesting_fd = -1; + rlen = event_queue_wait(master_queue, check_interval, &interesting_fd); + if (rlen < 0) { uwsgi_error("poll()"); } else if (rlen > 0) { - for(i=0;ibuffer, uwsgi.buffer_size, 0, (struct sockaddr *) &udp_client, &udp_len); - if (uwsgi_poll[i].revents & POLLIN) { + if (rlen < 0) { + uwsgi_error("recvfrom()"); + } + else if (rlen > 0) { - if (uwsgi_poll[i].fd == udp_fd) { - udp_len = sizeof(udp_client); - rlen = recvfrom(udp_fd, uwsgi.wsgi_req->buffer, uwsgi.buffer_size, 0, (struct sockaddr *) &udp_client, &udp_len); - if (rlen < 0) { - uwsgi_error("recvfrom()"); + memset(udp_client_addr, 0, 16); + if (inet_ntop(AF_INET, &udp_client.sin_addr.s_addr, udp_client_addr, 16)) { + if (uwsgi.wsgi_req->buffer[0] == UWSGI_MODIFIER_MULTICAST_ANNOUNCE) { } - else if (rlen > 0) { - memset(udp_client_addr, 0, 16); - if (inet_ntop(AF_INET, &udp_client.sin_addr.s_addr, udp_client_addr, 16)) { - if (uwsgi.wsgi_req->buffer[0] == UWSGI_MODIFIER_MULTICAST_ANNOUNCE) { - } #ifdef UWSGI_SNMP - else if (uwsgi.wsgi_req->buffer[0] == 0x30 && uwsgi.snmp) { - manage_snmp(udp_fd, (uint8_t *) uwsgi.wsgi_req->buffer, rlen, &udp_client); - } + else if (uwsgi.wsgi_req->buffer[0] == 0x30 && uwsgi.snmp) { + manage_snmp(udp_fd, (uint8_t *) uwsgi.wsgi_req->buffer, rlen, &udp_client); + } #endif - else { + else { - // loop the various udp manager until one returns true - udp_managed = 0; - for(i=0;i<0xFF;i++) { - if (uwsgi.p[i]->manage_udp) { - if (uwsgi.p[i]->manage_udp(udp_client_addr, udp_client.sin_port, uwsgi.wsgi_req->buffer, rlen)) { - udp_managed = 1; - break; - } - } - } - /* - if (udp_callable && udp_callable_args) { - UWSGI_GET_GIL - PyTuple_SetItem(udp_callable_args, 0, PyString_FromString(udp_client_addr)); - PyTuple_SetItem(udp_callable_args, 1, PyInt_FromLong(ntohs(udp_client.sin_port))); - PyTuple_SetItem(udp_callable_args, 2, PyString_FromStringAndSize(uwsgi.wsgi_req->buffer, rlen)); - udp_response = python_call(udp_callable, udp_callable_args, 0); - if (udp_response) { - Py_DECREF(udp_response); - } - if (PyErr_Occurred()) - PyErr_Print(); - - UWSGI_RELEASE_GIL - } - else { - // a simple udp logger - */ - if (!udp_managed) { - uwsgi_log( "[udp:%s:%d] %.*s", udp_client_addr, ntohs(udp_client.sin_port), rlen, uwsgi.wsgi_req->buffer); + // loop the various udp manager until one returns true + udp_managed = 0; + for(i=0;i<0xFF;i++) { + if (uwsgi.p[i]->manage_udp) { + if (uwsgi.p[i]->manage_udp(udp_client_addr, udp_client.sin_port, uwsgi.wsgi_req->buffer, rlen)) { + udp_managed = 1; + break; } } } - else { - uwsgi_error("inet_ntop()"); + + // else a simple udp logger + if (!udp_managed) { + uwsgi_log( "[udp:%s:%d] %.*s", udp_client_addr, ntohs(udp_client.sin_port), rlen, uwsgi.wsgi_req->buffer); } } } - - if (uwsgi_poll[i].fd == uwsgi.cluster_fd) { - - if (uwsgi_get_dgram(uwsgi.cluster_fd, uwsgi.wsgi_requests[0])) { - continue; - } - - switch(uwsgi.wsgi_requests[0]->uh.modifier1) { - case 95: - new_cluster_hostname = NULL; - new_cluster_address = NULL; - new_cluster_workers = NULL; - uwsgi_hooked_parse(uwsgi.wsgi_requests[0]->buffer, uwsgi.wsgi_requests[0]->uh.pktsize, print_dict, NULL); - if (new_cluster_hostname && new_cluster_address && new_cluster_workers) { - uwsgi_cluster_add_node(new_cluster_address, atoi(new_cluster_workers), CLUSTER_NODE_DYNAMIC); - } - break; - case 96: - uwsgi_log_verbose("%.*s\n", uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); - break; - case 98: - if (kill(getpid(), SIGHUP)) { - uwsgi_error("kill()"); - } - break; - case 99: - if (uwsgi.wsgi_requests[0]->uh.modifier2 == 0) { - uwsgi_log("requested configuration data, sending %d bytes\n", cluster_opt_size); - sendto(uwsgi.cluster_fd, cluster_opt_buf, cluster_opt_size, 0, (struct sockaddr *) &uwsgi.mc_cluster_addr, sizeof(uwsgi.mc_cluster_addr)); - } - break; - case 73: - uwsgi_log_verbose("[uWSGI cluster %s] new node available: %.*s\n", uwsgi.cluster, uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); - break; - } - } + else { + uwsgi_error("inet_ntop()"); + } } + + continue; } - } - } - else { -#endif - rlen = poll(uwsgi_signal_poll, uwsgi.numproc, check_interval*1000); - if (rlen < 0) { - uwsgi_error("poll()"); - continue; - } - else if (rlen > 0) { - for(i=0;iuh.modifier1) { + case 95: + new_cluster_hostname = NULL; + new_cluster_address = NULL; + new_cluster_workers = NULL; + uwsgi_hooked_parse(uwsgi.wsgi_requests[0]->buffer, uwsgi.wsgi_requests[0]->uh.pktsize, print_dict, NULL); + if (new_cluster_hostname && new_cluster_address && new_cluster_workers) { + uwsgi_cluster_add_node(new_cluster_address, atoi(new_cluster_workers), CLUSTER_NODE_DYNAMIC); + } + break; + case 96: + uwsgi_log_verbose("%.*s\n", uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); + break; + case 98: + if (kill(getpid(), SIGHUP)) { + uwsgi_error("kill()"); + } + break; + case 99: + if (uwsgi.wsgi_requests[0]->uh.modifier2 == 0) { + uwsgi_log("requested configuration data, sending %d bytes\n", cluster_opt_size); + sendto(uwsgi.cluster_fd, cluster_opt_buf, cluster_opt_size, 0, (struct sockaddr *) &uwsgi.mc_cluster_addr, sizeof(uwsgi.mc_cluster_addr)); + } + break; + case 73: + uwsgi_log_verbose("[uWSGI cluster %s] new node available: %.*s\n", uwsgi.cluster, uwsgi.wsgi_requests[0]->uh.pktsize, uwsgi.wsgi_requests[0]->buffer); + break; + } + continue; + } + + if (interesting_fd == fd_mon) { + uwsgi_log("FileSystem event !!!\n"); + continue; + } + + if (interesting_fd == 0xFFFF) { + uwsgi_log("timer elapsed !!!\n"); + continue; + } + + // finally check for uwsgi_signal + for(i=1;i<=uwsgi.numproc;i++) { + if (interesting_fd == uwsgi.workers[i].pipe[0]) { + rlen = read(interesting_fd, &uwsgi_signal, 1); if (rlen < 0) { uwsgi_error("read()"); } else if (rlen > 0) { - uwsgi_log("received uwsgi signal %d from worker %d\n", uwsgi_signal, i+1); + uwsgi_log("received uwsgi signal %d from worker %d\n", uwsgi_signal, i); } else { - uwsgi_log_verbose("lost connection with worker %d\n", i+1); - close(uwsgi_signal_poll[i].fd); + uwsgi_log_verbose("lost connection with worker %d\n", i); + close(interesting_fd); } } } } -#ifdef UWSGI_UDP - } -#endif current_time = time(NULL); // checking logsize @@ -673,7 +651,7 @@ void master_loop(char **argv, char **environ) { else { uwsgi_log( "Respawned uWSGI worker (new pid: %d)\n", pid); close(uwsgi.workers[uwsgi.mywid].pipe[1]); - uwsgi_signal_poll[uwsgi.mywid-1].fd = uwsgi.workers[uwsgi.mywid].pipe[0]; + event_queue_add_fd_read(master_queue, uwsgi.workers[uwsgi.mywid].pipe[0]); #ifdef UWSGI_SPOOLER if (uwsgi.mywid <= 0 && diedpid != uwsgi.shared->spooler_pid) { #else diff --git a/uwsgi.h b/uwsgi.h index 7d201720..47e9bd65 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -1258,3 +1258,12 @@ void uwsgi_lock(void *); void uwsgi_unlock(void *); inline void *uwsgi_malloc(size_t); + + +int event_queue_init(void); +int event_queue_add_fd_read(int, int); +int event_queue_wait(int, int, int *); + +int event_queue_add_timer(int, int, int); + +int event_queue_add_file_monitor(int, int); diff --git a/uwsgi_API.txt b/uwsgi_API.txt index a6129e73..0fb8a761 100644 --- a/uwsgi_API.txt +++ b/uwsgi_API.txt @@ -50,3 +50,4 @@ sharedarea_writebyte sharedarea_readlong sharedarea_writelong sharedarea_inclong +register_signal(num, handler) diff --git a/uwsgiconfig.py b/uwsgiconfig.py index 8b7bf8a7..0ecb5b19 100644 --- a/uwsgiconfig.py +++ b/uwsgiconfig.py @@ -136,7 +136,7 @@ class uConf(): def __init__(self, filename): self.config = ConfigParser.ConfigParser() self.config.read(filename) - self.gcc_list = ['utils', 'protocol', 'socket', 'logging', 'master', 'plugins', 'lock', 'cache', 'loop', 'uwsgi'] + self.gcc_list = ['utils', 'protocol', 'socket', 'logging', 'master', 'plugins', 'lock', 'cache', 'event', 'loop', 'uwsgi'] self.cflags = ['-O2', '-Wall', '-Werror', '-D_LARGEFILE_SOURCE', '-D_FILE_OFFSET_BITS=64'] + os.environ.get("CFLAGS", "").split() gcc_version = str(spcall2("%s -v" % GCC)).split('\n')[-1].split()[2] gcc_major = int(gcc_version.split('.')[0]) @@ -191,7 +191,6 @@ class uConf(): # set locking subsystem locking_mode = self.get('locking','auto') - print locking_mode, uwsgi_os if locking_mode == 'auto': if uwsgi_os == 'Linux': locking_mode = 'pthread_mutex' @@ -208,6 +207,48 @@ class uConf(): self.cflags.append('-DUWSGI_LOCK_USE_OSX_SPINLOCK') else: self.cflags.append('-DUWSGI_LOCK_USE_FLOCK') + + # set event subsystem + event_mode = self.get('event','auto') + + if event_mode == 'auto': + if uwsgi_os == 'Linux': + event_mode = 'epoll' + elif uwsgi_os in ('Darwin', 'FreeBSD'): + event_mode = 'kqueue' + + if event_mode == 'epoll': + self.cflags.append('-DUWSGI_EVENT_USE_EPOLL') + elif event_mode == 'kqueue': + self.cflags.append('-DUWSGI_EVENT_USE_KQUEUE') + + # set timer subsystem + timer_mode = self.get('timer','auto') + + if timer_mode == 'auto': + if uwsgi_os == 'Linux': + timer_mode = 'timerfd' + elif uwsgi_os in ('Darwin', 'FreeBSD'): + timer_mode = 'kqueue' + + if timer_mode == 'timerfd': + self.cflags.append('-DUWSGI_EVENT_TIMER_USE_TIMERFD') + elif timer_mode == 'kqueue': + self.cflags.append('-DUWSGI_EVENT_TIMER_USE_KQUEUE') + + # set timer subsystem + filemonitor_mode = self.get('filemonitor','auto') + + if filemonitor_mode == 'auto': + if uwsgi_os == 'Linux': + filemonitor_mode = 'inotify' + elif uwsgi_os in ('Darwin', 'FreeBSD'): + filemonitor_mode = 'kqueue' + + if filemonitor_mode == 'inotify': + self.cflags.append('-DUWSGI_EVENT_FILEMONITOR_USE_INOTIFY') + elif filemonitor_mode == 'kqueue': + self.cflags.append('-DUWSGI_EVENT_FILEMONITOR_USE_KQUEUE') if self.get('embedded'):