From 70db2befb5ec226f764b5843d087378349a3abb9 Mon Sep 17 00:00:00 2001 From: "roberto@fiorenzo" Date: Sun, 19 Dec 2010 12:50:36 +0100 Subject: [PATCH] fully working file monitor and timer support for Linux/epoll --- event.c | 111 +++++++++++++++++++++++++++----- master.c | 72 +++++++++++++++------ plugins/python/uwsgi_pymodule.c | 31 ++++++++- signal.c | 43 +++++++++++++ uwsgi.c | 2 + uwsgi.h | 27 +++++++- uwsgiconfig.py | 2 +- 7 files changed, 250 insertions(+), 38 deletions(-) create mode 100644 signal.c diff --git a/event.c b/event.c index 4d069e9c..40836911 100644 --- a/event.c +++ b/event.c @@ -1,5 +1,7 @@ #include "uwsgi.h" +extern struct uwsgi_server uwsgi; + #ifdef UWSGI_EVENT_USE_EPOLL #include @@ -32,7 +34,7 @@ int event_queue_add_fd_read(int eq, int fd) { return -1; } - return 0; + return fd; } int event_queue_wait(int eq, int timeout, int *interesting_fd) { @@ -81,7 +83,7 @@ int event_queue_add_fd_read(int eq, int fd) { return -1; } - return 0; + return fd; } int event_queue_add_fd_write(int eq, int fd) { @@ -145,30 +147,93 @@ int event_queue_add_file_monitor(int eq, int fd) { int event_queue_add_file_monitor(int eq, char *filename, int *id) { - int ifd = inotify_init(); - if (ifd < 0) { - uwsgi_error("inotify_init()"); - return -1; - } + int ifd = -1; + int i; + int add_to_queue = 0; - *id = ifd; + for (i=0;i sizeof(struct inotify_event)) { + bie = uwsgi_malloc(isize); + rlen = read(id, bie, isize); + } + else { + rlen = read(id, &ie, sizeof(struct inotify_event)); + bie = &ie; + } if (rlen < 0) { uwsgi_error("read()"); } + else { + items = isize/(sizeof(struct inotify_event)); + uwsgi_log("inotify returned %d items\n", items); + for(j=0;jwd) { + if (hook) { + hook(uwsgi.files_monitored[i].filename, iie->mask, iie->name); + } + else { + uf = &uwsgi.files_monitored[i]; + } + } + } + } + + } + + if (items > 1) { + free(bie); + } + + return uf; + } + + return NULL; } #endif @@ -203,17 +268,33 @@ int event_queue_add_timer(int eq, int *id, int sec) { return event_queue_add_fd_read(eq, tfd); } -void event_queue_ack_timer(int id) { +struct uwsgi_timer *event_queue_ack_timer(int id, void hook(int, int)) { + int i; ssize_t rlen; uint64_t counter; + struct uwsgi_timer *ut = NULL; + + for(i=0;ivalue, ut->id); + } + } + return ut; } #endif diff --git a/master.c b/master.c index b9739307..1c85dc30 100644 --- a/master.c +++ b/master.c @@ -98,11 +98,11 @@ void master_loop(char **argv, char **environ) { signal(SIGUSR1, (void *) &stats); - int master_queue = event_queue_init(); + uwsgi.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]); - event_queue_add_fd_read(master_queue, uwsgi.workers[i].pipe[0]); + event_queue_add_fd_read(uwsgi.master_queue, uwsgi.workers[i].pipe[0]); } uwsgi.wsgi_req->buffer = uwsgi.async_buf[0]; @@ -200,10 +200,24 @@ void master_loop(char **argv, char **environ) { /* int fake_timer = 0xFFFF; event_queue_add_timer(master_queue, &fake_timer, 5); - - int fd_mon = -1; - event_queue_add_file_monitor(master_queue, "/tmp/topolino", &fd_mon); */ + + + // add unregistered file monitors + for(i=0;i 0) { @@ -395,19 +409,41 @@ void master_loop(char **argv, char **environ) { continue; } - /* - if (interesting_fd == fd_mon) { - uwsgi_log("FileSystem event !!!\n"); - event_queue_ack_file_monitor(fd_mon); - continue; - } + + //event_queue_ack_timer(fake_timer); - if (interesting_fd == fake_timer) { - uwsgi_log("timer elapsed !!!\n"); - event_queue_ack_timer(fake_timer); - continue; + int next_iteration = 0; + + for(i=0;ifilename); + } + break; + } + } } - */ + if (next_iteration) continue; + + next_iteration = 0; + + for(i=0;ivalue); + } + break; + } + } + } + if (next_iteration) continue; + // finally check for uwsgi_signal for(i=1;i<=uwsgi.numproc;i++) { @@ -652,7 +688,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]); - event_queue_add_fd_read(master_queue, uwsgi.workers[uwsgi.mywid].pipe[0]); + event_queue_add_fd_read(uwsgi.master_queue, uwsgi.workers[uwsgi.mywid].pipe[0]); #ifdef UWSGI_SPOOLER if (uwsgi.mywid <= 0 && diedpid != uwsgi.shared->spooler_pid) { #else diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index 45b909ec..8f513ae6 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -159,13 +159,40 @@ PyObject *py_uwsgi_signal(PyObject * self, PyObject * args) { char uwsgi_signal; ssize_t rlen; + char *payload = NULL; + struct uwsgi_header uh; - if (!PyArg_ParseTuple(args, "B:signal", &uwsgi_signal)) { + if (!PyArg_ParseTuple(args, "B|s:signal", &uwsgi_signal, &payload)) { return NULL; } + // am i the master ? + if (uwsgi.mywid == 0) { + uwsgi_log("i am the master !!!\n"); + register_signal(uwsgi_signal, payload); + goto done; + } + uwsgi_log("sending %d to master\n", uwsgi_signal); - rlen = write(uwsgi.workers[uwsgi.mywid].pipe[1], &uwsgi_signal, 1); + uh.modifier1 = 110; + if (payload) { + uh.pktsize = strlen(payload); + } + else { + uh.pktsize = 0; + } + uh.modifier2 = uwsgi_signal; + + rlen = write(uwsgi.workers[uwsgi.mywid].pipe[1], &uh, 4); + if (rlen == 4) { + if (uh.pktsize) { + if (write(uwsgi.workers[uwsgi.mywid].pipe[1], payload, uh.pktsize) < 0) { + uwsgi_error("write()"); + } + } + } + +done: Py_INCREF(Py_None); return Py_None; diff --git a/signal.c b/signal.c new file mode 100644 index 00000000..38ec5d9d --- /dev/null +++ b/signal.c @@ -0,0 +1,43 @@ +#include "uwsgi.h" + +extern struct uwsgi_server uwsgi; + +int register_signal(uint8_t sig, char *payload) { + + uwsgi_log("SIGNAL %d %s\n", sig, payload); + + switch(sig) { + + case 10: + if (uwsgi.files_monitored_cnt < 64) { + uwsgi.files_monitored[uwsgi.files_monitored_cnt].filename = uwsgi_concat2(payload,""); + uwsgi.files_monitored[uwsgi.files_monitored_cnt].registered = 0; + // master is not running + if (uwsgi.master_queue != -1) { + uwsgi.files_monitored[uwsgi.files_monitored_cnt].fd = event_queue_add_file_monitor(uwsgi.master_queue, payload, &uwsgi.files_monitored[uwsgi.files_monitored_cnt].id); + uwsgi.files_monitored[uwsgi.files_monitored_cnt].registered = 1; + } + uwsgi.files_monitored_cnt++; + } + else { + uwsgi_log("you can register max 64 file monitors !!!\n"); + } + + case 11: + if (uwsgi.timers_cnt < 64) { + uwsgi.timers[uwsgi.timers_cnt].value = atoi(payload); + uwsgi.timers[uwsgi.timers_cnt].registered = 0; + // master is not running + if (uwsgi.master_queue != -1) { + uwsgi.timers[uwsgi.timers_cnt].fd = event_queue_add_timer(uwsgi.master_queue, &uwsgi.timers[uwsgi.timers_cnt].id, uwsgi.timers[uwsgi.timers_cnt].value); + uwsgi.timers[uwsgi.timers_cnt].registered = 1; + } + uwsgi.timers_cnt++; + } + else { + uwsgi_log("you can register max 64 timers !!!\n"); + } + } + + return 0; +} diff --git a/uwsgi.c b/uwsgi.c index 76b005e9..90484276 100644 --- a/uwsgi.c +++ b/uwsgi.c @@ -418,6 +418,8 @@ int main(int argc, char *argv[], char *envp[]) uwsgi.p[i] = &unconfigured_plugin; } + uwsgi.master_queue = -1; + uwsgi.cluster_fd = -1; uwsgi.cores = 1; diff --git a/uwsgi.h b/uwsgi.h index 5d0ba6f0..8b32fa54 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -588,6 +588,20 @@ struct wsgi_request { #define LOADER_MAX 8 +struct uwsgi_fmon { + char *filename; + int fd; + int id; + int registered; +}; + +struct uwsgi_timer { + int value; + int fd; + int id; + int registered; +}; + struct uwsgi_server { @@ -696,6 +710,7 @@ struct uwsgi_server { int post_buffering_bufsize; int master_process; + int master_queue; int no_defer_accept; @@ -843,6 +858,12 @@ struct uwsgi_server { void *cache_lock; void *user_lock; + + struct uwsgi_fmon files_monitored[64]; + int files_monitored_cnt; + + struct uwsgi_timer timers[64]; + int timers_cnt; }; struct uwsgi_lb_group { @@ -1265,8 +1286,10 @@ int event_queue_add_fd_read(int, int); int event_queue_wait(int, int, int *); int event_queue_add_timer(int, int *, int); -void event_queue_ack_timer(int); +struct uwsgi_timer *event_queue_ack_timer(int, void (*)(int, int)); int event_queue_add_file_monitor(int, char *, int *); -void event_queue_ack_file_monitor(int); +struct uwsgi_fmon *event_queue_ack_file_monitor(int, void (*)(char *, uint32_t, char *)); + +int register_signal(uint8_t, char *); diff --git a/uwsgiconfig.py b/uwsgiconfig.py index a36c21e5..45d0d3f9 100644 --- a/uwsgiconfig.py +++ b/uwsgiconfig.py @@ -137,7 +137,7 @@ class uConf(object): def __init__(self, filename): self.config = ConfigParser.ConfigParser() self.config.read(filename) - self.gcc_list = ['utils', 'protocol', 'socket', 'logging', 'master', 'plugins', 'lock', 'cache', 'event', 'loop', 'uwsgi'] + self.gcc_list = ['utils', 'protocol', 'socket', 'logging', 'master', 'plugins', 'lock', 'cache', 'event', 'signal', '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])