diff --git a/plugins/cache/cache.c b/plugins/cache/cache.c index 687891aa..22ec375d 100644 --- a/plugins/cache/cache.c +++ b/plugins/cache/cache.c @@ -34,6 +34,7 @@ int uwsgi_cache_request(struct wsgi_request *wsgi_req) { uint64_t vallen = 0; char *value; char *argv[3]; + uint16_t argvs[3]; uint8_t argc = 0; switch(wsgi_req->uh.modifier2) { @@ -52,9 +53,9 @@ int uwsgi_cache_request(struct wsgi_request *wsgi_req) { // set if (wsgi_req->uh.pktsize > 0) { argc = 3; - if (!uwsgi_parse_array(wsgi_req->buffer, wsgi_req->uh.pktsize, argv, &argc)) { + if (!uwsgi_parse_array(wsgi_req->buffer, wsgi_req->uh.pktsize, argv, argvs, &argc)) { if (argc > 1) { - uwsgi_cache_set(argv[0], strlen(argv[0]), argv[1], strlen(argv[1]), 0, 0); + uwsgi_cache_set(argv[0], argvs[0], argv[1], argvs[1], 0, 0); } } } diff --git a/plugins/erlang/erlang.c b/plugins/erlang/erlang.c index 135a88c9..dfc26dba 100644 --- a/plugins/erlang/erlang.c +++ b/plugins/erlang/erlang.c @@ -94,7 +94,8 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) { char *call; char buffer[0xffff]; - char *argv[0xff] ; + char *argv[256] ; + uint16_t argvs[256] ; int argc = 0; uint16_t ret; ei_x_buff xr; @@ -215,14 +216,16 @@ void uwsgi_erlang_rpc(int fd, erlang_pid *from, ei_x_buff *x) { argc = 1; argv[0] = uwsgi_malloc(esize+1); ei_decode_atom(x->buff, &x->index, argv[0]); + argvs[1] = esize; } else if (etype == ERL_STRING_EXT) { argc = 1; argv[0] = uwsgi_malloc(esize+1); ei_decode_string(x->buff, &x->index, argv[0]); + argvs[1] = esize; } - ret = uwsgi_rpc(call, argc, argv, buffer); + ret = uwsgi_rpc(call, argc, argv, argvs, buffer); #ifdef UWSGI_DEBUG uwsgi_log("buffer: %.*s\n", ret, buffer); diff --git a/plugins/lua/lua_plugin.c b/plugins/lua/lua_plugin.c index 92afc35b..040cdd23 100644 --- a/plugins/lua/lua_plugin.c +++ b/plugins/lua/lua_plugin.c @@ -682,7 +682,7 @@ int uwsgi_lua_signal_handler(uint8_t sig, void *handler) { } -uint16_t uwsgi_lua_rpc(void * func, uint8_t argc, char **argv, char *buffer) { +uint16_t uwsgi_lua_rpc(void * func, uint8_t argc, char **argv, uint16_t argvs[], char *buffer) { uint8_t i; const char *sv; @@ -698,7 +698,7 @@ uint16_t uwsgi_lua_rpc(void * func, uint8_t argc, char **argv, char *buffer) { lua_rawgeti(L, LUA_REGISTRYINDEX, ifunc); for(i=0;i 0) { - ST(0) = newSVpv(buffer, size); + ST(0) = newSVpv(response, size); sv_2mortal(ST(0)); - + free(response); XSRETURN(1); } diff --git a/plugins/python/python_plugin.c b/plugins/python/python_plugin.c index 19a1fc6f..8d42c368 100644 --- a/plugins/python/python_plugin.c +++ b/plugins/python/python_plugin.c @@ -1291,36 +1291,38 @@ clear: return -1; } -uint16_t uwsgi_python_rpc(void *func, uint8_t argc, char **argv, char *buffer) { +uint16_t uwsgi_python_rpc(void *func, uint8_t argc, char **argv, uint16_t argvs[], char *buffer) { UWSGI_GET_GIL; uint8_t i; - PyObject *pyargs = PyTuple_New(argc); - PyObject *ret; char *rv; size_t rl; + PyObject *pyargs = PyTuple_New(argc); + PyObject *ret; + if (!pyargs) return 0; for (i = 0; i < argc; i++) { - PyTuple_SetItem(pyargs, i, PyString_FromString(argv[i])); + PyTuple_SetItem(pyargs, i, PyString_FromStringAndSize(argv[i], argvs[i])); } ret = python_call((PyObject *) func, pyargs, 0, NULL); - + Py_DECREF(pyargs); if (ret) { if (PyString_Check(ret)) { rv = PyString_AsString(ret); rl = PyString_Size(ret); - if (rl <= 0xffff) { + if (rl <= 65536) { memcpy(buffer, rv, rl); Py_DECREF(ret); UWSGI_RELEASE_GIL; return rl; } } + Py_DECREF(ret); } if (PyErr_Occurred()) diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index ecc2c05a..869b98f9 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -294,13 +294,13 @@ PyObject *py_uwsgi_add_file_monitor(PyObject * self, PyObject * args) { PyObject *py_uwsgi_call(PyObject * self, PyObject * args) { - char buffer[0xffff]; char *func; uint16_t size = 0; PyObject *py_func; int argc = PyTuple_Size(args); int i; - char *argv[0xff]; + char *argv[256]; + uint16_t argvs[256]; // TODO better error reporting if (argc < 1) @@ -314,19 +314,29 @@ PyObject *py_uwsgi_call(PyObject * self, PyObject * args) { func = PyString_AsString(py_func); for (i = 0; i < (argc - 1); i++) { - argv[i] = PyString_AsString(PyTuple_GetItem(args, i + 1)); + PyObject *py_str = PyTuple_GetItem(args, i + 1); + if (!PyString_Check(py_str)) { + goto clear; + } + argv[i] = PyString_AsString(py_str); + argvs[i] = PyString_Size(py_str); } - size = uwsgi_rpc(func, argc - 1, argv, buffer); + char *response = uwsgi_do_rpc(NULL, func, argc - 1, argv, argvs, &size); if (size > 0) { - return PyString_FromStringAndSize(buffer, size); + PyObject *ret = PyString_FromStringAndSize(response, size); + free(response); + return ret; } - clear: - Py_INCREF(Py_None); return Py_None; + + clear: + + return PyErr_Format(PyExc_ValueError, "unable to call rpc function"); + } PyObject *py_uwsgi_rpc_list(PyObject * self, PyObject * args) { @@ -346,19 +356,15 @@ PyObject *py_uwsgi_rpc_list(PyObject * self, PyObject * args) { PyObject *py_uwsgi_rpc(PyObject * self, PyObject * args) { - char buffer[0xffff]; - char *node, *func; + char *node = NULL, *func; uint16_t size = 0; PyObject *py_node, *py_func; - struct wsgi_request rpc_req; - int argc = PyTuple_Size(args); - char *argv[0xff]; - int i, fd; - uint16_t pktsize = 0, ulen; - char *bufptr; - int rlen; - int rpc_args = 0; + int argc = PyTuple_Size(args); + char *argv[256]; + uint16_t argvs[256]; + + int i; // TODO better error reporting if (argc < 2) @@ -369,9 +375,6 @@ PyObject *py_uwsgi_rpc(PyObject * self, PyObject * args) { if (PyString_Check(py_node)) { node = PyString_AsString(py_node); } - else { - node = ""; - } py_func = PyTuple_GetItem(args, 1); @@ -381,87 +384,28 @@ PyObject *py_uwsgi_rpc(PyObject * self, PyObject * args) { func = PyString_AsString(py_func); for (i = 0; i < (argc - 2); i++) { - argv[i] = PyString_AsString(PyTuple_GetItem(args, i + 2)); - rpc_args++; + PyObject *py_str = PyTuple_GetItem(args, i + 2); + if (!PyString_Check(py_str)) + goto clear; + argv[i] = PyString_AsString(py_str); + argvs[i] = PyString_Size(py_str); } - if (!strcmp(node, "")) { - if (!rpc_args) { - size = uwsgi_rpc(func, 0, NULL, buffer); - } - else { - size = uwsgi_rpc(func, rpc_args, argv, buffer); - } - } - else { - - - // connect to node - fd = uwsgi_connect(node, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], 0); - - if (fd < 0) - goto clear; - // prepare a uwsgi array - - pktsize = 2 + strlen(func); - for (i = 0; i < argc - 2; i++) { - pktsize += 2 + strlen(argv[i]); - } - - memset(&rpc_req, 0, sizeof(struct wsgi_request)); - - rpc_req.uh.modifier1 = 173; - rpc_req.uh.pktsize = pktsize; - rpc_req.uh.modifier2 = 0; - - bufptr = buffer; - - ulen = strlen(func); - *bufptr++ = (uint8_t) (ulen & 0xff); - *bufptr++ = (uint8_t) ((ulen >> 8) & 0xff); - memcpy(bufptr, func, ulen); - bufptr += ulen; - - for (i = 0; i < argc - 2; i++) { - ulen = strlen(argv[i]); - *bufptr++ = (uint8_t) (ulen & 0xff); - *bufptr++ = (uint8_t) ((ulen >> 8) & 0xff); - memcpy(bufptr, argv[i], ulen); - bufptr += ulen; - } - - if (write(fd, &rpc_req.uh, 4) != 4) { - uwsgi_error("write()"); - close(fd); - goto clear; - } - - if (write(fd, buffer, pktsize) != pktsize) { - uwsgi_error("write()"); - close(fd); - goto clear; - } - - rlen = uwsgi_waitfd(fd, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]); - if (rlen > 0) { - rpc_req.poll.fd = fd; - rpc_req.poll.events = POLLIN; - rpc_req.buffer = buffer; - if (uwsgi_parse_packet(&rpc_req, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT])) { - size = rpc_req.uh.pktsize; - } - } - - } + char *response = uwsgi_do_rpc(node, func, argc - 2, argv, argvs, &size); if (size > 0) { - return PyString_FromStringAndSize(buffer, size); - } + PyObject *ret = PyString_FromStringAndSize(response, size); + free(response); + return ret; + } + + Py_INCREF(Py_None); + return Py_None; clear: - Py_INCREF(Py_None); - return Py_None; + return PyErr_Format(PyExc_ValueError, "unable to call rpc function"); + } PyObject *py_uwsgi_register_rpc(PyObject * self, PyObject * args) { diff --git a/plugins/rack/rack_api.c b/plugins/rack/rack_api.c index cfab385e..1e012e31 100644 --- a/plugins/rack/rack_api.c +++ b/plugins/rack/rack_api.c @@ -595,6 +595,59 @@ VALUE uwsgi_ruby_signal_registered(VALUE *class, VALUE signum) { return Qfalse; } +VALUE uwsgi_ruby_do_rpc(int argc, VALUE *rpc_argv, VALUE *class) { + + char *node = NULL, *func; + uint16_t size = 0; + + char *argv[256]; + uint16_t argvs[256]; + + int i; + + + // TODO better error reporting + if (argc < 2) + goto clear; + + + VALUE rpc_node = rpc_argv[0]; + + if (TYPE(rpc_node) == T_STRING) { + node = RSTRING_PTR(rpc_node); + } + + + VALUE rpc_func = rpc_argv[1]; + + if (TYPE(rpc_func) != T_STRING) + goto clear; + + func = RSTRING_PTR(rpc_func); + + for (i = 0; i < (argc - 2); i++) { + VALUE rpc_str = rpc_argv[i + 2]; + if (TYPE(rpc_str) != T_STRING) + goto clear; + argv[i] = RSTRING_PTR(rpc_str); + argvs[i] = RSTRING_LEN(rpc_str); + } + + char *response = uwsgi_do_rpc(node, func, argc - 2, argv, argvs, &size); + + if (size > 0) { + VALUE ret = rb_str_new(response, size); + free(response); + return ret; + } + + +clear: + + rb_raise(rb_eRuntimeError, "unable to call rpc function"); + return Qnil; +} + VALUE uwsgi_ruby_register_rpc(int argc, VALUE *argv, VALUE *class) { int rb_argc = 0; @@ -712,6 +765,8 @@ void uwsgi_rack_init_api() { uwsgi_rack_api("set_warning_message", rack_uwsgi_warning, 1); + uwsgi_rack_api("rpc", uwsgi_ruby_do_rpc, -1); + if (uwsgi.cache_max_items > 0) { diff --git a/plugins/rack/rack_plugin.c b/plugins/rack/rack_plugin.c index 46e0aabc..3d042b02 100644 --- a/plugins/rack/rack_plugin.c +++ b/plugins/rack/rack_plugin.c @@ -157,7 +157,7 @@ VALUE rack_call_rpc_handler(VALUE args) { } -uint16_t uwsgi_ruby_rpc(void *func, uint8_t argc, char **argv, char *buffer) { +uint16_t uwsgi_ruby_rpc(void *func, uint8_t argc, char **argv, uint16_t argvs[], char *buffer) { uint8_t i; VALUE rb_args = rb_ary_new2(2); @@ -170,7 +170,7 @@ uint16_t uwsgi_ruby_rpc(void *func, uint8_t argc, char **argv, char *buffer) { rb_ary_store(rb_args, 0, (VALUE) func); for (i = 0; i < argc; i++) { - rb_ary_store(rb_rpc_args, i, rb_str_new2(argv[i])); + rb_ary_store(rb_rpc_args, i, rb_str_new(argv[i], argvs[i])); } rb_ary_store(rb_args, 1, rb_rpc_args); diff --git a/plugins/rpc/rpc_plugin.c b/plugins/rpc/rpc_plugin.c index 3c8c1d35..505bd6ea 100644 --- a/plugins/rpc/rpc_plugin.c +++ b/plugins/rpc/rpc_plugin.c @@ -6,6 +6,7 @@ extern struct uwsgi_server uwsgi; int uwsgi_rpc_request(struct wsgi_request *wsgi_req) { char *argv[256]; + uint16_t argvs[256]; uint8_t argc = 0xff; /* Standard RPC request */ @@ -24,7 +25,7 @@ int uwsgi_rpc_request(struct wsgi_request *wsgi_req) { uwsgi_log("RPC pktsize %d\n", wsgi_req->uh.pktsize); #endif - if (uwsgi_parse_array(wsgi_req->buffer, wsgi_req->uh.pktsize, argv, &argc)) { + if (uwsgi_parse_array(wsgi_req->buffer, wsgi_req->uh.pktsize, argv, argvs, &argc)) { uwsgi_log("Invalid RPC request. skip.\n"); return -1; } @@ -33,7 +34,7 @@ int uwsgi_rpc_request(struct wsgi_request *wsgi_req) { uwsgi_log("RPC args %d\n", argc-1); #endif - wsgi_req->uh.pktsize = uwsgi_rpc(argv[0], argc-1, argv+1, wsgi_req->buffer); + wsgi_req->uh.pktsize = uwsgi_rpc(argv[0], argc-1, argv+1, argvs+1, wsgi_req->buffer); if (wsgi_req->uh.modifier2 == 0) { wsgi_req->headers_size = wsgi_req->socket->proto_write_header(wsgi_req, (char *)&wsgi_req->uh, 4); diff --git a/protocol.c b/protocol.c index d0d9a866..bddf0328 100644 --- a/protocol.c +++ b/protocol.c @@ -371,7 +371,8 @@ int uwsgi_read_response(int fd, struct uwsgi_header *uh, int timeout, char **buf } if (buf && uh->pktsize > 0) { - *buf = uwsgi_malloc(uh->pktsize); + if (*buf == NULL) + *buf = uwsgi_malloc(uh->pktsize); remains = uh->pktsize; ptr = *buf; ret = -1; @@ -434,7 +435,7 @@ int uwsgi_parse_packet(struct wsgi_request *wsgi_req, int timeout) { return 1; } -int uwsgi_parse_array(char *buffer, uint16_t size, char **argv, uint8_t *argc) { +int uwsgi_parse_array(char *buffer, uint16_t size, char **argv, uint16_t argvs[], uint8_t *argc) { char *ptrbuf, *bufferend; uint16_t strsize = 0; @@ -459,6 +460,7 @@ int uwsgi_parse_array(char *buffer, uint16_t size, char **argv, uint8_t *argc) { if (ptrbuf + strsize <= bufferend) { // item argv[*argc] = uwsgi_cheap_string(ptrbuf, strsize); + argvs[*argc] = strsize; #ifdef UWSGI_DEBUG uwsgi_log("arg %s\n", argv[*argc]); #endif diff --git a/rpc.c b/rpc.c index 492b9278..aa7c5b4f 100644 --- a/rpc.c +++ b/rpc.c @@ -33,7 +33,7 @@ int uwsgi_register_rpc(char *name, uint8_t modifier1, uint8_t args, void *func) return ret; } -uint16_t uwsgi_rpc(char *name, uint8_t argc, char *argv[], char *output) { +uint16_t uwsgi_rpc(char *name, uint8_t argc, char *argv[], uint16_t argvs[], char *output) { struct uwsgi_rpc *urpc = NULL; int i; @@ -50,9 +50,93 @@ uint16_t uwsgi_rpc(char *name, uint8_t argc, char *argv[], char *output) { if (urpc) { if (uwsgi.p[urpc->modifier1]->rpc) { - ret = uwsgi.p[urpc->modifier1]->rpc(urpc->func, argc, argv, output); + ret = uwsgi.p[urpc->modifier1]->rpc(urpc->func, argc, argv, argvs, output); } } return ret; } + + +char *uwsgi_do_rpc(char *node, char *func, uint8_t argc, char *argv[], uint16_t argvs[], uint16_t *len) { + + uint8_t i; + uint16_t ulen; + struct uwsgi_header uh; + char *buffer = NULL; + + *len = 0; + + if (node == NULL || !strcmp(node, "")) { + // allocate the whole buffer + buffer = uwsgi_malloc(65536); + *len = uwsgi_rpc(func, argc, argv, argvs, buffer); + return buffer; + } + + + // connect to node + int fd = uwsgi_connect(node, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], 0); + + if (fd < 0) return NULL; + + // prepare a uwsgi array + uint16_t buffer_size = 2 + strlen(func); + + for (i = 0; i < argc; i++) { + buffer_size += 2 + argvs[i]; + } + + // allocate the whole buffer + buffer = uwsgi_malloc(65536); + + uh.modifier1 = 173; + uh.pktsize = buffer_size; + uh.modifier2 = 0; + + // add func to the array + char *bufptr = buffer; + ulen = strlen(func); + *bufptr++ = (uint8_t) (ulen & 0xff); + *bufptr++ = (uint8_t) ((ulen >> 8) & 0xff); + memcpy(bufptr, func, ulen); + bufptr += ulen; + + for (i = 0; i < argc; i++) { + ulen = argvs[i]; + *bufptr++ = (uint8_t) (ulen & 0xff); + *bufptr++ = (uint8_t) ((ulen >> 8) & 0xff); + memcpy(bufptr, argv[i], ulen); + bufptr += ulen; + } + + if (write(fd, &uh, 4) != 4) { + uwsgi_error("write()"); + close(fd); + free(buffer); + return NULL; + } + + if (write(fd, buffer, buffer_size) != buffer_size) { + uwsgi_error("write()"); + close(fd); + free(buffer); + return NULL; + } + + if (uwsgi_read_response(fd, &uh, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], &buffer) < 0) { + close(fd); + free(buffer); + return NULL; + } + + close(fd); + + *len = uh.pktsize; + if (*len == 0) { + free(buffer); + return NULL; + } + return buffer; + +} diff --git a/uwsgi.h b/uwsgi.h index 70a307f5..e68901bd 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -678,7 +678,7 @@ struct uwsgi_plugin { int (*spooler) (char *, char *, uint16_t, char *, size_t); - uint16_t(*rpc) (void *, uint8_t, char **, char *); + uint16_t(*rpc) (void *, uint8_t, char **, uint16_t *, char *); void (*jail) (int (*)(void *), char **); @@ -2103,11 +2103,12 @@ void uwsgi_route_signal(uint8_t); int uwsgi_start(void *); int uwsgi_register_rpc(char *, uint8_t, uint8_t, void *); -uint16_t uwsgi_rpc(char *, uint8_t, char **, char *); +uint16_t uwsgi_rpc(char *, uint8_t, char **, uint16_t *, char *); +char *uwsgi_do_rpc(char *, char *, uint8_t, char **, uint16_t *, uint16_t *); char *uwsgi_cheap_string(char *, int); -int uwsgi_parse_array(char *, uint16_t, char **, uint8_t *); +int uwsgi_parse_array(char *, uint16_t, char **, uint16_t *, uint8_t *); void log_syslog(char *);