diff --git a/event.c b/event.c index 61089fb7..a1b3f606 100644 --- a/event.c +++ b/event.c @@ -151,25 +151,19 @@ int event_queue_add_file_monitor(int eq, char *filename, int *id) { return fd; } -struct uwsgi_fmon *event_queue_ack_file_monitor(int id, void hook(char *, uint32_t, char *)) { +struct uwsgi_fmon *event_queue_ack_file_monitor(int id) { int i; - struct uwsgi_fmon *uf = NULL; - for(i=0;ifiles_monitored_cnt;i++) { - if (uwsgi.shared->files_monitored[i].registered) { - if (uwsgi.shared->files_monitored[i].fd == id) { - if (hook) { - hook(uwsgi.shared->files_monitored[i].filename, 0, NULL); - } - else { - uf = &uwsgi.shared->files_monitored[i]; - } + for(i=0;ifiles_monitored_cnt;i++) { + if (ushared->files_monitored[i].registered) { + if (ushared->files_monitored[i].fd == id) { + return &ushared->files_monitored[i]; } } } - return uf; + return NULL } @@ -184,9 +178,9 @@ int event_queue_add_file_monitor(int eq, char *filename, int *id) { int i; int add_to_queue = 0; - for (i=0;ifiles_monitored_cnt;i++) { + if (ushared->files_monitored[i].registered) { + ifd = ushared->files_monitored[0].fd; break; } } @@ -212,7 +206,7 @@ int event_queue_add_file_monitor(int eq, char *filename, int *id) { } } -struct uwsgi_fmon *event_queue_ack_file_monitor(int id, void hook(char *, uint32_t, char *)) { +struct uwsgi_fmon *event_queue_ack_file_monitor(int id) { ssize_t rlen = 0; struct inotify_event ie, *bie, *iie; @@ -241,18 +235,13 @@ struct uwsgi_fmon *event_queue_ack_file_monitor(int id, void hook(char *, uint32 } else { items = isize/(sizeof(struct inotify_event)); - //uwsgi_log("inotify returned %d items\n", items); + 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]; - } + for(i=0;ifiles_monitored_cnt;i++) { + if (ushared->files_monitored[i].registered) { + if (ushared->files_monitored[i].fd == id && ushared->files_monitored[i].id == iie->wd) { + uf = &ushared->files_monitored[i]; } } } @@ -313,17 +302,17 @@ int event_queue_add_timer(int eq, int *id, int sec) { return event_queue_add_fd_read(eq, tfd); } -struct uwsgi_timer *event_queue_ack_timer(int id, void hook(int, int)) { +struct uwsgi_timer *event_queue_ack_timer(int id) { int i; ssize_t rlen; uint64_t counter; struct uwsgi_timer *ut = NULL; - for(i=0;itimers_cnt;i++) { + if (ushared->timers[i].registered) { + if (ushared->timers[i].id == id) { + ut = &ushared->timers[i]; } } } @@ -333,11 +322,6 @@ struct uwsgi_timer *event_queue_ack_timer(int id, void hook(int, int)) { if (rlen < 0) { uwsgi_error("read()"); } - else { - if (hook && ut) { - hook(ut->value, ut->id); - } - } return ut; } @@ -345,7 +329,7 @@ struct uwsgi_timer *event_queue_ack_timer(int id, void hook(int, int)) { #ifdef UWSGI_EVENT_TIMER_USE_NONE int event_queue_add_timer(int eq, int *id, int sec) { return -1; } -struct uwsgi_timer *event_queue_ack_timer(int id, void hook(int ,int)) { return NULL;} +struct uwsgi_timer *event_queue_ack_timer(int id) { return NULL;} #endif #ifdef UWSGI_EVENT_TIMER_USE_KQUEUE @@ -368,7 +352,7 @@ int event_queue_add_timer(int eq, int *id, int sec) { return *id; } -struct uwsgi_timer *event_queue_ack_timer(int id, void hook(int ,int)) { +struct uwsgi_timer *event_queue_ack_timer(int id) { int i; struct uwsgi_timer *ut = NULL; @@ -381,10 +365,6 @@ struct uwsgi_timer *event_queue_ack_timer(int id, void hook(int ,int)) { } } - if (hook) { - hook(ut->value, ut->id); - } - return ut; } diff --git a/loop.c b/loop.c index 66a8477f..8ea9cffb 100644 --- a/loop.c +++ b/loop.c @@ -66,16 +66,6 @@ void *simple_loop(void *arg1) { while (uwsgi.workers[uwsgi.mywid].manage_next_request) { -#ifndef __linux__ - if (uwsgi.no_orphans && uwsgi.master_process) { - // am i a son of init ? - if (getppid() == 1) { - uwsgi_log("UAAAAAAH my parent died :( i will follow him...\n"); - exit(1); - } - } -#endif - UWSGI_CLEAR_STATUS; diff --git a/master.c b/master.c index 4d76ec20..8e05e6e6 100644 --- a/master.c +++ b/master.c @@ -196,13 +196,6 @@ void master_loop(char **argv, char **environ) { /* - // add unregistered timers - for(i=0;ifiles_monitored_cnt;i++) { if (!ushared->files_monitored[i].registered) { ushared->files_monitored[i].fd = event_queue_add_file_monitor(uwsgi.master_queue, ushared->files_monitored[i].filename, &ushared->files_monitored[i].id); ushared->files_monitored[i].registered = 1; } } - uwsgi_unlock(uwsgi.fmon_table_lock); + + + // add unregistered timers + // locking is not needed as monitors can only increase + for(i=0;itimers_cnt;i++) { + if (!ushared->timers[i].registered) { + ushared->timers[i].fd = event_queue_add_timer(uwsgi.master_queue, &ushared->timers[i].id, ushared->timers[i].value); + ushared->timers[i].registered = 1; + } + } int interesting_fd = -1; rlen = event_queue_wait(uwsgi.master_queue, check_interval, &interesting_fd); @@ -411,15 +413,13 @@ void master_loop(char **argv, char **environ) { #endif - //event_queue_ack_timer(fake_timer); - int next_iteration = 0; uwsgi_lock(uwsgi.fmon_table_lock); for(i=0;ifiles_monitored_cnt;i++) { if (ushared->files_monitored[i].registered) { if (interesting_fd == ushared->files_monitored[i].fd) { - struct uwsgi_fmon *uf = event_queue_ack_file_monitor(interesting_fd, NULL); + struct uwsgi_fmon *uf = event_queue_ack_file_monitor(interesting_fd); // now call the file_monitor handler if (uf) { uwsgi_log("fd event for %s (signal %d)\n", uf->filename, uf->sig); @@ -437,24 +437,26 @@ void master_loop(char **argv, char **environ) { uwsgi_unlock(uwsgi.fmon_table_lock); if (next_iteration) continue; - /* next_iteration = 0; - for(i=0;itimers_cnt;i++) { + if (ushared->timers[i].registered) { //uwsgi_log("%d = %d\n", interesting_fd, uwsgi.timers[i].fd); - if (interesting_fd == uwsgi.timers[i].fd) { - struct uwsgi_timer *ut = event_queue_ack_timer(interesting_fd, NULL); + if (interesting_fd == ushared->timers[i].fd) { + struct uwsgi_timer *ut = event_queue_ack_timer(interesting_fd); // now call the file_monitor handler if (ut) { uwsgi_log("fd event for timer %d\n", ut->value); + struct uwsgi_signal_entry *use = &ushared->signal_table[ut->sig]; + if (use->kind == SIGNAL_KIND_WORKER) { + uwsgi_log("write signal returned %d\n", write(ushared->worker_signal_pipe[0], &ut->sig, 1)); + } } break; } } } if (next_iteration) continue; - */ // check for worker signal diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index 90ce04bc..e4b0aba2 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -155,6 +155,26 @@ PyObject *py_uwsgi_close(PyObject * self, PyObject * args) { } +PyObject *py_uwsgi_register_timer(PyObject * self, PyObject * args) { + + uint8_t uwsgi_signal; + uint8_t signal_kind; + PyObject *handler; + int secs; + + if (!PyArg_ParseTuple(args, "BiBO:register_timer", &uwsgi_signal, &secs, &signal_kind, &handler)) { + return NULL; + } + + uwsgi_log("signal_kind %d\n", signal_kind); + + uwsgi_register_timer(uwsgi_signal, secs, signal_kind, handler, 0); + + Py_INCREF(Py_None); + return Py_None; +} + + PyObject *py_uwsgi_register_file_monitor(PyObject * self, PyObject * args) { uint8_t uwsgi_signal; @@ -1893,7 +1913,7 @@ static PyMethodDef uwsgi_advanced_methods[] = { {"register_signal", py_uwsgi_register_signal, METH_VARARGS, ""}, {"signal", py_uwsgi_signal, METH_VARARGS, ""}, {"register_file_monitor", py_uwsgi_register_file_monitor, METH_VARARGS, ""}, - //{"register_timer", py_uwsgi_register_timer, METH_VARARGS, ""}, + {"register_timer", py_uwsgi_register_timer, METH_VARARGS, ""}, #ifdef UWSGI_SENDFILE {"sendfile", py_uwsgi_advanced_sendfile, METH_VARARGS, ""}, #endif diff --git a/signal.c b/signal.c index 1d51f9d3..498c213a 100644 --- a/signal.c +++ b/signal.c @@ -93,8 +93,8 @@ void uwsgi_register_file_monitor(uint8_t sig, char *filename, uint8_t kind, void memcpy(ushared->files_monitored[ushared->files_monitored_cnt].filename, filename, strlen(filename)); ushared->files_monitored[ushared->files_monitored_cnt].registered = 0; ushared->files_monitored[ushared->files_monitored_cnt].sig = sig; - ushared->files_monitored_cnt++; uwsgi_register_signal(sig, kind, handler, modifier1, filename, strlen(filename)); + ushared->files_monitored_cnt++; } else { uwsgi_log("you can register max 64 file monitors !!!\n"); @@ -103,3 +103,25 @@ void uwsgi_register_file_monitor(uint8_t sig, char *filename, uint8_t kind, void uwsgi_unlock(uwsgi.fmon_table_lock); } + +void uwsgi_register_timer(uint8_t sig, int secs, uint8_t kind, void *handler, uint8_t modifier1) { + + uwsgi_lock(uwsgi.timer_table_lock); + + if (ushared->timers_cnt < 64) { + + // fill the timer table, the master will use it to add items to the event queue + ushared->timers[ushared->timers_cnt].value = secs; + snprintf(ushared->timers[ushared->timers_cnt].svalue, 0xff, "%d", secs); + ushared->timers[ushared->timers_cnt].registered = 0; + ushared->timers[ushared->timers_cnt].sig = sig; + uwsgi_register_signal(sig, kind, handler, modifier1, ushared->timers[ushared->timers_cnt].svalue, strlen(ushared->timers[ushared->timers_cnt].svalue)); + ushared->timers_cnt++; + } + else { + uwsgi_log("you can register max 64 timers !!!\n"); + } + + uwsgi_unlock(uwsgi.timer_table_lock); + +} diff --git a/tests/signals.py b/tests/signals.py new file mode 100644 index 00000000..c2a905ca --- /dev/null +++ b/tests/signals.py @@ -0,0 +1,34 @@ +import uwsgi + + +def hello_signal(num, payload): + print "i am the signal %d" % num + +def hello_signal2(num, payload): + print "i am the signal %d with payload: %s" % (num, payload) + +def hello_file(num, filename): + print "file %s has been modified !!!" % filename + +def hello_timer(num, secs): + print "%s seconds elapsed" % secs + +#uwsgi.register_signal(30, uwsgi.SIGNAL_KIND_WORKER, hello_signal) +uwsgi.register_signal(30, 1, hello_signal) +uwsgi.register_signal(22, 1, hello_signal2, "*** PAYLOAD FOO ***") + +uwsgi.register_file_monitor(17, "/tmp", 1, hello_file) +uwsgi.register_timer(26, 2, 1, hello_timer) +uwsgi.register_timer(17, 4, 1, hello_timer) +uwsgi.register_timer(5, 8, 1, hello_timer) + + +def application(env, start_response): + + start_response('200 Ok', [('Content-Type', 'text/html')] ) + + # this will send a signal to the master that will report it to the first available worker + uwsgi.signal(30) + uwsgi.signal(22) + + return "signals sent to workers" diff --git a/utils.c b/utils.c index 65820392..ff47e350 100644 --- a/utils.c +++ b/utils.c @@ -463,7 +463,7 @@ int wsgi_req_accept(struct wsgi_request *wsgi_req) { if (wsgi_req->leave_open) return 0; polling: - ret = poll(uwsgi.sockets_poll, uwsgi.sockets_cnt+uwsgi.no_orphans, -1); + ret = poll(uwsgi.sockets_poll, uwsgi.sockets_cnt+uwsgi.master_process, -1); if (ret < 0) { uwsgi_error("poll()"); @@ -471,9 +471,11 @@ polling: } if (uwsgi.master_process && uwsgi.sockets_poll[uwsgi.sockets_cnt].revents) { - if (read(uwsgi.sockets_poll[uwsgi.sockets_cnt].fd, &uwsgi_signal, 1) <= 0 && uwsgi.no_orphans) { - uwsgi_log_verbose("uWSGI worker %d screams: UAAAAAAH my master died, i will follow him...\n", uwsgi.mywid); - end_me(); + if (read(uwsgi.sockets_poll[uwsgi.sockets_cnt].fd, &uwsgi_signal, 1) <= 0) { + if (uwsgi.no_orphans) { + uwsgi_log_verbose("uWSGI worker %d screams: UAAAAAAH my master died, i will follow him...\n", uwsgi.mywid); + end_me(); + } } else { uwsgi_log_verbose("master sent signal %b to worker %d\n", uwsgi_signal, uwsgi.mywid); diff --git a/uwsgi.h b/uwsgi.h index 6556ad36..d659e482 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -601,10 +601,12 @@ struct uwsgi_fmon { }; struct uwsgi_timer { + char svalue[0xff]; int value; int fd; int id; int registered; + uint8_t sig; }; struct uwsgi_server { @@ -1309,15 +1311,15 @@ int event_queue_add_fd_read(int, int); int event_queue_wait(int, int, int *); int event_queue_add_timer(int, int *, int); -struct uwsgi_timer *event_queue_ack_timer(int, void (*)(int, int)); +struct uwsgi_timer *event_queue_ack_timer(int); int event_queue_add_file_monitor(int, char *, int *); -struct uwsgi_fmon *event_queue_ack_file_monitor(int, void (*)(char *, uint32_t, char *)); +struct uwsgi_fmon *event_queue_ack_file_monitor(int); void *uwsgi_mmap_shared_lock(void); void uwsgi_register_signal(uint8_t, uint8_t, void *, uint8_t, char *, uint8_t); -int uwsgi_signal_handler(uint8_t); - void uwsgi_register_file_monitor(uint8_t, char *, uint8_t, void *, uint8_t); +void uwsgi_register_timer(uint8_t, int, uint8_t, void *, uint8_t); +int uwsgi_signal_handler(uint8_t);