From dce3263ee35863343c9ec2d99f26e91779013f7a Mon Sep 17 00:00:00 2001 From: "roberto@debian32" Date: Mon, 3 Oct 2011 07:45:18 +0200 Subject: [PATCH] added mule messaging subsystem --- master.c | 2 ++ mule.c | 63 +++++++++++++++++++++++++++------ plugins/python/python_plugin.c | 25 +++++++------ plugins/python/uwsgi_pymodule.c | 49 +++++++++++++++++++++++++ utils.c | 19 ++++++++++ uwsgi.c | 31 ++++++++++++++-- uwsgi.h | 11 ++++-- welcome.py | 2 ++ 8 files changed, 177 insertions(+), 25 deletions(-) diff --git a/master.c b/master.c index dd373148..440d7ca5 100644 --- a/master.c +++ b/master.c @@ -482,6 +482,8 @@ int master_loop(char **argv, char **environ) { // spawn mules for(i=0;imule_signal_pipe[1]; uwsgi_close_all_sockets(); @@ -40,6 +42,17 @@ void uwsgi_mule(int id) { } + if (uwsgi.mules[id-1].patch) { + uwsgi_log("loading patch %s\n", uwsgi.mules[id-1].patch); + for (i = 0; i < 0xFF; i++) { + if (uwsgi.p[i]->mule) { + if (uwsgi.p[i]->mule(uwsgi.mules[id-1].patch) == 1) { + break; + } + } + } + } + uwsgi_mule_handler(); } else if (pid > 0) { @@ -54,19 +67,49 @@ void uwsgi_mule_handler() { ssize_t len; uint8_t uwsgi_signal; + int rlen; + int interesting_fd; + + // this must be configurable + char message[65536]; + + int mule_queue = event_queue_init(); + + event_queue_add_fd_read(mule_queue, uwsgi.signal_socket); + event_queue_add_fd_read(mule_queue, uwsgi.my_signal_socket); + event_queue_add_fd_read(mule_queue, uwsgi.mules[uwsgi.muleid-1].queue_pipe[1]); for(;;) { - len = read(uwsgi.signal_socket, &uwsgi_signal, 1); - if (len <= 0) { - uwsgi_log_verbose("uWSGI mule %d braying: my master died, i will follow him...\n", uwsgi.muleid); - end_me(0); - } + rlen = event_queue_wait(mule_queue, -1, &interesting_fd); + if (rlen <= 0) { + continue; + } + + if (interesting_fd == uwsgi.signal_socket || interesting_fd == uwsgi.my_signal_socket) { + len = read(interesting_fd, &uwsgi_signal, 1); + if (len <= 0) { + uwsgi_log_verbose("uWSGI mule %d braying: my master died, i will follow him...\n", uwsgi.muleid); + end_me(0); + } #ifdef UWSGI_DEBUG - uwsgi_log_verbose("master sent signal %d to mule %d\n", uwsgi_signal, uwsgi.muleid); + uwsgi_log_verbose("master sent signal %d to mule %d\n", uwsgi_signal, uwsgi.muleid); #endif - if (uwsgi_signal_handler(uwsgi_signal)) { - uwsgi_log_verbose("error managing signal %d on mule %d\n", uwsgi_signal, uwsgi.mywid); - } + if (uwsgi_signal_handler(uwsgi_signal)) { + uwsgi_log_verbose("error managing signal %d on mule %d\n", uwsgi_signal, uwsgi.mywid); + } + } + else if (interesting_fd == uwsgi.mules[uwsgi.muleid-1].queue_pipe[1]) { + len = read(uwsgi.mules[uwsgi.muleid-1].queue_pipe[1], message, 65536); + if (len < 0) { + uwsgi_error("read()"); + } + else if (len == 0) { + exit(1); + } + else { + uwsgi_log("*** mule %d received a %d bytes message\n", uwsgi.muleid, len); + } + } } } diff --git a/plugins/python/python_plugin.c b/plugins/python/python_plugin.c index 218e0278..6d92c627 100644 --- a/plugins/python/python_plugin.c +++ b/plugins/python/python_plugin.c @@ -525,11 +525,8 @@ void init_uwsgi_embedded_module() { } } - uwsgi_log("getting exported opts\n"); - PyObject *py_opt_dict = PyDict_New(); for (i = 0; i < uwsgi.exported_opts_cnt; i++) { - uwsgi_log("%s = %s\n", uwsgi.exported_opts[i]->key, uwsgi.exported_opts[i]->value); if (PyDict_Contains(py_opt_dict, PyString_FromString(uwsgi.exported_opts[i]->key))) { PyObject *py_opt_item = PyDict_GetItemString(py_opt_dict, uwsgi.exported_opts[i]->key); if (PyList_Check(py_opt_item)) { @@ -563,8 +560,6 @@ void init_uwsgi_embedded_module() { } } - uwsgi_log("DONE\n"); - if (PyDict_SetItemString(up.embedded_dict, "opt", py_opt_dict)) { PyErr_Print(); exit(1); @@ -910,9 +905,7 @@ void uwsgi_python_init_apps() { #endif } - uwsgi_log("init_pyargv\n"); init_pyargv(); - uwsgi_log("init_pyargv\n"); #ifndef UWSGI_PYPY #ifdef UWSGI_EMBEDDED @@ -920,7 +913,6 @@ void uwsgi_python_init_apps() { #endif #endif - uwsgi_log("init_pyargv\n"); #ifdef __linux__ #if !defined(PYTHREE) && !defined(UWSGI_PYPY) @@ -935,8 +927,6 @@ void uwsgi_python_init_apps() { exit(1); } - uwsgi_log("init_pyargv\n"); - init_uwsgi_vars(); // setup app loaders @@ -1404,6 +1394,19 @@ void uwsgi_python_hijack(void) { #endif } +int uwsgi_python_mule(char *opt) { + + if (uwsgi_endswith(opt, ".py")) { + UWSGI_GET_GIL; + uwsgi_pyimport_by_filename("__main__", opt); + UWSGI_RELEASE_GIL; + return 1; + } + + return 0; + +} + struct uwsgi_plugin python_plugin = { .name = "python", @@ -1440,6 +1443,8 @@ struct uwsgi_plugin python_plugin = { .signal_handler = uwsgi_python_signal_handler, .rpc = uwsgi_python_rpc, + .mule = uwsgi_python_mule, + .spooler = uwsgi_python_spooler, .code_string = uwsgi_python_code_string, diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index 45fcea46..5d6854be 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -1050,6 +1050,52 @@ PyObject *py_uwsgi_embedded_data(PyObject * self, PyObject * args) { } +PyObject *py_uwsgi_mule_msg(PyObject * self, PyObject * args) { + + char *message = NULL; + Py_ssize_t message_len = 0; + int mule_id = 0; + ssize_t len; + + if (!PyArg_ParseTuple(args, "s#|i:mule_msg", &message, &message_len, &mule_id)) { + return NULL; + } + + if (mule_id == 0) { + } + else if (mule_id > 0 && mule_id <= uwsgi.mules_cnt) { + len = write(uwsgi.mules[mule_id-1].queue_pipe[0], message, message_len); + if (len <= 0) { + uwsgi_error("write()"); + } + } + + Py_INCREF(Py_None); + return Py_None; + +} + +PyObject *py_uwsgi_mule_get_msg(PyObject * self, PyObject * args) { + + ssize_t len; + // this buffer will be configurable + char message[65536]; + + if (uwsgi.muleid == 0) { + return PyErr_Format(PyExc_ValueError, "you can receive mule messages only in a mule !!!"); + } + UWSGI_RELEASE_GIL; + len = read(uwsgi.mules[uwsgi.muleid-1].queue_pipe[1], message, 65536); + UWSGI_GET_GIL; + if (len <= 0) { + uwsgi_error("read()"); + Py_INCREF(Py_None); + return Py_None; + } + + return PyString_FromStringAndSize(message, len); +} + PyObject *py_uwsgi_extract(PyObject * self, PyObject * args) { char *name; @@ -2784,6 +2830,9 @@ static PyMethodDef uwsgi_advanced_methods[] = { {"parsefile", py_uwsgi_parse_file, METH_VARARGS, ""}, {"embedded_data", py_uwsgi_embedded_data, METH_VARARGS, ""}, {"extract", py_uwsgi_extract, METH_VARARGS, ""}, + + {"mule_msg", py_uwsgi_mule_msg, METH_VARARGS, ""}, + {"mule_get_msg", py_uwsgi_mule_get_msg, METH_VARARGS, ""}, //{"call_hook", py_uwsgi_call_hook, METH_VARARGS, ""}, {NULL, NULL}, diff --git a/utils.c b/utils.c index ca0830f6..872b1bfa 100644 --- a/utils.c +++ b/utils.c @@ -2613,6 +2613,25 @@ struct uwsgi_string_list *uwsgi_string_new_list(struct uwsgi_string_list **list, return uwsgi_string; } +char *uwsgi_string_get_list(struct uwsgi_string_list **list, int pos, size_t *len) { + + struct uwsgi_string_list *uwsgi_string = *list; + int counter = 0; + + while(uwsgi_string) { + if (counter == pos) { + *len = uwsgi_string->len; + return uwsgi_string->value; + } + uwsgi_string = uwsgi_string->next; + counter++; + } + + *len = 0; + return NULL; + +} + void uwsgi_string_del_list(struct uwsgi_string_list **list, struct uwsgi_string_list *item) { diff --git a/uwsgi.c b/uwsgi.c index 97d89a02..f21bda23 100644 --- a/uwsgi.c +++ b/uwsgi.c @@ -520,6 +520,11 @@ void kill_them_all(int signum) { kill(uwsgi.gateways[i].pid, SIGKILL); } + for (i = 0; i < uwsgi.mules_cnt; i++) { + if (uwsgi.mules[i].pid > 0) + kill(uwsgi.mules[i].pid, SIGKILL); + } + } void grace_them_all(int signum) { @@ -563,6 +568,11 @@ void grace_them_all(int signum) { kill(uwsgi.gateways[i].pid, SIGKILL); } + for (i = 0; i < uwsgi.mules_cnt; i++) { + if (uwsgi.mules[i].pid > 0) + kill(uwsgi.mules[i].pid, SIGKILL); + } + uwsgi_log("...gracefully killing workers...\n"); for (i = 1; i <= uwsgi.numproc; i++) { @@ -630,6 +640,11 @@ void reap_them_all(int signum) { kill(uwsgi.gateways[i].pid, SIGKILL); } + for (i = 0; i < uwsgi.mules_cnt; i++) { + if (uwsgi.mules[i].pid > 0) + kill(uwsgi.mules[i].pid, SIGKILL); + } + if (uwsgi.emperor_pid >= 0) { kill(uwsgi.emperor_pid, SIGKILL); waitpid(uwsgi.emperor_pid, &i, 0); @@ -928,6 +943,8 @@ int main(int argc, char *argv[], char *envp[]) { uwsgi.shared->spooler_signal_pipe[0] = -1; uwsgi.shared->spooler_signal_pipe[1] = -1; #endif + uwsgi.shared->mule_signal_pipe[0] = -1; + uwsgi.shared->mule_signal_pipe[1] = -1; uwsgi.mime_file = "/etc/mime.types"; @@ -2132,6 +2149,18 @@ skipzero: exit(1); } memset(uwsgi.mules, 0, sizeof(struct uwsgi_mule) * uwsgi.mules_cnt); + + if (socketpair(AF_UNIX, SOCK_STREAM, 0, uwsgi.shared->mule_signal_pipe)) { + uwsgi_error("socketpair()"); + exit(1); + } + + for(i=0;i