improved signal handling and rack mule_msg

This commit is contained in:
roberto@debian32
2011-11-13 15:08:29 +01:00
parent c86b757904
commit e195a5d72c
8 changed files with 359 additions and 206 deletions
+20 -111
View File
@@ -518,7 +518,6 @@ PyObject *py_uwsgi_register_signal(PyObject * self, PyObject * args) {
PyObject *py_uwsgi_signal(PyObject * self, PyObject * args) {
uint8_t uwsgi_signal;
ssize_t rlen;
if (!PyArg_ParseTuple(args, "B:signal", &uwsgi_signal)) {
return NULL;
@@ -528,11 +527,7 @@ PyObject *py_uwsgi_signal(PyObject * self, PyObject * args) {
uwsgi_log("sending %d to master\n", uwsgi_signal);
#endif
rlen = write(uwsgi.signal_socket, &uwsgi_signal, 1);
if (rlen != 1) {
uwsgi_error("write()");
}
uwsgi_signal_send(uwsgi.signal_socket, uwsgi_signal);
Py_INCREF(Py_None);
return Py_None;
@@ -1110,7 +1105,6 @@ PyObject *py_uwsgi_mule_msg(PyObject * self, PyObject * args) {
char *message = NULL;
Py_ssize_t message_len = 0;
PyObject *mule_obj = NULL;
ssize_t len;
int fd = -1;
int mule_id = -1;
@@ -1122,10 +1116,7 @@ PyObject *py_uwsgi_mule_msg(PyObject * self, PyObject * args) {
return PyErr_Format(PyExc_ValueError, "no mule configured");
if (mule_obj == NULL) {
len = write(uwsgi.shared->mule_queue_pipe[0], message, message_len);
if (len <= 0) {
uwsgi_error("write()");
}
mule_send_msg(uwsgi.shared->mule_queue_pipe[0], message, message_len);
}
else {
if (PyString_Check(mule_obj)) {
@@ -1152,10 +1143,7 @@ PyObject *py_uwsgi_mule_msg(PyObject * self, PyObject * args) {
}
if (fd > -1) {
len = write(fd, message, message_len);
if (len < 0) {
uwsgi_error("write()");
}
mule_send_msg(fd, message, message_len);
}
}
@@ -1167,17 +1155,13 @@ PyObject *py_uwsgi_mule_msg(PyObject * self, PyObject * args) {
PyObject *py_uwsgi_mule_get_msg(PyObject * self, PyObject * args, PyObject *kwargs) {
ssize_t len = 0;
// this buffer will be configurable
// this buffer is configurable (default 64k)
char *message;
struct pollfd *mulepoll;
int count = 4;
int farms_count = 0;
uint8_t uwsgi_signal;
PyObject *manage_signals = NULL;
PyObject *manage_farms = NULL;
int buffer_size = 65536;
PyObject *py_manage_signals = NULL;
PyObject *py_manage_farms = NULL;
size_t buffer_size = 65536;
int timeout = -1;
int i;
int manage_signals = 1, manage_farms = 1;
static char *kwlist[] = {"signals", "buffer_size", "timeout", "farms", NULL};
@@ -1185,109 +1169,34 @@ PyObject *py_uwsgi_mule_get_msg(PyObject * self, PyObject * args, PyObject *kwar
return PyErr_Format(PyExc_ValueError, "you can receive mule messages only in a mule !!!");
}
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|OOii:mule_get_msg", kwlist, &manage_signals, &manage_farms, &buffer_size, &timeout)) {
if (!PyArg_ParseTupleAndKeywords(args, kwargs, "|OOii:mule_get_msg", kwlist, &py_manage_signals, &py_manage_farms, &buffer_size, &timeout)) {
return NULL;
}
if (manage_signals == Py_None || manage_signals == Py_False) {
count = 2;
// signals and farms are managed by default
if (py_manage_signals == Py_None || py_manage_signals == Py_False) {
manage_signals = 0;
}
if (manage_farms == Py_None || manage_farms == Py_False) {
goto next;
if (py_manage_farms == Py_None || py_manage_farms == Py_False) {
manage_farms = 0;
}
for(i=0;i<uwsgi.farms_cnt;i++) {
if (uwsgi_farm_has_mule(&uwsgi.farms[i], uwsgi.muleid)) farms_count++;
}
next:
UWSGI_RELEASE_GIL;
if (timeout > -1) timeout = timeout*1000;
message = uwsgi_malloc(buffer_size);
mulepoll = uwsgi_malloc(sizeof(struct pollfd) * (count+farms_count));
mulepoll[0].fd = uwsgi.mules[uwsgi.muleid-1].queue_pipe[1];
mulepoll[0].events = POLLIN;
mulepoll[1].fd = uwsgi.shared->mule_queue_pipe[1];
mulepoll[1].events = POLLIN;
if (count > 2) {
mulepoll[2].fd = uwsgi.signal_socket;
mulepoll[2].events = POLLIN;
mulepoll[3].fd = uwsgi.my_signal_socket;
mulepoll[3].events = POLLIN;
}
for(i=0;i<farms_count;i++) {
mulepoll[count+i].fd = uwsgi.farms[i].queue_pipe[1];
mulepoll[count+i].events = POLLIN;
}
int ret = poll(mulepoll, count+farms_count, timeout);
if (ret <= 0) {
uwsgi_error("poll");
}
else {
if (mulepoll[0].revents & POLLIN) {
len = read(uwsgi.mules[uwsgi.muleid-1].queue_pipe[1], message, buffer_size);
}
else if (mulepoll[1].revents & POLLIN) {
len = read(uwsgi.shared->mule_queue_pipe[1], message, buffer_size);
}
else {
if (count > 2) {
int interesting_fd = -1;
if (mulepoll[2].revents & POLLIN) {
interesting_fd = mulepoll[2].fd;
}
else if (mulepoll[3].revents & POLLIN) {
interesting_fd = mulepoll[3].fd;
}
if (interesting_fd > -1) {
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);
#endif
if (uwsgi_signal_handler(uwsgi_signal)) {
uwsgi_log_verbose("error managing signal %d on mule %d\n", uwsgi_signal, uwsgi.mywid);
}
goto clear;
}
}
for(i=0;i<farms_count;i++) {
if (mulepoll[count+i].revents & POLLIN) {
len = read(mulepoll[count+i].fd, message, buffer_size);
break;
}
}
}
}
UWSGI_RELEASE_GIL;
len = uwsgi_mule_get_msg(manage_signals, manage_farms, message, buffer_size, timeout) ;
UWSGI_GET_GIL;
if (len < 0) {
uwsgi_error("read()");
goto clear2;
free(message);
Py_INCREF(Py_None);
return Py_None;
}
PyObject *msg = PyString_FromStringAndSize(message, len);
free(message);
free(mulepoll);
return msg;
clear:
UWSGI_GET_GIL;
clear2:
free(message);
free(mulepoll);
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_farm_get_msg(PyObject * self, PyObject * args) {