various async optimizations and routing python examples

This commit is contained in:
roberto@mrspurr
2010-12-04 13:43:12 +01:00
parent 47801964c4
commit 2f970ce585
16 changed files with 587 additions and 195 deletions
+412 -76
View File
@@ -20,15 +20,195 @@ extern struct uwsgi_python up;
#define UWSGI_LOGBASE "[- uWSGI -"
PyObject *py_uwsgi_send(PyObject * self, PyObject * args) {
char *uwsgi_encode_pydict(PyObject *pydict, uint16_t *size) {
char *data;
int i;
PyObject *zero, *key, *val;
uint16_t keysize, valsize;
if (!PyArg_ParseTuple(args, "s:send", &data)) {
char *buf, *bufptr;
PyObject *vars = PyDict_Items(pydict);
if (!vars) {
PyErr_Print();
return NULL;
}
*size = 0;
// calc the packet size
// try to fallback whenever possible
for (i = 0; i < PyList_Size(vars); i++) {
zero = PyList_GetItem(vars, i);
if (!zero) {
PyErr_Print();
continue;
}
if (!PyTuple_Check(zero)) {
uwsgi_log("invalid python dictionary item\n");
Py_DECREF(zero);
continue;
}
if (PyTuple_Size(zero) < 2) {
uwsgi_log("invalid python dictionary item\n");
Py_DECREF(zero);
continue;
}
key = PyTuple_GetItem(zero, 0);
val = PyTuple_GetItem(zero, 1);
if (!PyString_Check(key) || !PyString_Check(val)) {
Py_DECREF(zero);
continue;
}
keysize = PyString_Size(key);
valsize = PyString_Size(val);
*size += (keysize + 2 + valsize + 2);
// do not DECREF here !!!
//Py_DECREF(zero);
}
if (*size <= 4) {
uwsgi_log("empty python dictionary\n");
return NULL;
}
// remember to free this memory !!!
buf = malloc(*size);
if (!buf) {
uwsgi_error("malloc()");
return NULL;
}
if (write(uwsgi.wsgi_req->poll.fd, data, strlen(data)) < 0) {
bufptr = buf;
for (i = 0; i < PyList_Size(vars); i++) {
zero = PyList_GetItem(vars, i);
if (!zero) {
PyErr_Print();
continue;
}
if (!PyTuple_Check(zero)) {
uwsgi_log("invalid python dictionary item\n");
Py_DECREF(zero);
continue;
}
if (PyTuple_Size(zero) < 2) {
uwsgi_log("invalid python dictionary item\n");
Py_DECREF(zero);
continue;
}
key = PyTuple_GetItem(zero, 0);
val = PyTuple_GetItem(zero, 1);
if (!PyString_Check(key) || !PyString_Check(val)) {
Py_DECREF(zero);
continue;
}
keysize = PyString_Size(key);
valsize = PyString_Size(val);
if (bufptr + keysize + 2 + valsize + 2 <= buf + *size) {
#ifdef __BIG_ENDIAN__
keysize = uwsgi_swap16(keysize);
#endif
memcpy(bufptr, &keysize, 2);
bufptr += 2;
#ifdef __BIG_ENDIAN__
keysize = uwsgi_swap16(keysize);
#endif
memcpy(bufptr, PyString_AsString(key), keysize);
bufptr += keysize;
#ifdef __BIG_ENDIAN__
valsize = uwsgi_swap16(valsize);
#endif
memcpy(bufptr, &valsize, 2);
bufptr += 2;
#ifdef __BIG_ENDIAN__
valsize = uwsgi_swap16(valsize);
#endif
memcpy(bufptr, PyString_AsString(val), valsize);
bufptr += valsize;
}
Py_DECREF(zero);
}
return buf;
}
PyObject *py_uwsgi_close(PyObject * self, PyObject * args) {
int fd;
if (!PyArg_ParseTuple(args, "i:close", &fd)) {
return NULL;
}
close(fd);
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_recv(PyObject * self, PyObject * args) {
int fd, max_size = 4096;
char buf[4096];
ssize_t rlen ;
if (!PyArg_ParseTuple(args, "i|i:recv", &fd, &max_size)) {
return NULL;
}
// security check
if (max_size > 4096) max_size = 4096;
rlen = read(fd, buf, max_size) ;
if ( rlen > 0) {
return PyString_FromStringAndSize(buf, rlen);
}
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_send(PyObject * self, PyObject * args) {
PyObject *data;
PyObject *arg1, *arg2;
int uwsgi_fd = uwsgi.wsgi_req->poll.fd ;
if (!PyArg_ParseTuple(args, "O|O:send", &arg1, &arg2)) {
return NULL;
}
if (PyTuple_Size(args) > 1) {
uwsgi_fd = PyInt_AsLong(arg1);
data = arg2;
}
else {
data = arg1;
}
if (write(uwsgi_fd, PyString_AsString(data), PyString_Size(data)) < 0) {
uwsgi_error("write()");
Py_INCREF(Py_None);
return Py_None;
@@ -765,64 +945,205 @@ clear:
return Py_None;
}
PyObject *py_uwsgi_send_message(PyObject * self, PyObject * args) {
typedef struct {
PyObject_HEAD
int fd;
int timeout;
} uwsgi_Iter;
PyObject *arg_message = NULL;
const char *arg_host = NULL;
int arg_port = 0;
int arg_modifier1 = 0;
int arg_modifier2 = 0;
int arg_timeout = 0;
PyObject* uwsgi_Iter_iter(PyObject *self) {
Py_INCREF(self);
return self;
}
//PyObject *marshalled;
//PyObject *retobject;
PyObject* uwsgi_Iter_next(PyObject *self) {
int rlen;
uwsgi_Iter *ui = (uwsgi_Iter *)self;
char buf[4096];
if (!PyArg_ParseTuple(args, "siiiO|i:send_uwsgi_message", &arg_host, &arg_port, &arg_modifier1, &arg_modifier2, &arg_message, &arg_timeout)) {
return NULL;
rlen = uwsgi_waitfd(ui->fd, ui->timeout);
if (rlen > 0) {
rlen = read(ui->fd, buf, 4096);
if (rlen < 0) {
uwsgi_error("read()");
}
else if (rlen > 0) {
return PyString_FromStringAndSize(buf, rlen);
}
/*
switch (arg_modifier1) {
case UWSGI_MODIFIER_MESSAGE_MARSHAL:
marshalled = PyMarshal_WriteObjectToString(arg_message, 1);
if (!marshalled) {
PyErr_Print();
Py_INCREF(Py_None);
return Py_None;
}
retobject = uwsgi_send_message(arg_host, arg_port, arg_modifier1, arg_modifier2, PyString_AsString(marshalled), PyString_Size(marshalled), arg_timeout);
Py_DECREF(marshalled);
if (!retobject) {
PyErr_Print();
PyErr_Clear();
}
else {
return retobject;
}
break;
case UWSGI_MODIFIER_ADMIN_REQUEST:
if (PyString_Check(arg_message)) {
retobject = uwsgi_send_message(arg_host, arg_port, arg_modifier1, arg_modifier2, PyString_AsString(arg_message), PyString_Size(arg_message), arg_timeout);
if (!retobject) {
PyErr_Print();
PyErr_Clear();
}
else {
return retobject;
}
}
break;
default:
break;
}
*/
Py_INCREF(Py_None);
return Py_None;
}
else if (rlen == 0) {
uwsgi_log("uwsgi request timed out waiting for response\n");
}
PyErr_SetNone(PyExc_StopIteration);
return NULL;
}
static PyTypeObject uwsgi_IterType = {
PyObject_HEAD_INIT(NULL)
0, /*ob_size*/
"uwsgi._Iter", /*tp_name*/
sizeof(uwsgi_Iter), /*tp_basicsize*/
0, /*tp_itemsize*/
0, /*tp_dealloc*/
0, /*tp_print*/
0, /*tp_getattr*/
0, /*tp_setattr*/
0, /*tp_compare*/
0, /*tp_repr*/
0, /*tp_as_number*/
0, /*tp_as_sequence*/
0, /*tp_as_mapping*/
0, /*tp_hash */
0, /*tp_call*/
0, /*tp_str*/
0, /*tp_getattro*/
0, /*tp_setattro*/
0, /*tp_as_buffer*/
Py_TPFLAGS_DEFAULT | Py_TPFLAGS_HAVE_ITER,
"uwsgi response iterator object.", /* tp_doc */
0, /* tp_traverse */
0, /* tp_clear */
0, /* tp_richcompare */
0, /* tp_weaklistoffset */
uwsgi_Iter_iter, /* tp_iter: __iter__() method */
uwsgi_Iter_next /* tp_iternext: next() method */
};
PyObject *py_uwsgi_async_connect(PyObject * self, PyObject * args) {
char *socket_name = NULL;
if (!PyArg_ParseTuple(args, "s:async_connect", &socket_name)) {
return NULL;
}
return PyInt_FromLong(uwsgi_connect(socket_name, 0, 1));
}
PyObject *py_uwsgi_async_send_message(PyObject * self, PyObject * args) {
PyObject *pyobj = NULL, *marshalled = NULL;
int uwsgi_fd;
int modifier1 = 0;
int modifier2 = 0;
ssize_t ret ;
char *encoded;
uint16_t esize = 0;
if (!PyArg_ParseTuple(args, "iiiO:async_send_message", &uwsgi_fd, &modifier1, &modifier2, &pyobj)) {
return NULL;
}
if (uwsgi_fd < 0) goto clear;
// now check for the type of object to send (fallback to marshal)
if (PyDict_Check(pyobj)) {
encoded = uwsgi_encode_pydict(pyobj, &esize);
if (esize > 0) {
ret = uwsgi_send_message(uwsgi_fd, (uint8_t) modifier1, (uint8_t) modifier2, encoded, esize, -1, 0, 0);
free(encoded);
}
}
else if (PyString_Check(pyobj)) {
ret = uwsgi_send_message(uwsgi_fd, (uint8_t) modifier1, (uint8_t) modifier2, PyString_AsString(pyobj), PyString_Size(pyobj), -1, 0, 0);
}
else {
marshalled = PyMarshal_WriteObjectToString(pyobj, 1);
if (!marshalled) {
PyErr_Print();
goto clear;
}
ret = uwsgi_send_message(uwsgi_fd, (uint8_t) modifier1, (uint8_t) modifier2, PyString_AsString(marshalled), PyString_Size(marshalled), -1, 0, 0);
}
clear:
Py_INCREF(Py_None);
return Py_None;
}
PyObject *py_uwsgi_send_message(PyObject * self, PyObject * args) {
PyObject *destination = NULL, *pyobj = NULL, *marshalled = NULL;
int modifier1 = 0;
int modifier2 = 0;
int timeout = uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT];
int fd = -1;
int cl = 0;
ssize_t ret ;
int uwsgi_fd = -1;
char *encoded;
uint16_t esize = 0;
int close_fd = 0;
uwsgi_Iter *ui;
if (!PyArg_ParseTuple(args, "OiiO|iii:send_message", &destination, &modifier1, &modifier2, &pyobj, &timeout, &fd, &cl)) {
return NULL;
}
// first of all get the fd for the destination
if (PyInt_Check(destination)) {
uwsgi_fd = PyInt_AsLong(destination);
}
else if (PyString_Check(destination)) {
uwsgi_fd = uwsgi_connect(PyString_AsString(destination), timeout, 0);
close_fd = 1;
}
if (uwsgi_fd < 0) goto clear;
// now check for the type of object to send (fallback to marshal)
if (PyDict_Check(pyobj)) {
encoded = uwsgi_encode_pydict(pyobj, &esize);
if (esize > 0) {
ret = uwsgi_send_message(uwsgi_fd, (uint8_t) modifier1, (uint8_t) modifier2, encoded, esize, fd, cl, timeout);
free(encoded);
}
}
else if (PyString_Check(pyobj)) {
ret = uwsgi_send_message(uwsgi_fd, (uint8_t) modifier1, (uint8_t) modifier2, PyString_AsString(pyobj), PyString_Size(pyobj), fd, cl, timeout);
}
else {
marshalled = PyMarshal_WriteObjectToString(pyobj, 1);
if (!marshalled) {
PyErr_Print();
goto clear2;
}
ret = uwsgi_send_message(uwsgi_fd, (uint8_t) modifier1, (uint8_t) modifier2, PyString_AsString(marshalled), PyString_Size(marshalled), fd, cl, timeout);
}
// request sent, return the iterator response
ui = PyObject_New(uwsgi_Iter, &uwsgi_IterType);
if (!ui) {
PyErr_Print();
goto clear2;
}
ui->fd = uwsgi_fd;
ui->timeout = timeout;
return (PyObject *) ui;
clear2:
if (close_fd) close(uwsgi_fd);
clear:
Py_INCREF(Py_None);
return Py_None;
}
/* uWSGI masterpid */
PyObject *py_uwsgi_masterpid(PyObject * self, PyObject * args) {
@@ -1163,9 +1484,9 @@ PyObject *py_uwsgi_suspend(PyObject * self, PyObject * args) {
}
static PyMethodDef uwsgi_advanced_methods[] = {
{"send_uwsgi_message", py_uwsgi_send_message, METH_VARARGS, ""},
{"send_multi_uwsgi_message", py_uwsgi_send_multi_message, METH_VARARGS, ""},
static PyMethodDef uwsgi_advanced_methods[] = {
{"send_message", py_uwsgi_send_message, METH_VARARGS, ""},
{"send_multi_message", py_uwsgi_send_multi_message, METH_VARARGS, ""},
{"reload", py_uwsgi_reload, METH_VARARGS, ""},
{"workers", py_uwsgi_workers, METH_VARARGS, ""},
{"masterpid", py_uwsgi_masterpid, METH_VARARGS, ""},
@@ -1196,27 +1517,35 @@ PyObject *py_uwsgi_suspend(PyObject * self, PyObject * args) {
#endif
#ifdef UWSGI_ASYNC
{"async_sleep", py_uwsgi_async_sleep, METH_VARARGS, ""},
{"async_connect", py_uwsgi_async_connect, METH_VARARGS, ""},
{"async_send_message", py_uwsgi_async_send_message, METH_VARARGS, ""},
{"green_schedule", py_uwsgi_suspend, METH_VARARGS, ""},
{"suspend", py_uwsgi_suspend, METH_VARARGS, ""},
{"wait_fd_read", py_eventfd_read, METH_VARARGS, ""},
{"wait_fd_write", py_eventfd_write, METH_VARARGS, ""},
#endif
{"send", py_uwsgi_send, METH_VARARGS, ""},
{"recv", py_uwsgi_recv, METH_VARARGS, ""},
{"close", py_uwsgi_close, METH_VARARGS, ""},
{"parsefile", py_uwsgi_parse_file, METH_VARARGS, ""},
//{"call_hook", py_uwsgi_call_hook, METH_VARARGS, ""},
{NULL, NULL},
};
};
static PyMethodDef uwsgi_sa_methods[] = {
{"sharedarea_read", py_uwsgi_sharedarea_read, METH_VARARGS, ""},
{"sharedarea_write", py_uwsgi_sharedarea_write, METH_VARARGS, ""},
{"sharedarea_readbyte", py_uwsgi_sharedarea_readbyte, METH_VARARGS, ""},
{"sharedarea_writebyte", py_uwsgi_sharedarea_writebyte, METH_VARARGS, ""},
{"sharedarea_readlong", py_uwsgi_sharedarea_readlong, METH_VARARGS, ""},
{"sharedarea_writelong", py_uwsgi_sharedarea_writelong, METH_VARARGS, ""},
{"sharedarea_inclong", py_uwsgi_sharedarea_inclong, METH_VARARGS, ""},
{NULL, NULL},
};
static PyMethodDef uwsgi_sa_methods[] = {
{"sharedarea_read", py_uwsgi_sharedarea_read, METH_VARARGS, ""},
{"sharedarea_write", py_uwsgi_sharedarea_write, METH_VARARGS, ""},
{"sharedarea_readbyte", py_uwsgi_sharedarea_readbyte, METH_VARARGS, ""},
{"sharedarea_writebyte", py_uwsgi_sharedarea_writebyte, METH_VARARGS, ""},
{"sharedarea_readlong", py_uwsgi_sharedarea_readlong, METH_VARARGS, ""},
{"sharedarea_writelong", py_uwsgi_sharedarea_writelong, METH_VARARGS, ""},
{"sharedarea_inclong", py_uwsgi_sharedarea_inclong, METH_VARARGS, ""},
{NULL, NULL},
};
@@ -1246,7 +1575,7 @@ PyObject *py_uwsgi_suspend(PyObject * self, PyObject * args) {
}
#endif
void init_uwsgi_module_advanced(PyObject * current_uwsgi_module) {
void init_uwsgi_module_advanced(PyObject * current_uwsgi_module) {
PyMethodDef *uwsgi_function;
PyObject *uwsgi_module_dict;
@@ -1256,13 +1585,20 @@ PyObject *py_uwsgi_suspend(PyObject * self, PyObject * args) {
exit(1);
}
for (uwsgi_function = uwsgi_advanced_methods; uwsgi_function->ml_name != NULL; uwsgi_function++) {
PyObject *func = PyCFunction_New(uwsgi_function, NULL);
PyDict_SetItemString(uwsgi_module_dict, uwsgi_function->ml_name, func);
Py_DECREF(func);
}
uwsgi_IterType.tp_new = PyType_GenericNew;
if (PyType_Ready(&uwsgi_IterType) < 0) {
PyErr_Print();
exit(1);
}
for (uwsgi_function = uwsgi_advanced_methods; uwsgi_function->ml_name != NULL; uwsgi_function++) {
PyObject *func = PyCFunction_New(uwsgi_function, NULL);
PyDict_SetItemString(uwsgi_module_dict, uwsgi_function->ml_name, func);
Py_DECREF(func);
}
}
void init_uwsgi_module_sharedarea(PyObject * current_uwsgi_module) {
PyMethodDef *uwsgi_function;
PyObject *uwsgi_module_dict;