updated signal framework to the new api

This commit is contained in:
roberto@sirius
2011-02-24 09:48:28 +01:00
parent b752425223
commit 7133069ee9
5 changed files with 52 additions and 69 deletions
+5 -12
View File
@@ -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();
+2 -3
View File
@@ -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);
+15 -20
View File
@@ -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, ""},
+25 -27
View File
@@ -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")) {
+5 -7
View File
@@ -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);