diff --git a/async.c b/async.c index 38ad3fa1..f72e9be7 100644 --- a/async.c +++ b/async.c @@ -2,6 +2,10 @@ extern struct uwsgi_server uwsgi; +struct wsgi_request *find_wsgi_req_proto_by_fd(int fd) { + return uwsgi.async_proto_fd_table[fd]; +} + struct wsgi_request *find_wsgi_req_by_fd(int fd) { return uwsgi.async_waiting_fd_table[fd]; } @@ -195,6 +199,9 @@ void *async_loop(void *arg1) { int interesting_fd, i; struct uwsgi_rb_timer *min_timeout; int timeout; + int j; + int is_a_new_connection; + int proto_parser_status; static struct uwsgi_async_request *current_request = NULL, *next_async_request = NULL; @@ -235,49 +242,86 @@ void *async_loop(void *arg1) { async_expire_timeouts(); } + for(i=0;iasync_id ); - if (wsgi_req_simple_accept(uwsgi.wsgi_req, interesting_fd)) { + // new request coming in ? + + for(j=0;jasync_id ); + if (wsgi_req_simple_accept(uwsgi.wsgi_req, interesting_fd)) { +#ifdef UWSGI_EVENT_USE_PORT + event_queue_add_fd_read(uwsgi.async_queue, interesting_fd); +#endif + break; + } #ifdef UWSGI_EVENT_USE_PORT event_queue_add_fd_read(uwsgi.async_queue, interesting_fd); #endif - continue; - } -#ifdef UWSGI_EVENT_USE_PORT - event_queue_add_fd_read(uwsgi.async_queue, interesting_fd); -#endif // on linux we do not need to reset the socket to blocking state #ifndef __linux__ - if (uwsgi.numproc > 1) { - /* re-set blocking socket */ - if (fcntl(uwsgi.wsgi_req->poll.fd, F_SETFL, uwsgi.sockets[0].arg) < 0) { - uwsgi_error("fcntl()"); - continue; + if (uwsgi.numproc > 1) { + /* re-set blocking socket */ + if (fcntl(uwsgi.wsgi_req->poll.fd, F_SETFL, uwsgi.sockets[j].arg) < 0) { + uwsgi_error("fcntl()"); + break; + } } - } #endif + if (wsgi_req_async_recv(uwsgi.wsgi_req, j)) { + break; + } - if (wsgi_req_simple_recv(uwsgi.wsgi_req)) { + break; + } + } + + if (!is_a_new_connection) { + // proto event + uwsgi.wsgi_req = find_wsgi_req_proto_by_fd(interesting_fd); + if (uwsgi.wsgi_req) { + proto_parser_status = uwsgi.wsgi_req->socket_proto(uwsgi.wsgi_req); + // reset timeout + rb_erase(&uwsgi.wsgi_req->async_timeout->rbt, uwsgi.rb_async_timeouts); + free(uwsgi.wsgi_req->async_timeout); + uwsgi.wsgi_req->async_timeout = NULL; + // parsing complete + if (!proto_parser_status) { + // remove fd from event poll and fd proto table + event_queue_del_fd(uwsgi.async_queue, interesting_fd, event_queue_read()); + uwsgi.async_proto_fd_table[interesting_fd] = NULL; + // put request in the runqueue + runqueue_push(uwsgi.wsgi_req); + continue; + } + else if (proto_parser_status == -1) { + uwsgi_log("error parsing request\n"); + uwsgi.async_proto_fd_table[interesting_fd] = NULL; + close(interesting_fd); + continue; + } + // re-add timer + async_add_timeout(uwsgi.wsgi_req, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]); continue; } - // put request in the runqueue - runqueue_push(uwsgi.wsgi_req); - } - else { // app event uwsgi.wsgi_req = find_wsgi_req_by_fd(interesting_fd); // unknown fd, remove it (for safety) @@ -324,7 +368,7 @@ void *async_loop(void *arg1) { next_async_request = current_request->next; // request ended ? - if (uwsgi.wsgi_req->async_status == UWSGI_OK) { + if (uwsgi.wsgi_req->async_status <= UWSGI_OK) { // remove all the monitored fds and timeout while(uwsgi.wsgi_req->waiting_fds) { #ifndef UWSGI_EVENT_USE_PORT diff --git a/cache.c b/cache.c index f8d251bf..03197719 100644 --- a/cache.c +++ b/cache.c @@ -497,17 +497,17 @@ void cache_command(char *key, uint16_t keylen, char *val, uint16_t vallen, void if (!uwsgi_strncmp(key, keylen, "key", 3)) { val = uwsgi_cache_get(val, vallen, &tmp_vallen); if (val && tmp_vallen > 0) { - wsgi_req->response_size = write(wsgi_req->poll.fd, val, tmp_vallen); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, val, tmp_vallen); } } else if (!uwsgi_strncmp(key, keylen, "get", 3)) { val = uwsgi_cache_get(val, vallen, &tmp_vallen); if (val && vallen > 0) { - wsgi_req->response_size = write(wsgi_req->poll.fd, val, tmp_vallen); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, val, tmp_vallen); } else { - wsgi_req->response_size = write(wsgi_req->poll.fd, "HTTP/1.0 404 Not Found\r\n\r\n

Not Found

", 44); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, "HTTP/1.0 404 Not Found\r\n\r\n

Not Found

", 44); } } } @@ -527,8 +527,8 @@ int uwsgi_cache_request(struct wsgi_request *wsgi_req) { value = uwsgi_cache_get(wsgi_req->buffer, wsgi_req->uh.pktsize, &vallen); if (value && vallen > 0) { wsgi_req->uh.pktsize = vallen; - wsgi_req->response_size = write(wsgi_req->poll.fd, &wsgi_req->uh, 4); - wsgi_req->response_size += write(wsgi_req->poll.fd, value, vallen); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, (char *)&wsgi_req->uh, 4); + wsgi_req->response_size += wsgi_req->socket_proto_write(wsgi_req, value, vallen); } } break; @@ -563,13 +563,13 @@ int uwsgi_cache_request(struct wsgi_request *wsgi_req) { if (value && vallen > 0) { wsgi_req->uh.pktsize = 0; wsgi_req->uh.modifier2 = 1; - wsgi_req->response_size = write(wsgi_req->poll.fd, &wsgi_req->uh, 4); - wsgi_req->response_size += write(wsgi_req->poll.fd, value, vallen); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, (char *)&wsgi_req->uh, 4); + wsgi_req->response_size += wsgi_req->socket_proto_write(wsgi_req, value, vallen); } else { wsgi_req->uh.pktsize = 0; wsgi_req->uh.modifier2 = 0; - wsgi_req->response_size = write(wsgi_req->poll.fd, &wsgi_req->uh, 4); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, (char *)&wsgi_req->uh, 4); } } break; diff --git a/plugins/echo/echo_plugin.c b/plugins/echo/echo_plugin.c index 2a4282f8..d4c8b8a0 100644 --- a/plugins/echo/echo_plugin.c +++ b/plugins/echo/echo_plugin.c @@ -5,7 +5,7 @@ extern struct uwsgi_server uwsgi; int uwsgi_echo_request(struct wsgi_request *wsgi_req) { - wsgi_req->response_size = write(wsgi_req->poll.fd, wsgi_req->buffer, wsgi_req->uh.pktsize); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, wsgi_req->buffer, wsgi_req->uh.pktsize); return 0; } diff --git a/plugins/example/example_plugin.c b/plugins/example/example_plugin.c index f89cb74b..728afd0d 100644 --- a/plugins/example/example_plugin.c +++ b/plugins/example/example_plugin.c @@ -11,7 +11,7 @@ int uwsgi_request(struct wsgi_request *wsgi_req) { char *http = "HTTP/1.1 200 Ok\r\nContent-type: text/html\r\n\r\n

Hello World

"; - wsgi_req->response_size += write(wsgi_req->poll.fd, http, strlen(http)); + wsgi_req->response_size += wsgi_req->socket_proto_write(wsgi_req, http, strlen(http)); return 0; } diff --git a/plugins/nagios/nagios.c b/plugins/nagios/nagios.c index 383bcb7a..cb5771b1 100644 --- a/plugins/nagios/nagios.c +++ b/plugins/nagios/nagios.c @@ -48,7 +48,7 @@ int nagios() { exit(2); } nagios_poll.events = POLLIN; - if (!uwsgi_parse_response(&nagios_poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) uwsgi.wsgi_req, uwsgi.wsgi_req->buffer)) { + if (!uwsgi_parse_response(&nagios_poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) uwsgi.wsgi_req, uwsgi.wsgi_req->buffer, uwsgi_proto_uwsgi_parser)) { fprintf(stdout, "UWSGI CRITICAL: timed out waiting for response\n"); exit(2); } diff --git a/plugins/ping/ping_plugin.c b/plugins/ping/ping_plugin.c index 6cca28f9..d82c30aa 100644 --- a/plugins/ping/ping_plugin.c +++ b/plugins/ping/ping_plugin.c @@ -36,7 +36,7 @@ static void ping() { exit(2); } uwsgi_poll.events = POLLIN; - if (!uwsgi_parse_response(&uwsgi_poll, uping.ping_timeout, &uh, NULL)) { + if (!uwsgi_parse_response(&uwsgi_poll, uping.ping_timeout, &uh, NULL, uwsgi_proto_uwsgi_parser)) { exit(1); } else { diff --git a/plugins/psgi/psgi_response.c b/plugins/psgi/psgi_response.c index 016f4cd0..29da9fc5 100644 --- a/plugins/psgi/psgi_response.c +++ b/plugins/psgi/psgi_response.c @@ -122,7 +122,7 @@ int psgi_response(struct wsgi_request *wsgi_req, PerlInterpreter *my_perl, AV *r vi = (i*2)+base; wsgi_req->hvec[vi].iov_base = "\r\n"; wsgi_req->hvec[vi].iov_len = 2; - if ( !(wsgi_req->headers_size = writev(wsgi_req->poll.fd, wsgi_req->hvec, vi+1)) ) { + if ( !(wsgi_req->headers_size = wsgi_req->socket_proto_writev_header(wsgi_req, wsgi_req->hvec, vi+1)) ) { uwsgi_error("writev()"); } @@ -147,7 +147,7 @@ int psgi_response(struct wsgi_request *wsgi_req, PerlInterpreter *my_perl, AV *r if (hlen <= 0) { break; } - wsgi_req->response_size = write(wsgi_req->poll.fd, chitem, hlen); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, chitem, hlen); } @@ -159,7 +159,7 @@ int psgi_response(struct wsgi_request *wsgi_req, PerlInterpreter *my_perl, AV *r for(i=0; i<=av_len(body); i++) { hitem = av_fetch(body,i,0); chitem = SvPV(*hitem, hlen); - wsgi_req->response_size = write(wsgi_req->poll.fd, chitem, hlen); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, chitem, hlen); } } diff --git a/plugins/python/uwsgi_pymodule.c b/plugins/python/uwsgi_pymodule.c index 0eaca7d2..7243ca79 100644 --- a/plugins/python/uwsgi_pymodule.c +++ b/plugins/python/uwsgi_pymodule.c @@ -397,7 +397,7 @@ PyObject *py_uwsgi_rpc(PyObject * self, PyObject * args) { if (rlen > 0) { upoll.fd = fd; upoll.events = POLLIN; - if (uwsgi_parse_response(&upoll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], &uh, buffer)) { + if (uwsgi_parse_response(&upoll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], &uh, buffer, uwsgi_proto_uwsgi_parser)) { size = uh.pktsize; } } @@ -1353,7 +1353,7 @@ PyObject *py_uwsgi_send_multi_message(PyObject * self, PyObject * args) { else { for (i = 0; i < clen; i++) { if (multipoll[i].revents & POLLIN) { - if (!uwsgi_parse_response(&multipoll[i], PyInt_AsLong(arg_timeout), &uh, &buffer[i])) { + if (!uwsgi_parse_response(&multipoll[i], PyInt_AsLong(arg_timeout), &uh, &buffer[i], uwsgi_proto_uwsgi_parser)) { goto megamulticlear; } else { diff --git a/plugins/python/wsgi_handlers.c b/plugins/python/wsgi_handlers.c index d02447a9..02ac0aac 100644 --- a/plugins/python/wsgi_handlers.c +++ b/plugins/python/wsgi_handlers.c @@ -171,7 +171,7 @@ PyObject *py_uwsgi_write(PyObject * self, PyObject * args) { content = PyString_AsString(data); len = PyString_Size(data); UWSGI_RELEASE_GIL - wsgi_req->response_size = write(wsgi_req->poll.fd, content, len); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, content, len); UWSGI_GET_GIL } @@ -354,6 +354,7 @@ int uwsgi_request_wsgi(struct wsgi_request *wsgi_req) { goto clear; } + if (strncmp(wsgi_req->protocol, "HTTP/", 5)) { uwsgi_log( "INVALID PROTOCOL: %.*s\n", wsgi_req->protocol_len, wsgi_req->protocol); internal_server_error(wsgi_req->poll.fd, "invalid HTTP protocol !!!"); @@ -410,37 +411,40 @@ int uwsgi_request_wsgi(struct wsgi_request *wsgi_req) { - if (uwsgi.post_buffering > 0) { - UWSGI_RELEASE_GIL - // read to disk if post_cl > post_buffering - if (!up.pep3333_input) { - if (wsgi_req->post_cl >= (size_t) uwsgi.post_buffering) { - if (!uwsgi_read_whole_body(wsgi_req, wsgi_req->post_buffering_buf, uwsgi.post_buffering_bufsize)) { - goto clear; + // some protocol (http included) pass the body directly as a FILE object + if (!wsgi_req->async_post) { + if (uwsgi.post_buffering > 0) { + UWSGI_RELEASE_GIL + // read to disk if post_cl > post_buffering + if (!up.pep3333_input) { + if (wsgi_req->post_cl >= (size_t) uwsgi.post_buffering) { + if (!uwsgi_read_whole_body(wsgi_req, wsgi_req->post_buffering_buf, uwsgi.post_buffering_bufsize)) { + goto clear; + } + } + else { + wsgi_req->async_post = fdopen(wsgi_req->poll.fd, "r"); } } else { - wsgi_req->async_post = fdopen(wsgi_req->poll.fd, "r"); + // read to disk if post_cl > post_buffering + if (wsgi_req->post_cl >= (size_t) uwsgi.post_buffering) { + if (!uwsgi_read_whole_body(wsgi_req, wsgi_req->post_buffering_buf, uwsgi.post_buffering_bufsize)) { + goto clear; + } + } + // on tiny post use memory + else { + if (!uwsgi_read_whole_body_in_mem(wsgi_req, wsgi_req->post_buffering_buf)) { + goto clear; + } + } } + UWSGI_GET_GIL } else { - // read to disk if post_cl > post_buffering - if (wsgi_req->post_cl >= (size_t) uwsgi.post_buffering) { - if (!uwsgi_read_whole_body(wsgi_req, wsgi_req->post_buffering_buf, uwsgi.post_buffering_bufsize)) { - goto clear; - } - } - // on tiny post use memory - else { - if (!uwsgi_read_whole_body_in_mem(wsgi_req, wsgi_req->post_buffering_buf)) { - goto clear; - } - } + wsgi_req->async_post = fdopen(wsgi_req->poll.fd, "r"); } - UWSGI_GET_GIL - } - else { - wsgi_req->async_post = fdopen(wsgi_req->poll.fd, "r"); } if (!up.pep3333_input) { @@ -488,9 +492,9 @@ int uwsgi_request_wsgi(struct wsgi_request *wsgi_req) { // LOCK THIS PART - wsgi_req->response_size += write(wsgi_req->poll.fd, wsgi_req->protocol, wsgi_req->protocol_len); - wsgi_req->response_size += write(wsgi_req->poll.fd, " 500 Internal Server Error\r\n", 28 ); - wsgi_req->response_size += write(wsgi_req->poll.fd, "Content-type: text/plain\r\n\r\n", 28 ); + wsgi_req->response_size += wsgi_req->socket_proto_write(wsgi_req, wsgi_req->protocol, wsgi_req->protocol_len); + wsgi_req->response_size += wsgi_req->socket_proto_write(wsgi_req, " 500 Internal Server Error\r\n", 28 ); + wsgi_req->response_size += wsgi_req->socket_proto_write(wsgi_req, "Content-type: text/plain\r\n\r\n", 28 ); wsgi_req->header_cnt = 1; /* diff --git a/plugins/python/wsgi_headers.c b/plugins/python/wsgi_headers.c index 7ba1c085..e9397922 100644 --- a/plugins/python/wsgi_headers.c +++ b/plugins/python/wsgi_headers.c @@ -176,7 +176,7 @@ PyObject *py_uwsgi_spit(PyObject * self, PyObject * args) { uh.pktsize += wsgi_req->hvec[j].iov_len; UWSGI_RELEASE_GIL - wsgi_req->headers_size = writev(wsgi_req->poll.fd, wsgi_req->hvec, j + 1); + wsgi_req->headers_size = wsgi_req->socket_proto_writev(wsgi_req, wsgi_req->hvec, j + 1); UWSGI_GET_GIL if (wsgi_req->headers_size < 0) { uwsgi_error("writev()"); diff --git a/plugins/python/wsgi_subhandler.c b/plugins/python/wsgi_subhandler.c index 818ed6c7..a95f0670 100644 --- a/plugins/python/wsgi_subhandler.c +++ b/plugins/python/wsgi_subhandler.c @@ -124,7 +124,7 @@ int uwsgi_response_subhandler_wsgi(struct wsgi_request *wsgi_req) { // return or yield ? if (PyString_Check((PyObject *)wsgi_req->async_result)) { - if ((wsize = write(wsgi_req->poll.fd, PyString_AsString(wsgi_req->async_result), PyString_Size(wsgi_req->async_result))) < 0) { + if ((wsize = wsgi_req->socket_proto_write(wsgi_req, PyString_AsString(wsgi_req->async_result), PyString_Size(wsgi_req->async_result))) < 0) { uwsgi_error("write()"); goto clear; } @@ -185,7 +185,7 @@ int uwsgi_response_subhandler_wsgi(struct wsgi_request *wsgi_req) { if (PyString_Check(pychunk)) { - if ((wsize = write(wsgi_req->poll.fd, PyString_AsString(pychunk), PyString_Size(pychunk))) < 0) { + if ((wsize = wsgi_req->socket_proto_write(wsgi_req, PyString_AsString(pychunk), PyString_Size(pychunk))) < 0) { uwsgi_error("write()"); Py_DECREF(pychunk); goto clear; @@ -216,7 +216,7 @@ clear: } if (wsgi_req->async_post && !wsgi_req->fd_closed) { fclose(wsgi_req->async_post); - if (!uwsgi.post_buffering || wsgi_req->post_cl <= (size_t) uwsgi.post_buffering) { + if ( (!uwsgi.post_buffering || wsgi_req->post_cl <= (size_t) uwsgi.post_buffering) && !wsgi_req->body_as_file) { wsgi_req->fd_closed = 1; } } diff --git a/plugins/rack/rack_plugin.c b/plugins/rack/rack_plugin.c index 850561b3..22b078d5 100644 --- a/plugins/rack/rack_plugin.c +++ b/plugins/rack/rack_plugin.c @@ -398,11 +398,10 @@ VALUE send_body(VALUE obj) { struct wsgi_request *wsgi_req = current_wsgi_req(); ssize_t len = 0; - int fd = wsgi_req->poll.fd; //uwsgi_log("sending body\n"); if (TYPE(obj) == T_STRING) { - len = write( fd, RSTRING_PTR(obj), RSTRING_LEN(obj)); + len = wsgi_req->socket_proto_write( wsgi_req, RSTRING_PTR(obj), RSTRING_LEN(obj)); } else { uwsgi_log("UNMANAGED BODY TYPE %d\n", TYPE(obj)); @@ -460,9 +459,9 @@ VALUE send_header(VALUE obj, VALUE headers) { //uwsgi_log("header: %.*s: %.*s\n", RSTRING_LEN(hkey), RSTRING_PTR(hkey), RSTRING_LEN(hval), RSTRING_PTR(hval)); - len = write( fd, RSTRING_PTR(hkey), RSTRING_LEN(hkey)); + len = wsgi_req->socket_proto_write_header( wsgi_req, RSTRING_PTR(hkey), RSTRING_LEN(hkey)); wsgi_req->headers_size += len; - len = write( fd, ": ", 2); + len = wsgi_req->socket_proto_write_header( wsgi_req, ": ", 2); wsgi_req->headers_size += len; char *header_value = RSTRING_PTR(hval); @@ -471,17 +470,17 @@ VALUE send_header(VALUE obj, VALUE headers) { char *header_value_splitted = memchr(header_value, '\n', header_value_len); if (!header_value_splitted) { - len = write( fd, header_value, header_value_len); + len = wsgi_req->socket_proto_write_header( wsgi_req, header_value, header_value_len); wsgi_req->headers_size += len; - len = write( fd, "\r\n", 2); + len = wsgi_req->socket_proto_write_header( wsgi_req, "\r\n", 2); wsgi_req->headers_size += len; wsgi_req->header_cnt++; } else { header_value_splitted[0] = 0; - len = write( fd, header_value, header_value_splitted-header_value); + len = wsgi_req->socket_proto_write_header( wsgi_req, header_value, header_value_splitted-header_value); wsgi_req->headers_size += len; - len = write( fd, "\r\n", 2); + len = wsgi_req->socket_proto_write_header( wsgi_req, "\r\n", 2); wsgi_req->headers_size += len; wsgi_req->header_cnt++; @@ -491,14 +490,14 @@ VALUE send_header(VALUE obj, VALUE headers) { while(header_value_len && (header_value_splitted = memchr(header_value, '\n', header_value_len))) { header_value_splitted[0] = 0; - len = write( fd, RSTRING_PTR(hkey), RSTRING_LEN(hkey)); + len = wsgi_req->socket_proto_write( wsgi_req, RSTRING_PTR(hkey), RSTRING_LEN(hkey)); wsgi_req->headers_size += len; - len = write( fd, ": ", 2); + len = wsgi_req->socket_proto_write( wsgi_req, ": ", 2); wsgi_req->headers_size += len; - len = write( fd, header_value, header_value_splitted-header_value); + len = wsgi_req->socket_proto_write( wsgi_req, header_value, header_value_splitted-header_value); wsgi_req->headers_size += len; - len = write( fd, "\r\n", 2); + len = wsgi_req->socket_proto_write( wsgi_req, "\r\n", 2); wsgi_req->headers_size += len; wsgi_req->header_cnt++; @@ -628,7 +627,7 @@ int uwsgi_rack_request(struct wsgi_request *wsgi_req) { wsgi_req->hvec[5].iov_len = 2 ; //RUBY_GVL_UNLOCK - if ( !(wsgi_req->headers_size = writev(wsgi_req->poll.fd, wsgi_req->hvec, 6)) ) { + if ( !(wsgi_req->headers_size = wsgi_req->socket_proto_writev_header(wsgi_req, wsgi_req->hvec, 6)) ) { uwsgi_error("writev()"); } //RUBY_GVL_LOCK @@ -639,7 +638,7 @@ int uwsgi_rack_request(struct wsgi_request *wsgi_req) { } //RUBY_GVL_UNLOCK - if (write(wsgi_req->poll.fd, "\r\n", 2) != 2) { + if (wsgi_req->socket_proto_write(wsgi_req, "\r\n", 2) != 2) { uwsgi_error("write()"); } //RUBY_GVL_LOCK diff --git a/plugins/rpc/rpc_plugin.c b/plugins/rpc/rpc_plugin.c index 65fda85a..0424dc25 100644 --- a/plugins/rpc/rpc_plugin.c +++ b/plugins/rpc/rpc_plugin.c @@ -29,10 +29,10 @@ int uwsgi_rpc_request(struct wsgi_request *wsgi_req) { wsgi_req->uh.pktsize = uwsgi_rpc(argv[0], argc-1, argv+1, wsgi_req->buffer); if (wsgi_req->uh.modifier2 == 0) { - wsgi_req->headers_size = write(wsgi_req->poll.fd, &wsgi_req->uh, 4); + wsgi_req->headers_size = wsgi_req->socket_proto_write_header(wsgi_req, (char *)&wsgi_req->uh, 4); } - wsgi_req->response_size = write(wsgi_req->poll.fd, wsgi_req->buffer, wsgi_req->uh.pktsize); + wsgi_req->response_size = wsgi_req->socket_proto_write(wsgi_req, wsgi_req->buffer, wsgi_req->uh.pktsize); wsgi_req->status = 0; return 0; diff --git a/proto/http.c b/proto/http.c new file mode 100644 index 00000000..890f9bd9 --- /dev/null +++ b/proto/http.c @@ -0,0 +1,328 @@ +/* async uwsgi protocol parser */ + +#include "../uwsgi.h" + +extern struct uwsgi_server uwsgi; + +static uint16_t http_add_uwsgi_header(struct wsgi_request *wsgi_req, char *hh, int hhlen) { + + char *buffer = wsgi_req->buffer+wsgi_req->uh.pktsize; + char *watermark = wsgi_req->buffer+uwsgi.buffer_size; + + int i; + int status = 0; + char *val = hh; + uint16_t keylen = 0, vallen = 0; + int prefix = 0; + char *ptr = buffer; + + for(i=0;ipost_cl = uwsgi_str_num(val, vallen); + } + + if (buffer+keylen+vallen+2+2 >= watermark) { + if (prefix) { + uwsgi_log("[WARNING] unable to add HTTP_%.*s=%.*s to uwsgi packet, consider increasing buffer size\n", keylen, hh, vallen, val); + } + else { + uwsgi_log("[WARNING] unable to add %.*s=%.*s to uwsgi packet, consider increasing buffer size\n", keylen, hh, vallen, val); + } + return 0; + } + + + *ptr++= (uint8_t) (keylen & 0xff); + *ptr++= (uint8_t) ((keylen >> 8) & 0xff); + + if (prefix) { + memcpy(ptr, "HTTP_", 5); + ptr+=5; + memcpy(ptr, hh, keylen-5); ptr+=(keylen-5); + } + else { + memcpy(ptr, hh, keylen); ptr+=keylen; + } + + *ptr++= (uint8_t) (vallen & 0xff); + *ptr++= (uint8_t) ((vallen >> 8) & 0xff); + memcpy(ptr, val, vallen); + +#ifdef UWSGI_DEBUG + uwsgi_log("add uwsgi var: %.*s = %.*s\n", keylen-(prefix*5), hh, vallen, val); +#endif + + return 2+keylen+2+vallen; +} + + +static uint16_t http_add_uwsgi_var(struct wsgi_request *wsgi_req, char *key, uint16_t keylen, char *val, uint16_t vallen) { + + + char *buffer = wsgi_req->buffer+wsgi_req->uh.pktsize; + char *watermark = wsgi_req->buffer+uwsgi.buffer_size; + char *ptr = buffer; + + if (buffer+keylen+vallen+2+2 >= watermark) { + uwsgi_log("[WARNING] unable to add %.*s=%.*s to uwsgi packet, consider increasing buffer size\n", keylen, key, vallen, val); + return 0; + } + + + *ptr++= (uint8_t) (keylen & 0xff); + *ptr++= (uint8_t) ((keylen >> 8) & 0xff); + memcpy(ptr, key, keylen); ptr+=keylen; + + *ptr++= (uint8_t) (vallen & 0xff); + *ptr++= (uint8_t) ((vallen >> 8) & 0xff); + memcpy(ptr, val, vallen); + +#ifdef UWSGI_DEBUG + uwsgi_log("add uwsgi var: %.*s = %.*s\n", keylen, key, vallen, val); +#endif + + return keylen+vallen+2+2; +} + +static int http_parse(struct wsgi_request *wsgi_req, char *watermark) { + + char *ptr = wsgi_req->proto_parser_buf; + char *base = ptr; + char *query_string = NULL; + + // REQUEST_METHOD + while(ptr < watermark) { + if (*ptr == ' ') { + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "REQUEST_METHOD", 14, base, ptr-base); + ptr++; + break; + } + ptr++; + } + + // REQUEST_URI / PATH_INFO / QUERY_STRING + base = ptr; + while(ptr < watermark) { + if (*ptr == '?' && !query_string) { + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, base, ptr-base); + query_string = ptr+1; + } + else if (*ptr == ' ') { + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "REQUEST_URI", 11, base, ptr-base); + if (!query_string) { + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "PATH_INFO", 9, base, ptr-base); + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "QUERY_STRING", 12, "", 0); + } + else { + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "QUERY_STRING", 12, query_string, ptr-query_string); + } + ptr++; + break; + } + ptr++; + } + + // SERVER_PROTOCOL + base = ptr; + while(ptr < watermark) { + if (*ptr == '\r') { + if (ptr + 1 >= watermark) return 0; + if (*(ptr+1) != '\n') return 0; + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "SERVER_PROTOCOL", 15, base, ptr-base); + ptr+=2; + break; + } + ptr++; + } + + // SCRIPT_NAME + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "SCRIPT_NAME", 11, "", 0); + + // SERVER_NAME + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "SERVER_NAME", 11, uwsgi.hostname, uwsgi.hostname_len); + + // SERVER_PORT + wsgi_req->uh.pktsize += http_add_uwsgi_var(wsgi_req, "SERVER_PORT", 11, "3031", 4); + + /* + // REMOTE_ADDR + if (inet_ntop(AF_INET, &h_session->ip_addr, h_session->ip, INET_ADDRSTRLEN)) { + h_session->uh.pktsize += http_add_uwsgi_var(h_session->iov, h_session->uss+c, h_session->uss+c+2, "REMOTE_ADDR", 11, h_session->ip, strlen(h_session->ip)); + } + else { + uwsgi_error("inet_ntop()"); + } + */ + + //HEADERS + base = ptr; + + while(ptr < watermark) { + if (*ptr == '\r') { + if (ptr + 1 >= watermark) return 0; + if (*(ptr+1) != '\n') return 0; + // multiline header ? + if (ptr+2 < watermark) { + if (*(ptr+2) == ' ' || *(ptr+2) == '\t') { + ptr+=2; + continue; + } + } + wsgi_req->uh.pktsize += http_add_uwsgi_header(wsgi_req, base, ptr-base); + ptr++; + base = ptr+1; + } + ptr++; + } + + return 0; + +} + + + +int uwsgi_proto_http_parser(struct wsgi_request *wsgi_req) { + + ssize_t len; + int j; + char *ptr; + int ret; + ssize_t remains; + // make this buffer configurable + char post_buf[8192]; + + // first round ? this memory area will be freed by async_loop + if (!wsgi_req->proto_parser_buf) { + wsgi_req->proto_parser_buf = uwsgi_malloc(uwsgi.buffer_size); + } + + if (wsgi_req->post_cl) { + remains = wsgi_req->post_cl-wsgi_req->proto_parser_pos; + if (remains > 0) { + remains = UMIN(remains, 8192); + len = read(wsgi_req->poll.fd, post_buf, remains); + if (len <= 0) { + uwsgi_error("read()"); + fclose(wsgi_req->async_post); + return -1; + } + + if (!fwrite(post_buf, len, 1, wsgi_req->async_post)) { + uwsgi_error("fwrite()"); + fclose(wsgi_req->async_post); + return -1; + } + wsgi_req->proto_parser_pos += len; + + if (wsgi_req->proto_parser_pos < wsgi_req->post_cl) + return UWSGI_AGAIN; + + } + rewind(wsgi_req->async_post); + return UWSGI_OK; + } + + len = read(wsgi_req->poll.fd, wsgi_req->proto_parser_buf+wsgi_req->proto_parser_pos, uwsgi.buffer_size-wsgi_req->proto_parser_pos); + if (len <= 0) { + free(wsgi_req->proto_parser_buf); + uwsgi_error("recv()"); + return -1; + } + + ptr = wsgi_req->proto_parser_buf+wsgi_req->proto_parser_pos; + + wsgi_req->proto_parser_pos+=len; + + for(j=0;jproto_parser_status == 0 || wsgi_req->proto_parser_status == 2)) { + wsgi_req->proto_parser_status++; + } + else if (*ptr == '\r') { + wsgi_req->proto_parser_status = 1; + } + else if (*ptr == '\n' && wsgi_req->proto_parser_status == 1) { + wsgi_req->proto_parser_status = 2; + } + else if (*ptr == '\n' && wsgi_req->proto_parser_status == 3) { + ptr++; + remains = len-(j+1); + ret = http_parse(wsgi_req, ptr); + //is there a Content_Length ? + if (wsgi_req->post_cl) { + wsgi_req->body_as_file = 1; + wsgi_req->async_post = tmpfile(); + if (!wsgi_req->async_post) { + free(wsgi_req->proto_parser_buf); + uwsgi_error("tmpfile()"); + return -1; + } + wsgi_req->proto_parser_pos = 0; + remains = UMIN((size_t)remains, wsgi_req->post_cl); + if (remains) { + if (!fwrite(ptr, remains, 1, wsgi_req->async_post)) { + free(wsgi_req->proto_parser_buf); + uwsgi_error("fwrite()"); + fclose(wsgi_req->async_post); + return -1; + } + wsgi_req->proto_parser_pos += remains; + if (wsgi_req->proto_parser_pos >= wsgi_req->post_cl) { + free(wsgi_req->proto_parser_buf); + rewind(wsgi_req->async_post); + return UWSGI_OK; + } + } + return UWSGI_AGAIN; + } + free(wsgi_req->proto_parser_buf); + return UWSGI_OK; + } + else { + wsgi_req->proto_parser_status = 0; + } + ptr++; + } + + return UWSGI_AGAIN; +} + +ssize_t uwsgi_proto_http_writev_header(struct wsgi_request *wsgi_req, struct iovec *iovec, size_t iov_len) { + return writev(wsgi_req->poll.fd, iovec, iov_len); +} + +ssize_t uwsgi_proto_http_writev(struct wsgi_request *wsgi_req, struct iovec *iovec, size_t iov_len) { + return writev(wsgi_req->poll.fd, iovec, iov_len); +} + +ssize_t uwsgi_proto_http_write(struct wsgi_request *wsgi_req, char *buf, size_t len) { + return write(wsgi_req->poll.fd, buf, len); +} + +ssize_t uwsgi_proto_http_write_header(struct wsgi_request *wsgi_req, char *buf, size_t len) { + return write(wsgi_req->poll.fd, buf, len); +} + diff --git a/proto/uwsgi.c b/proto/uwsgi.c new file mode 100644 index 00000000..24d7d93a --- /dev/null +++ b/proto/uwsgi.c @@ -0,0 +1,128 @@ +/* async uwsgi protocol parser */ + +#include "../uwsgi.h" + +extern struct uwsgi_server uwsgi; + +#define PROTO_STATUS_RECV_HDR 0 +#define PROTO_STATUS_RECV_VARS 1 + +int uwsgi_proto_uwsgi_parser(struct wsgi_request *wsgi_req) { + + uint8_t *hdr_buf = (uint8_t *) &wsgi_req->uh; + ssize_t len; + + struct iovec iov [1]; + struct cmsghdr *cmsg; + + if (wsgi_req->proto_parser_status == PROTO_STATUS_RECV_HDR) { + + + if (wsgi_req->proto_parser_pos > 0) { + len = read(wsgi_req->poll.fd, hdr_buf+wsgi_req->proto_parser_pos, 4-wsgi_req->proto_parser_pos); + } + else { + iov [0].iov_base = hdr_buf; + iov [0].iov_len = 4; + + wsgi_req->msg.msg_name = NULL; + wsgi_req->msg.msg_namelen = 0; + wsgi_req->msg.msg_iov = iov; + wsgi_req->msg.msg_iovlen = 1; + wsgi_req->msg.msg_control = &wsgi_req->msg_control; + wsgi_req->msg.msg_controllen = sizeof(wsgi_req->msg_control); + wsgi_req->msg.msg_flags = 0; + + len = recvmsg(wsgi_req->poll.fd, &wsgi_req->msg, 0); + } + + if (len <= 0) { + uwsgi_error("read()"); + return -1; + } + wsgi_req->proto_parser_pos+=len; + // header ready ? + if (wsgi_req->proto_parser_pos == 4) { + wsgi_req->proto_parser_status = PROTO_STATUS_RECV_VARS; + wsgi_req->proto_parser_pos = 0; +/* big endian ? */ +#ifdef __BIG_ENDIAN__ + wsgi_req->uh.pktsize = uwsgi_swap16(wsgi_req->uh.pktsize); +#endif + +#ifdef UWSGI_DEBUG + uwsgi_debug("uwsgi payload size: %d (0x%X) modifier1: %d modifier2: %d\n", wsgi_req->uh.pktsize, wsgi_req->uh.pktsize, wsgi_req->uh.modifier1, wsgi_req->uh.modifier2); +#endif + + /* check for max buffer size */ + if (wsgi_req->uh.pktsize > uwsgi.buffer_size) { + uwsgi_log( "invalid request block size: %d...skip\n", wsgi_req->uh.pktsize); + return -1; + } + + if (!wsgi_req->uh.pktsize) return UWSGI_OK; + + } + return UWSGI_AGAIN; + } + + else if (wsgi_req->proto_parser_status == PROTO_STATUS_RECV_VARS) { + len = read(wsgi_req->poll.fd, wsgi_req->buffer+wsgi_req->proto_parser_pos, wsgi_req->uh.pktsize-wsgi_req->proto_parser_pos); + if (len <= 0) { + uwsgi_error("read()"); + return -1; + } + wsgi_req->proto_parser_pos+=len; + + // body ready ? + if (wsgi_req->proto_parser_pos >= wsgi_req->uh.pktsize) { + + // older OSX versions make mess with CMSG_FIRSTHDR +#ifdef __APPLE__ + if (!wsgi_req->msg.msg_controllen) return UWSGI_OK; +#endif + + cmsg = CMSG_FIRSTHDR (&wsgi_req->msg); + while(cmsg != NULL) { + if (cmsg->cmsg_len == CMSG_LEN(sizeof(int)) && + cmsg->cmsg_level == SOL_SOCKET && + cmsg->cmsg_type && SCM_RIGHTS) { + + // upgrade connection to the new socket +#ifdef UWSGI_DEBUG + uwsgi_log("upgrading fd %d to ", wsgi_req->poll.fd); +#endif + close(wsgi_req->poll.fd); + memcpy(&wsgi_req->poll.fd, CMSG_DATA(cmsg), sizeof(int)); +#ifdef UWSGI_DEBUG + uwsgi_log("%d\n", wsgi_req->poll.fd); +#endif + } + cmsg = CMSG_NXTHDR (&wsgi_req->msg, cmsg); + } + + return UWSGI_OK; + } + return UWSGI_AGAIN; + } + + // never here + + return -1; +} + +ssize_t uwsgi_proto_uwsgi_writev_header(struct wsgi_request *wsgi_req, struct iovec *iovec, size_t iov_len) { + return writev(wsgi_req->poll.fd, iovec, iov_len); +} + +ssize_t uwsgi_proto_uwsgi_writev(struct wsgi_request *wsgi_req, struct iovec *iovec, size_t iov_len) { + return writev(wsgi_req->poll.fd, iovec, iov_len); +} + +ssize_t uwsgi_proto_uwsgi_write(struct wsgi_request *wsgi_req, char *buf, size_t len) { + return write(wsgi_req->poll.fd, buf, len); +} + +ssize_t uwsgi_proto_uwsgi_write_header(struct wsgi_request *wsgi_req, char *buf, size_t len) { + return write(wsgi_req->poll.fd, buf, len); +} diff --git a/protocol.c b/protocol.c index ae1dc6c1..845e60f3 100644 --- a/protocol.c +++ b/protocol.c @@ -339,141 +339,29 @@ ssize_t uwsgi_send_message(int fd, uint8_t modifier1, uint8_t modifier2, char *m } -int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header *uh, char *buffer) { - int rlen, i; - struct msghdr msg; - struct iovec iov [1]; - struct cmsghdr *cmsg; - union { - struct cmsghdr cmsg; - char control [CMSG_SPACE (sizeof (int))]; - } msg_control; +int uwsgi_parse_response(struct pollfd *upoll, int timeout, struct uwsgi_header *uh, char *buffer, int (*socket_proto)(struct wsgi_request *)) { + int rlen; + int status = UWSGI_AGAIN; if (!timeout) timeout = 1; - /* first 4 byte header */ - rlen = poll(upoll, 1, timeout * 1000); - if (rlen < 0) { - uwsgi_error("poll()"); - exit(1); - } - else if (rlen == 0) { - uwsgi_log( "timeout. skip request\n"); - close(upoll->fd); - return 0; - } - - iov [0].iov_base = uh; - iov [0].iov_len = 4; - - msg.msg_name = NULL; - msg.msg_namelen = 0; - msg.msg_iov = iov; - msg.msg_iovlen = 1; - msg.msg_control = &msg_control; - msg.msg_controllen = sizeof (msg_control); - msg.msg_flags = 0; - - //rlen = read(upoll->fd, uh, 4); - rlen = recvmsg(upoll->fd, &msg, 0); - if (rlen > 0 && rlen < 4) { - i = rlen; - while (i < 4) { - rlen = poll(upoll, 1, timeout * 1000); - if (rlen < 0) { - uwsgi_error("poll()"); - exit(1); - } - else if (rlen == 0) { - uwsgi_log( "timeout waiting for header. skip request.\n"); - close(upoll->fd); - break; - } - rlen = read(upoll->fd, (char *) (uh) + i, 4 - i); - if (rlen <= 0) { - uwsgi_log( "broken header. skip request.\n"); - close(upoll->fd); - break; - } - i += rlen; - } - if (i < 4) { - return 0; - } - } - else if (rlen <= 0) { - uwsgi_log( "invalid request header size: %d...skip\n", rlen); - close(upoll->fd); - return 0; - } - /* big endian ? */ -#ifdef __BIG_ENDIAN__ - uh->pktsize = uwsgi_swap16(uh->pktsize); -#endif - -#ifdef UWSGI_DEBUG - uwsgi_debug("uwsgi payload size: %d (0x%X) modifier1: %d modifier2: %d\n", uh->pktsize, uh->pktsize, uh->modifier1, uh->modifier2); -#endif - - /* check for max buffer size */ - if (uh->pktsize > uwsgi.buffer_size) { - uwsgi_log( "invalid request block size: %d...skip\n", uh->pktsize); - close(upoll->fd); - return 0; - } - - - //uwsgi_log("ready for reading %d bytes\n", wsgi_req.size); - - i = 0; - while (i < uh->pktsize) { + while(status == UWSGI_AGAIN) { rlen = poll(upoll, 1, timeout * 1000); if (rlen < 0) { uwsgi_error("poll()"); exit(1); } else if (rlen == 0) { - uwsgi_log( "timeout. skip request. (expecting %d bytes, got %d)\n", uh->pktsize, i); + uwsgi_log( "timeout waiting for header. skip request.\n"); close(upoll->fd); - break; + return 0; } - rlen = read(upoll->fd, buffer + i, uh->pktsize - i); - if (rlen <= 0) { - uwsgi_log( "broken vars. skip request.\n"); + status = socket_proto((struct wsgi_request *) uh); + if (status < 0) { close(upoll->fd); - break; + return 0; } - i += rlen; - } - - - if (i < uh->pktsize) { - return 0; - } - - // older OSX versions make mess with CMSG_FIRSTHDR -#ifdef __APPLE__ - if (!msg.msg_controllen) return 1; -#endif - - cmsg = CMSG_FIRSTHDR (&msg); - while(cmsg != NULL) { - if (cmsg->cmsg_len == CMSG_LEN(sizeof(int)) && - cmsg->cmsg_level == SOL_SOCKET && - cmsg->cmsg_type && SCM_RIGHTS) { - - // upgrade connection to the new socket -#ifdef UWSGI_DEBUG - uwsgi_log("upgrading fd %d to ", upoll->fd); -#endif - close(upoll->fd); - memcpy(&upoll->fd, CMSG_DATA(cmsg), sizeof(int)); -#ifdef UWSGI_DEBUG - uwsgi_log("%d\n", upoll->fd); -#endif - } - cmsg = CMSG_NXTHDR (&msg, cmsg); } return 1; @@ -705,18 +593,18 @@ int uwsgi_parse_vars(struct wsgi_request *wsgi_req) { ptrbuf += strsize; } else { - uwsgi_log("Invalid uwsgi request. skip.\n"); + uwsgi_log("invalid uwsgi request (current strsize: %d). skip.\n", strsize); return -1; } } else { - uwsgi_log("Invalid uwsgi request. skip.\n"); + uwsgi_log("invalid uwsgi request (current strsize: %d). skip.\n", strsize); return -1; } } } else { - uwsgi_log("Invalid uwsgi request. skip.\n"); + uwsgi_log("invalid uwsgi request (current strsize: %d). skip.\n", strsize); return -1; } } @@ -846,7 +734,7 @@ int uwsgi_ping_node(int node, struct wsgi_request *wsgi_req) { } uwsgi_poll.events = POLLIN; - if (!uwsgi_parse_response(&uwsgi_poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, wsgi_req->buffer)) { + if (!uwsgi_parse_response(&uwsgi_poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, wsgi_req->buffer, uwsgi_proto_uwsgi_parser)) { return -1; } @@ -1211,7 +1099,7 @@ char *uwsgi_simple_message_string(char *socket_name, uint8_t modifier1, uint8_t upoll.events = POLLIN; if (buffer) { - if (!uwsgi_parse_response(&upoll, timeout, &uh, buffer)) { + if (!uwsgi_parse_response(&upoll, timeout, &uh, buffer, uwsgi_proto_uwsgi_parser)) { close(fd); if (response_len) *response_len = 0; return NULL; diff --git a/setup.py b/setup.py index cb996737..8cd10df8 100644 --- a/setup.py +++ b/setup.py @@ -51,7 +51,7 @@ class uWSGIDistribution(Distribution): setup(name='uWSGI', - version='0.9.7.2', + version='0.9.8-dev', description='The uWSGI server', author='Unbit', author_email='info@unbit.it', diff --git a/utils.c b/utils.c index f2761191..e1d0f383 100644 --- a/utils.c +++ b/utils.c @@ -486,16 +486,23 @@ void wsgi_req_setup(struct wsgi_request *wsgi_req, int async_id) { } -int wsgi_req_simple_recv(struct wsgi_request *wsgi_req) { +int wsgi_req_async_recv(struct wsgi_request *wsgi_req, int socket_id) { UWSGI_SET_IN_REQUEST; gettimeofday(&wsgi_req->start_of_request, NULL); + if (event_queue_add_fd_read( uwsgi.async_queue, wsgi_req->poll.fd ) < 0) return -1; - if (!uwsgi_parse_response(&wsgi_req->poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, wsgi_req->buffer)) { - return -1; - } + async_add_timeout(wsgi_req, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT]); + + uwsgi.async_proto_fd_table[wsgi_req->poll.fd] = wsgi_req; + + wsgi_req->socket_proto = uwsgi.sockets[socket_id].proto; + wsgi_req->socket_proto_write = uwsgi.sockets[socket_id].proto_write; + wsgi_req->socket_proto_writev = uwsgi.sockets[socket_id].proto_writev; + wsgi_req->socket_proto_write_header = uwsgi.sockets[socket_id].proto_write_header; + wsgi_req->socket_proto_writev_header = uwsgi.sockets[socket_id].proto_writev_header; // enter harakiri mode if (uwsgi.shared->options[UWSGI_OPTION_HARAKIRI] > 0) { @@ -512,7 +519,7 @@ int wsgi_req_recv(struct wsgi_request *wsgi_req) { gettimeofday(&wsgi_req->start_of_request, NULL); - if (!uwsgi_parse_response(&wsgi_req->poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, wsgi_req->buffer)) { + if (!uwsgi_parse_response(&wsgi_req->poll, uwsgi.shared->options[UWSGI_OPTION_SOCKET_TIMEOUT], (struct uwsgi_header *) wsgi_req, wsgi_req->buffer, wsgi_req->socket_proto)) { return -1; } @@ -602,6 +609,12 @@ polling: fcntl(wsgi_req->poll.fd, F_SETFD, FD_CLOEXEC); } + // set socket protocol + wsgi_req->socket_proto = uwsgi.sockets[i].proto; + wsgi_req->socket_proto_write = uwsgi.sockets[i].proto_write; + wsgi_req->socket_proto_writev = uwsgi.sockets[i].proto_writev; + wsgi_req->socket_proto_write_header = uwsgi.sockets[i].proto_write_header; + wsgi_req->socket_proto_writev_header = uwsgi.sockets[i].proto_writev_header; return 0; } } @@ -1807,3 +1820,18 @@ int uwsgi_str4_num(char *str) { return num; } + +int uwsgi_str_num(char *str, int len) { + + int i; + int num = 0; + + int delta = pow(10, len); + + for(i=0;i 0) { @@ -1745,7 +1753,7 @@ int uwsgi_start(void *v_argv) { } - // put listening socket in non-blocking state + // put listening socket in non-blocking state and set the protocol for (i = 0; i < uwsgi.sockets_cnt; i++) { uwsgi.sockets[i].arg = fcntl(uwsgi.sockets[i].fd, F_GETFL, NULL); if (uwsgi.sockets[i].arg < 0) { @@ -1757,6 +1765,21 @@ int uwsgi_start(void *v_argv) { uwsgi_error("fcntl()"); exit(1); } + + if (!strcmp("http", uwsgi.protocol)) { + uwsgi.sockets[i].proto = uwsgi_proto_http_parser; + uwsgi.sockets[i].proto_write = uwsgi_proto_http_write; + uwsgi.sockets[i].proto_writev = uwsgi_proto_http_writev; + uwsgi.sockets[i].proto_write_header = uwsgi_proto_http_write_header; + uwsgi.sockets[i].proto_writev_header = uwsgi_proto_http_writev_header; + } + else { + uwsgi.sockets[i].proto = uwsgi_proto_uwsgi_parser; + uwsgi.sockets[i].proto_write = uwsgi_proto_uwsgi_write; + uwsgi.sockets[i].proto_writev = uwsgi_proto_uwsgi_writev; + uwsgi.sockets[i].proto_write_header = uwsgi_proto_uwsgi_write_header; + uwsgi.sockets[i].proto_writev_header = uwsgi_proto_uwsgi_writev_header; + } } } @@ -1801,7 +1824,7 @@ int uwsgi_start(void *v_argv) { } #endif - if (!uwsgi.sockets_cnt && !uwsgi.gateways_cnt && !uwsgi.no_server) { + if (!uwsgi.sockets_cnt && !uwsgi.gateways_cnt && !uwsgi.no_server && !uwsgi.udp_socket) { uwsgi_log("The -s/--socket option is missing and stdin is not a socket.\n"); exit(1); } @@ -2374,6 +2397,9 @@ end: uwsgi.threads = atoi(optarg); return 1; #endif + case LONG_ARGS_PROTOCOL: + uwsgi.protocol = optarg; + return 1; #ifdef UWSGI_ASYNC case LONG_ARGS_ASYNC: uwsgi.async = atoi(optarg); diff --git a/uwsgi.h b/uwsgi.h index d87340de..84f93e4b 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -2,7 +2,7 @@ /* indent -i8 -br -brs -brf -l0 -npsl -nip -npcs -npsl -di1 */ -#define UWSGI_VERSION "0.9.7.2" +#define UWSGI_VERSION "0.9.8-dev" #define UMAX16 65536 @@ -67,6 +67,7 @@ #include #include #include +#include #include #ifdef __sun__ @@ -402,6 +403,7 @@ struct uwsgi_opt { #define LONG_ARGS_EMPEROR_AMQP_VHOST 17096 #define LONG_ARGS_EMPEROR_AMQP_USERNAME 17097 #define LONG_ARGS_EMPEROR_AMQP_PASSWORD 17098 +#define LONG_ARGS_PROTOCOL 17099 #define UWSGI_OK 0 @@ -485,6 +487,8 @@ struct uwsgi_loop { void (*loop) (void); }; +struct wsgi_request; + struct uwsgi_socket { int fd; char *name; @@ -493,9 +497,14 @@ struct uwsgi_socket { int bound; int arg; void *ctx; + + int (*proto)(struct wsgi_request *); + ssize_t (*proto_write)(struct wsgi_request *, char *, size_t); + ssize_t (*proto_writev)(struct wsgi_request *, struct iovec *, size_t); + ssize_t (*proto_write_header)(struct wsgi_request *, char *, size_t); + ssize_t (*proto_writev_header)(struct wsgi_request *, struct iovec *, size_t); }; -struct wsgi_request; struct uwsgi_server; struct uwsgi_plugin { @@ -630,6 +639,12 @@ struct wsgi_request { //iovec struct iovec *hvec; + struct msghdr msg; + union { + struct cmsghdr cmsg; + char control [CMSG_SPACE (sizeof (int))]; + } msg_control; + struct timeval start_of_request; struct timeval end_of_request; @@ -728,9 +743,20 @@ struct wsgi_request { char *post_buffering_buf; uint64_t post_buffering_read; + int (*socket_proto)(struct wsgi_request *); + ssize_t (*socket_proto_write)(struct wsgi_request *, char *, size_t); + ssize_t (*socket_proto_writev)(struct wsgi_request *, struct iovec *, size_t); + ssize_t (*socket_proto_write_header)(struct wsgi_request *, char *, size_t); + ssize_t (*socket_proto_writev_header)(struct wsgi_request *, struct iovec *, size_t); + + int body_as_file; //for generic use off_t buf_pos; + uint64_t proto_parser_pos; + int proto_parser_status; + void *proto_parser_buf; + char *buffer; off_t frame_pos; @@ -849,6 +875,7 @@ struct uwsgi_server { char **async_post_buf; struct wsgi_request **async_waiting_fd_table; + struct wsgi_request **async_proto_fd_table; struct uwsgi_async_request *async_runqueue; struct uwsgi_async_request *async_runqueue_last; int async_runqueue_cnt; @@ -1010,6 +1037,8 @@ struct uwsgi_server { char *ns; char *ns_net; #endif + + char *protocol; int sockets_cnt; struct uwsgi_socket sockets[MAX_SOCKETS]; @@ -1351,7 +1380,7 @@ uint64_t uwsgi_swap64(uint64_t); ssize_t send_udp_message(uint8_t, char *, char *, uint16_t); #endif -int uwsgi_parse_response(struct pollfd *, int, struct uwsgi_header *, char *); +int uwsgi_parse_response(struct pollfd *, int, struct uwsgi_header *, char *, int (*)(struct wsgi_request *)); int uwsgi_parse_vars(struct wsgi_request *); int uwsgi_enqueue_message(char *, int, uint8_t, uint8_t, char *, int, int); @@ -1402,7 +1431,7 @@ void uwsgi_close_request(struct wsgi_request *); void wsgi_req_setup(struct wsgi_request *, int); int wsgi_req_recv(struct wsgi_request *); -int wsgi_req_simple_recv(struct wsgi_request *); +int wsgi_req_async_recv(struct wsgi_request *, int); int wsgi_req_accept(struct wsgi_request *); int wsgi_req_simple_accept(struct wsgi_request *, int); @@ -1760,3 +1789,21 @@ inline int uwsgi_starts_with(char *, int, char *, int); #ifdef __sun__ time_t timegm(struct tm *); #endif + + +int uwsgi_proto_uwsgi_parser(struct wsgi_request *); +int uwsgi_proto_http_parser(struct wsgi_request *); + +int uwsgi_str_num(char *, int); + + +ssize_t uwsgi_proto_uwsgi_writev_header(struct wsgi_request *, struct iovec *, size_t); +ssize_t uwsgi_proto_uwsgi_writev(struct wsgi_request *, struct iovec *, size_t); +ssize_t uwsgi_proto_uwsgi_write(struct wsgi_request *, char *, size_t); +ssize_t uwsgi_proto_uwsgi_write_header(struct wsgi_request *, char *, size_t); + +ssize_t uwsgi_proto_http_writev_header(struct wsgi_request *, struct iovec *, size_t); +ssize_t uwsgi_proto_http_writev(struct wsgi_request *, struct iovec *, size_t); +ssize_t uwsgi_proto_http_write(struct wsgi_request *, char *, size_t); +ssize_t uwsgi_proto_http_write_header(struct wsgi_request *, char *, size_t); + diff --git a/uwsgiconfig.py b/uwsgiconfig.py index 867c297d..d969874d 100644 --- a/uwsgiconfig.py +++ b/uwsgiconfig.py @@ -158,6 +158,9 @@ class uConf(object): self.config.read(filename) self.gcc_list = ['utils', 'protocol', 'socket', 'logging', 'master', 'emperor', 'plugins', 'lock', 'cache', 'queue', 'event', 'signal', 'rpc', 'gateway', 'loop', 'lib/rbtree', 'lib/amqp', 'rb_timers', 'uwsgi'] + # add protocols + self.gcc_list.append('proto/uwsgi') + self.gcc_list.append('proto/http') if uwsgi_os == 'Linux': self.gcc_list.append('lib/netlink') self.cflags = ['-O2', '-Wall', '-Werror', '-D_LARGEFILE_SOURCE', '-D_FILE_OFFSET_BITS=64'] + os.environ.get("CFLAGS", "").split()