diff --git a/plugins/lua/lua_plugin.c b/plugins/lua/lua_plugin.c index 20da1485..be08fcdf 100644 --- a/plugins/lua/lua_plugin.c +++ b/plugins/lua/lua_plugin.c @@ -165,25 +165,18 @@ static int uwsgi_api_cache_set(lua_State *L) { static int uwsgi_api_register_signal(lua_State *L) { int args = lua_gettop(L); - uint8_t sig, kind; + uint8_t sig; long lhandler; - const char *payload; - size_t payload_size; + const char *who; if (args >= 3) { sig = lua_tonumber(L, 1); - kind = lua_tonumber(L, 2); + who = lua_tostring(L, 2); lua_pushvalue(L, 3); lhandler = luaL_ref(L, LUA_REGISTRYINDEX); - if (args > 3) { - payload = lua_tolstring(L, 4, &payload_size); - uwsgi_register_signal(sig, kind, (void *) lhandler, 6, (char *) payload, payload_size); - } - else { - uwsgi_register_signal(sig, kind, (void *) lhandler, 6, NULL, 0); - } + uwsgi_register_signal(sig, kind, (void *) lhandler, 6); } lua_pushnil(L); @@ -597,7 +590,7 @@ int uwsgi_lua_magic(char *mountpoint, char *lazy) { return 0; } -int uwsgi_lua_signal_handler(uint8_t sig, void *handler, char *payload, uint8_t payload_size) { +int uwsgi_lua_signal_handler(uint8_t sig, void *handler) { struct wsgi_request *wsgi_req = current_wsgi_req(); diff --git a/plugins/python/python_plugin.c b/plugins/python/python_plugin.c index 8e4fddae..030fabbb 100644 --- a/plugins/python/python_plugin.c +++ b/plugins/python/python_plugin.c @@ -818,16 +818,15 @@ void uwsgi_python_suspend(struct wsgi_request *wsgi_req) { } -int uwsgi_python_signal_handler(uint8_t sig, void *handler, char *payload, uint8_t payload_size) { +int uwsgi_python_signal_handler(uint8_t sig, void *handler) { - PyObject *args = PyTuple_New(2); + PyObject *args = PyTuple_New(1); PyObject *ret; if (!args) return -1; PyTuple_SetItem(args, 0, PyInt_FromLong(sig)); - PyTuple_SetItem(args, 1, PyString_FromStringAndSize(payload, payload_size)); ret = python_call(handler, args, 0); diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index 7cc9ef86..c17c0f61 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -181,36 +181,34 @@ PyObject *py_uwsgi_close(PyObject * self, PyObject * args) { } -PyObject *py_uwsgi_register_timer(PyObject * self, PyObject * args) { +PyObject *py_uwsgi_add_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)) { + if (!PyArg_ParseTuple(args, "Bi:add_timer", &uwsgi_signal, &secs)) { return NULL; } - uwsgi_register_timer(uwsgi_signal, secs, signal_kind, handler, 0); + if (uwsgi_add_timer(uwsgi_signal, secs)) + return PyErr_Format(PyExc_ValueError, "unable to add timer"); Py_INCREF(Py_None); return Py_None; } -PyObject *py_uwsgi_register_file_monitor(PyObject * self, PyObject * args) { +PyObject *py_uwsgi_add_file_monitor(PyObject * self, PyObject * args) { uint8_t uwsgi_signal; - uint8_t signal_kind; - PyObject *handler; char *filename; - if (!PyArg_ParseTuple(args, "BsBO:register_file_monitor", &uwsgi_signal, &filename, &signal_kind, &handler)) { + if (!PyArg_ParseTuple(args, "Bs:add_file_monitor", &uwsgi_signal, &filename)) { return NULL; } - uwsgi_register_file_monitor(uwsgi_signal, filename, signal_kind, handler, 0); + if (uwsgi_add_file_monitor(uwsgi_signal, filename)) + return PyErr_Format(PyExc_ValueError, "unable to add file monitor"); Py_INCREF(Py_None); return Py_None; @@ -419,19 +417,15 @@ PyObject *py_uwsgi_attach_daemon(PyObject * self, PyObject * args) { PyObject *py_uwsgi_register_signal(PyObject * self, PyObject * args) { uint8_t uwsgi_signal; - uint8_t signal_kind; + char *signal_kind; PyObject *handler; - char *payload = NULL; - if (!PyArg_ParseTuple(args, "BBO|s:register_signal", &uwsgi_signal, &signal_kind, &handler, &payload)) { + if (!PyArg_ParseTuple(args, "BsO:register_signal", &uwsgi_signal, &signal_kind, &handler)) { return NULL; } - if (payload == NULL) { - //uwsgi_register_signal(uwsgi_signal, signal_kind, handler, 0, NULL, 0); - } - else { - //uwsgi_register_signal(uwsgi_signal, signal_kind, handler, 0, payload, strlen(payload)); + if (uwsgi_register_signal(uwsgi_signal, signal_kind, handler, 0)) { + return PyErr_Format(PyExc_ValueError, "unable to register signal"); } Py_INCREF(Py_None); @@ -2327,8 +2321,9 @@ static PyMethodDef uwsgi_advanced_methods[] = { {"signal", py_uwsgi_signal, METH_VARARGS, ""}, {"signal_wait", py_uwsgi_signal_wait, METH_VARARGS, ""}, {"signal_received", py_uwsgi_signal_received, METH_VARARGS, ""}, - {"register_file_monitor", py_uwsgi_register_file_monitor, METH_VARARGS, ""}, - {"register_timer", py_uwsgi_register_timer, METH_VARARGS, ""}, + {"add_file_monitor", py_uwsgi_add_file_monitor, METH_VARARGS, ""}, + {"add_timer", py_uwsgi_add_timer, METH_VARARGS, ""}, + //{"add_rb_timer", py_uwsgi_add_rb_timer, METH_VARARGS, ""}, {"register_rpc", py_uwsgi_register_rpc, METH_VARARGS, ""}, {"rpc", py_uwsgi_rpc, METH_VARARGS, ""}, diff --git a/signal.c b/signal.c index 3be8c279..1b865044 100644 --- a/signal.c +++ b/signal.c @@ -8,23 +8,19 @@ int uwsgi_signal_handler(uint8_t sig) { use = &uwsgi.shared->signal_table[sig]; - /* - if (!use->kind) { - return -1; - } - */ - if (!uwsgi.p[use->modifier1]->signal_handler) { return -1; } - return uwsgi.p[use->modifier1]->signal_handler(sig, use->handler, use->payload, use->payload_size); + return uwsgi.p[use->modifier1]->signal_handler(sig, use->handler); } -void uwsgi_register_signal(uint8_t sig, char *receiver, void *handler, uint8_t modifier1, char *payload, uint8_t payload_size) { +int uwsgi_register_signal(uint8_t sig, char *receiver, void *handler, uint8_t modifier1) { struct uwsgi_signal_entry *use = NULL; + if (strlen(receiver) > 63) return -1; + uwsgi_lock(uwsgi.signal_table_lock); use = &uwsgi.shared->signal_table[sig]; @@ -33,22 +29,19 @@ void uwsgi_register_signal(uint8_t sig, char *receiver, void *handler, uint8_t m use->handler = handler; use->modifier1 = modifier1; - memcpy(use->payload, payload, payload_size); - - use->payload_size = payload_size; - uwsgi_log("registered signal %d\n", sig); uwsgi_unlock(uwsgi.signal_table_lock); + return 0; } -void uwsgi_register_file_monitor(uint8_t sig, char *filename, uint8_t kind, void *handler, uint8_t modifier1) { +int uwsgi_add_file_monitor(uint8_t sig, char *filename) { if (strlen(filename) > (0xff-1)) { - uwsgi_log("uwsgi_register_file_monitor: invalid filename length\n"); - return; + uwsgi_log("uwsgi_add_file_monitor: invalid filename length\n"); + return -1; } uwsgi_lock(uwsgi.fmon_table_lock); @@ -60,19 +53,21 @@ void uwsgi_register_file_monitor(uint8_t sig, char *filename, uint8_t kind, void ushared->files_monitored[ushared->files_monitored_cnt].registered = 0; ushared->files_monitored[ushared->files_monitored_cnt].sig = sig; - //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"); + uwsgi_unlock(uwsgi.fmon_table_lock); + return -1; } uwsgi_unlock(uwsgi.fmon_table_lock); + return 0; + } -void uwsgi_register_timer(uint8_t sig, int secs, uint8_t kind, void *handler, uint8_t modifier1) { +int uwsgi_add_timer(uint8_t sig, int secs) { uwsgi_lock(uwsgi.timer_table_lock); @@ -80,36 +75,39 @@ void uwsgi_register_timer(uint8_t sig, int secs, uint8_t kind, void *handler, ui // 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); + return -1; } uwsgi_unlock(uwsgi.timer_table_lock); + return 0; + } void uwsgi_route_signal(uint8_t sig) { struct uwsgi_signal_entry *use = &ushared->signal_table[sig]; - // send to all workers - if (!strcmp(use->receiver, "workers")) { + // send to first available worker + if (use->receiver[0] == 0 || !strcmp(use->receiver, "worker") || !strcmp(use->receiver, "worker0")) { if (write(ushared->worker_signal_pipe[0], &sig, 1) != 1) { uwsgi_error("write()"); uwsgi_log("could not deliver signal %d to workers pool\n", sig); } } - // send to first available worker - else if (!strcmp(use->receiver, "worker")) { - } - // loop back to master - else if (!strcmp(use->receiver, "master")) { + // send to all workers + else if (!strcmp(use->receiver, "workers")) { + if (write(ushared->worker_signal_pipe[0], &sig, 1) != 1) { + uwsgi_error("write()"); + uwsgi_log("could not deliver signal %d to workers pool\n", sig); + } } // route to subscribed else if (!strcmp(use->receiver, "subscribed")) { diff --git a/uwsgi.h b/uwsgi.h index 68845cbe..17aea92b 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -470,7 +470,7 @@ struct uwsgi_plugin { void* (*encode_string)(char *); char* (*decode_string)(void *); - int (*signal_handler)(uint8_t, void *, char *, uint8_t); + int (*signal_handler)(uint8_t, void *); int (*spooler)(char *, uint16_t); @@ -989,10 +989,8 @@ struct uwsgi_lb_group { struct uwsgi_signal_entry { uint8_t modifier1; - uint8_t payload_size; - char receiver[0xff]; + char receiver[64]; void *handler; - char payload[0xff]; }; struct uwsgi_lb_node { @@ -1410,9 +1408,9 @@ struct uwsgi_fmon *event_queue_ack_file_monitor(int, int); void *uwsgi_mmap_shared_lock(void); void *uwsgi_mmap_shared_rwlock(void); -void uwsgi_register_signal(uint8_t, char *, void *, uint8_t, char *, 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_register_signal(uint8_t, char *, void *, uint8_t); +int uwsgi_add_file_monitor(uint8_t, char *); +int uwsgi_add_timer(uint8_t, int); int uwsgi_signal_handler(uint8_t); void uwsgi_route_signal(uint8_t);