diff --git a/core/init.c b/core/init.c index 9fb10d72..95e0848f 100644 --- a/core/init.c +++ b/core/init.c @@ -91,7 +91,7 @@ void uwsgi_init_default() { uwsgi.threads = 1; // default max number of rpc slot - uwsgi.max_rpc = 64; + uwsgi.rpc_max = 64; uwsgi.offload_threads_events = 64; diff --git a/core/io.c b/core/io.c index 3bbfaba0..7cab6506 100644 --- a/core/io.c +++ b/core/io.c @@ -2,6 +2,13 @@ extern struct uwsgi_server uwsgi; +/* + + poll based fd waiter. + Use it for blocking areas (like startup functions) + DO NOT USE IN REQUEST PLUGINS !!! + +*/ int uwsgi_waitfd_event(int fd, int timeout, int event) { int ret; @@ -32,6 +39,9 @@ int uwsgi_waitfd_event(int fd, int timeout, int event) { return ret; } +/* + consume data from an fd (blocking) +*/ char *uwsgi_read_fd(int fd, size_t *size, int add_zero) { char stack_buf[4096]; @@ -63,6 +73,7 @@ char *uwsgi_read_fd(int fd, size_t *size, int add_zero) { } +// simply read the whole content of a file char *uwsgi_simple_file_read(char *filename) { struct stat sb; @@ -101,6 +112,10 @@ end: } +/* + extremely complex function for reading resources (files, url...) + need a lot of refactoring... +*/ char *uwsgi_open_and_read(char *url, size_t *size, int add_zero, char *magic_table[]) { int fd; @@ -456,6 +471,7 @@ end: return buffer; } +// attach an fd using UNIX sockets int *uwsgi_attach_fd(int fd, int *count_ptr, char *code, size_t code_len) { struct msghdr msg; @@ -544,6 +560,7 @@ int *uwsgi_attach_fd(int fd, int *count_ptr, char *code, size_t code_len) { return ret; } +// signal free close void uwsgi_protected_close(int fd) { sigset_t mask, oset; @@ -559,6 +576,7 @@ void uwsgi_protected_close(int fd) { } } +// signal free read ssize_t uwsgi_protected_read(int fd, void *buf, size_t len) { sigset_t mask, oset; @@ -577,6 +595,8 @@ ssize_t uwsgi_protected_read(int fd, void *buf, size_t len) { return ret; } + +// pipe datas from a fd to another (blocking) ssize_t uwsgi_pipe(int src, int dst, int timeout) { char buf[8192]; size_t written = -1; @@ -633,6 +653,9 @@ timeout: return -1; } +/* + even if it is marked as non-blocking, so not use in request plugins as it uses poll() and not the hooks +*/ int uwsgi_write_nb(int fd, char *buf, size_t remains, int timeout) { char *ptr = buf; while(remains > 0) { @@ -652,6 +675,39 @@ int uwsgi_write_nb(int fd, char *buf, size_t remains, int timeout) { return 0; } +/* + this is like uwsgi_write_nb() but with fast initial write and hooked wait (use it in request plugin) +*/ +int uwsgi_write_true_nb(int fd, char *buf, size_t remains, int timeout) { + char *ptr = buf; + int ret; + + while(remains > 0) { + ssize_t len = write(fd, ptr, remains); + if (len > 0) goto written; + if (len == 0) return -1; + if (len < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) goto wait; + return -1; + } +wait: + ret = uwsgi.wait_write_hook(fd, timeout); + if (ret > 0) { + len = write(fd, ptr, remains); + if (len > 0) goto written; + } + return -1; +written: + ptr += len; + remains -= len; + continue; + } + + return 0; +} + + +// like uwsgi_pipe but with fixed size ssize_t uwsgi_pipe_sized(int src, int dst, size_t required, int timeout) { char buf[8192]; size_t written = 0; @@ -709,6 +765,7 @@ timeout: } +// check if an fd is valid int uwsgi_valid_fd(int fd) { int ret = fcntl(fd, F_GETFL); if (ret == 0) { @@ -766,3 +823,79 @@ int uwsgi_read_nb(int fd, char *buf, size_t remains, int timeout) { return 0; } + +/* + this is a pretty magic function used for read a full uwsgi response + it is true non blocking, so you can use it in request plugins + buffer is expected to be at least 4 bytes, rlen is a get/set value +*/ + +int uwsgi_read_with_realloc(int fd, char **buffer, size_t *rlen, int timeout) { + if (*rlen < 4) return -1; + char *buf = *buffer; + int ret; + + // start reading the header + char *ptr = buf; + size_t remains = 4; + while(remains > 0) { + ssize_t len = read(fd, ptr, remains); + if (len > 0) goto readok; + if (len == 0) return -1; + if (len < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) goto wait; + return -1; + } +wait: + ret = uwsgi.wait_read_hook(fd, timeout); + if (ret > 0) { + len = read(fd, ptr, remains); + if (len > 0) goto readok; + } + return -1; +readok: + ptr += len; + remains -= len; + continue; + } + + struct uwsgi_header *uh = (struct uwsgi_header *) buf; + uint16_t pktsize = uh->pktsize; + + if (pktsize > *rlen) { + char *tmp_buf = realloc(buf, pktsize); + if (!tmp_buf) { + uwsgi_error("uwsgi_read_with_realloc()/realloc()"); + return -1; + } + *buffer = tmp_buf; + } + + *rlen = pktsize; + // read the body + remains = pktsize; + ptr = buf; + while(remains > 0) { + ssize_t len = read(fd, ptr, remains); + if (len > 0) goto readok2; + if (len == 0) return -1; + if (len < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) goto wait2; + return -1; + } +wait2: + ret = uwsgi.wait_read_hook(fd, timeout); + if (ret > 0) { + len = read(fd, ptr, remains); + if (len > 0) goto readok2; + } + return -1; +readok2: + ptr += len; + remains -= len; + continue; + } + + return 0; + +} diff --git a/core/rpc.c b/core/rpc.c index 78b1e285..dc8efdc6 100644 --- a/core/rpc.c +++ b/core/rpc.c @@ -62,7 +62,7 @@ char *uwsgi_do_rpc(char *node, char *func, uint8_t argc, char *argv[], uint16_t uint8_t i; uint16_t ulen; - struct uwsgi_header uh; + struct uwsgi_header *uh = NULL; char *buffer = NULL; *len = 0; @@ -75,12 +75,17 @@ char *uwsgi_do_rpc(char *node, char *func, uint8_t argc, char *argv[], uint16_t } - // connect to node - int fd = uwsgi_connect(node, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], 0); - + // connect to node (async way) + int fd = uwsgi_connect(node, 0, 1); if (fd < 0) return NULL; + // wait for connection; + int ret = uwsgi.wait_write_hook(fd, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]); + if (ret <= 0) { + return NULL; + } + // prepare a uwsgi array uint16_t buffer_size = 2 + strlen(func); @@ -89,14 +94,16 @@ char *uwsgi_do_rpc(char *node, char *func, uint8_t argc, char *argv[], uint16_t } // allocate the whole buffer - buffer = uwsgi_malloc(65536); + buffer = uwsgi_malloc(4+buffer_size); - uh.modifier1 = 173; - uh.pktsize = buffer_size; - uh.modifier2 = 0; + // set the uwsgi header + uh = (struct uwsgi_header *) buffer; + uh->modifier1 = 173; + uh->pktsize = buffer_size; + uh->modifier2 = 0; // add func to the array - char *bufptr = buffer; + char *bufptr = buffer + 4; ulen = strlen(func); *bufptr++ = (uint8_t) (ulen & 0xff); *bufptr++ = (uint8_t) ((ulen >> 8) & 0xff); @@ -111,33 +118,32 @@ char *uwsgi_do_rpc(char *node, char *func, uint8_t argc, char *argv[], uint16_t bufptr += ulen; } - if (write(fd, &uh, 4) != 4) { - uwsgi_error("write()"); - close(fd); - free(buffer); - return NULL; + // ok the reuqest is ready, let's send it in non blocking way + if (uwsgi_write_true_nb(fd, buffer, buffer_size+4, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT])) { + goto error; } - 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; + // ok time to wait for the response in non blocking way + size_t rlen = buffer_size+4; + if (uwsgi_read_with_realloc(fd, &buffer, &rlen, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT])) { + goto error; } close(fd); - - *len = uh.pktsize; + *len = rlen; if (*len == 0) { - free(buffer); - return NULL; + goto error; } return buffer; +error: + close(fd); + free(buffer); + return NULL; + +} + + +void uwsgi_rpc_init() { + uwsgi.rpc_table = uwsgi_calloc_shared(sizeof(struct uwsgi_rpc) * uwsgi.rpc_max); } diff --git a/core/uwsgi.c b/core/uwsgi.c index 87d03e8f..3d329808 100644 --- a/core/uwsgi.c +++ b/core/uwsgi.c @@ -219,6 +219,8 @@ static struct uwsgi_option uwsgi_base_options[] = { {"signal-bufsize", required_argument, 0, "set buffer size for signal queue", uwsgi_opt_set_int, &uwsgi.signal_bufsize, 0}, {"signals-bufsize", required_argument, 0, "set buffer size for signal queue", uwsgi_opt_set_int, &uwsgi.signal_bufsize, 0}, + {"rpc-max", required_argument, 0, "maximum number of rpc slots (default: 64)", uwsgi_opt_set_64bit, &uwsgi.rpc_max, 0}, + {"disable-logging", no_argument, 'L', "disable request logging", uwsgi_opt_dyn_false, (void *) UWSGI_OPTION_LOGGING, 0}, {"flock", required_argument, 0, "lock the specified file before starting, exit if locked", uwsgi_opt_flock, NULL, UWSGI_OPT_IMMEDIATE}, @@ -2452,6 +2454,9 @@ unsafe: } } + // allocate rpc structures + uwsgi_rpc_init(); + // set masterpid uwsgi.mypid = getpid(); masterpid = uwsgi.mypid; diff --git a/plugins/php/php_plugin.c b/plugins/php/php_plugin.c index 1216918e..fcb20550 100644 --- a/plugins/php/php_plugin.c +++ b/plugins/php/php_plugin.c @@ -393,14 +393,12 @@ PHP_FUNCTION(uwsgi_rpc) { // response must always be freed char *response = uwsgi_do_rpc(node, func, num_args - 2, argv, argvs, &size); - - if (size > 0) { + if (response) { // here we do not free varargs for performance reasons char *ret = estrndup(response, size); free(response); RETURN_STRING(ret, 0); } - free(response); clear: efree(varargs); diff --git a/plugins/psgi/uwsgi_plmodule.c b/plugins/psgi/uwsgi_plmodule.c index a59d5607..2459f70a 100644 --- a/plugins/psgi/uwsgi_plmodule.c +++ b/plugins/psgi/uwsgi_plmodule.c @@ -230,14 +230,12 @@ XS(XS_call) { // response must be always freed char *response = uwsgi_do_rpc(NULL, func, items-1, argv, argvs, &size); - - if (size > 0) { + if (response) { ST(0) = newSVpv(response, size); sv_2mortal(ST(0)); free(response); XSRETURN(1); } - free(response); XSRETURN_UNDEF; } diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index 5b9bb3ab..563117f6 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -316,13 +316,12 @@ PyObject *py_uwsgi_call(PyObject * self, PyObject * args) { char *response = uwsgi_do_rpc(NULL, func, argc - 1, argv, argvs, &size); UWSGI_GET_GIL; - if (size > 0) { + if (response) { PyObject *ret = PyString_FromStringAndSize(response, size); free(response); return ret; } - free(response); Py_INCREF(Py_None); return Py_None; @@ -388,13 +387,12 @@ PyObject *py_uwsgi_rpc(PyObject * self, PyObject * args) { char *response = uwsgi_do_rpc(node, func, argc - 2, argv, argvs, &size); UWSGI_GET_GIL; - if (size > 0) { + if (response) { PyObject *ret = PyString_FromStringAndSize(response, size); free(response); return ret; } - free(response); Py_INCREF(Py_None); return Py_None; diff --git a/plugins/rack/rack_api.c b/plugins/rack/rack_api.c index 877f84e4..9c1f32ac 100644 --- a/plugins/rack/rack_api.c +++ b/plugins/rack/rack_api.c @@ -666,14 +666,11 @@ VALUE uwsgi_ruby_do_rpc(int argc, VALUE *rpc_argv, VALUE *class) { // response must always be freed char *response = uwsgi_do_rpc(node, func, argc - 2, argv, argvs, &size); - - if (size > 0) { + if (response) { VALUE ret = rb_str_new(response, size); free(response); return ret; } - free(response); - clear: rb_raise(rb_eRuntimeError, "unable to call rpc function"); diff --git a/uwsgi.h b/uwsgi.h index 618db088..9cfaa9f2 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -2183,12 +2183,12 @@ struct uwsgi_server { }; - struct uwsgi_rpc { - char name[0xff]; - void *func; - uint8_t args; - uint8_t modifier1; - }; +struct uwsgi_rpc { + char name[0xff]; + void *func; + uint8_t args; + uint8_t modifier1; +}; struct uwsgi_signal_entry { int wid; @@ -2675,6 +2675,7 @@ void uwsgi_cache_start_sync_servers(void); int uwsgi_register_rpc(char *, uint8_t, uint8_t, void *); uint16_t uwsgi_rpc(char *, uint8_t, char **, uint16_t *, char *); char *uwsgi_do_rpc(char *, char *, uint8_t, char **, uint16_t *, uint16_t *); +void uwsgi_rpc_init(void); char *uwsgi_cheap_string(char *, int); @@ -2921,13 +2922,16 @@ struct uwsgi_dict { void uwsgi_string_del_list(struct uwsgi_string_list **, struct uwsgi_string_list *); - void uwsgi_init_all_apps(void); - void uwsgi_init_worker_mount_apps(void); - void uwsgi_socket_nb(int); - void uwsgi_socket_b(int); - int uwsgi_write_nb(int, char *, size_t, int); - int uwsgi_read_nb(int, char *, size_t, int); - int uwsgi_read_uh(int fd, struct uwsgi_header *, int); +void uwsgi_init_all_apps(void); +void uwsgi_init_worker_mount_apps(void); +void uwsgi_socket_nb(int); +void uwsgi_socket_b(int); +int uwsgi_write_nb(int, char *, size_t, int); +int uwsgi_read_nb(int, char *, size_t, int); +int uwsgi_read_uh(int fd, struct uwsgi_header *, int); + +int uwsgi_read_with_realloc(int, char **, size_t *, int); +int uwsgi_write_true_nb(int, char *, size_t, int); void uwsgi_destroy_request(struct wsgi_request *);