diff --git a/core/writer.c b/core/writer.c new file mode 100644 index 00000000..3ffedb99 --- /dev/null +++ b/core/writer.c @@ -0,0 +1,67 @@ +// status could be NNN or NNN message +int uwsgi_response_prepare_headers(struct wsgi_request *wsgi_req, char *status, uint16_t status_len) { +} + +//each protocol has its header generator +int uwsgi_response_add_header(struct wsgi_request *wsgi_req, char *key, uint16_t key_len, char *value, uint16_t value_len) { + + if (wsgi_req->headers_sent) return -1; + + if (!wsgi_req->headers) { + wsgi_req->headers = uwsgi_buffer_new(uwsgi.page_size); + wsgi_req->headers->limit = UMAX16; + } + + struct uwsgi_buffer *hh = wsgi_req->socket->proto_add_header(wsgi_req, key, key_len, value, value_len); + if (!hh) return -1; + if (uwsgi_buffer_append(wsgi_req->headers, hh->buf, hh->pos)) goto error; + uwsgi_buffer_destroy(hh); + return 0; +error: + uwsgi_buffer_destroy(hh); + return -1; +} + +// send the headers in a non blocking-way +int uwsgi_response_commit_headers(struct wsgi_request *wsgi_req) { + struct uwsgi_buffer *ub = wsgi_req->headers; + if (!ub) return UWSGI_OK; + + ssize_t remains = ub->pos - wsgi_req->headers_write_pos; + if (!remains) return UWSGI_OK; + + // this is special, as it returns -1 on error and OK/AGAIN otherwise + int ret = wsgi_req->socket->proto_fix_headers(wsgi_req); + if (len < 0) { + wsgi_req->write_errors++; + } + if (len == remains) { + wsgi_req->headers_sent = 1; + return UWSGI_OK; + } +} + +int uwsgi_response_write_body_do(struct wsgi_request *wsgi_req, char *buf, size_t len) { + wsgi_req->response_write_body_pos = 0; + + for(;;) { + int ret = wsgi_req->socket->proto_write(wsgi_req, buf, len); + if (ret < 0) { + if (!uwsgi.ignore_write_errors) { + uwsgi_error("uwsgi_response_write_body_do()"); + } + return -1; + } + if (ret == UWSGI_OK) { + break; + } + ret = uwsgi.wait_write_hook(wsgi_req); + if (ret < 0) return -1; + // callback based hook... + if (ret == UWSGI_AGAIN) return UWSGI_AGAIN; + } + + wsgi_req->response_size += wsgi_req->response_write_body_pos; + + return UWSGI_OK; +} diff --git a/plugins/psgi/psgi_plugin.c b/plugins/psgi/psgi_plugin.c index 28418bff..0e863f94 100644 --- a/plugins/psgi/psgi_plugin.c +++ b/plugins/psgi/psgi_plugin.c @@ -23,8 +23,6 @@ struct uwsgi_option uwsgi_perl_options[] = { }; -extern struct http_status_codes hsc[]; - SV *uwsgi_perl_obj_new(char *class, size_t class_len) { SV *newobj; @@ -377,11 +375,6 @@ int uwsgi_perl_init(){ PERL_SET_CONTEXT(uperl.main[0]); - // filling http status codes - for (http_sc = hsc; http_sc->message != NULL; http_sc++) { - http_sc->message_size = strlen(http_sc->message); - } - #ifdef PERL_VERSION_STRING uwsgi_log_initial("initialized Perl %s main interpreter at %p\n", PERL_VERSION_STRING, uperl.main[0]); #else diff --git a/plugins/psgi/psgi_response.c b/plugins/psgi/psgi_response.c index f292a256..1beca059 100644 --- a/plugins/psgi/psgi_response.c +++ b/plugins/psgi/psgi_response.c @@ -67,36 +67,8 @@ int psgi_response(struct wsgi_request *wsgi_req, AV *response) { status_code = av_fetch(response, 0, 0); if (!status_code) { uwsgi_log("invalid PSGI status code\n"); return UWSGI_OK;} - wsgi_req->hvec[0].iov_base = "HTTP/1.1 "; - wsgi_req->hvec[0].iov_len = 9; - - wsgi_req->hvec[1].iov_base = SvPV(*status_code, hlen); - - wsgi_req->hvec[1].iov_len = 3; - - wsgi_req->status = atoi(wsgi_req->hvec[1].iov_base); - - wsgi_req->hvec[2].iov_base = " "; - wsgi_req->hvec[2].iov_len = 1; - - wsgi_req->hvec[3].iov_len = 0; - - // get the status code - for (http_sc = hsc; http_sc->message != NULL; http_sc++) { - if (!strncmp(http_sc->key, wsgi_req->hvec[1].iov_base, 3)) { - wsgi_req->hvec[3].iov_base = (char *) http_sc->message; - wsgi_req->hvec[3].iov_len = http_sc->message_size; - break; - } - } - - if (wsgi_req->hvec[3].iov_len == 0) { - wsgi_req->hvec[3].iov_base = "Unknown"; - wsgi_req->hvec[3].iov_len = 7; - } - - wsgi_req->hvec[4].iov_base = "\r\n"; - wsgi_req->hvec[4].iov_len = 2; + char *status_str = SvPV(*status_code, hlen); + if (uwsgi_response_prepare_headers(wsgi_req, status_str, hlen)) return UWSGI_OK; hitem = av_fetch(response, 1, 0); if (!hitem) { uwsgi_log("invalid PSGI headers\n"); return UWSGI_OK;} @@ -109,65 +81,14 @@ int psgi_response(struct wsgi_request *wsgi_req, AV *response) { // put them in hvec for(i=0; i<=av_len(headers); i++) { - if (wsgi_req->header_cnt+1 > uwsgi.max_vars) { - uwsgi_log("no more space in iovec. consider increasing max-vars...\n"); - break; - } vi = (i*2)+base; hitem = av_fetch(headers,i,0); chitem = SvPV(*hitem, hlen); - wsgi_req->hvec[vi].iov_base = chitem; wsgi_req->hvec[vi].iov_len = hlen; - - wsgi_req->hvec[vi+1].iov_base = ": "; wsgi_req->hvec[vi+1].iov_len = 2; - - hitem = av_fetch(headers,i+1,0); - chitem = SvPV(*hitem, hlen); - wsgi_req->hvec[vi+2].iov_base = chitem; wsgi_req->hvec[vi+2].iov_len = hlen; - - wsgi_req->hvec[vi+3].iov_base = "\r\n"; wsgi_req->hvec[vi+3].iov_len = 2; - - wsgi_req->header_cnt++; - - i++; + hitem2 = av_fetch(headers,i+1,0); + chitem2 = SvPV(*hitem, hlen); + if (uwsgi_response_add_header(wsgi_req, chitem, hlen, chitem2, hlen)) return UWSGI_OK; } - int j = (i*2)+base; - struct uwsgi_string_list *ah = uwsgi.additional_headers; - while(ah) { - if (wsgi_req->header_cnt+1 > uwsgi.max_vars) { - uwsgi_log("no more space in iovec. consider increasing max-vars...\n"); - break; - } - wsgi_req->header_cnt++; - wsgi_req->hvec[j].iov_base = ah->value; - wsgi_req->hvec[j].iov_len = ah->len; - j++; - wsgi_req->hvec[j].iov_base = "\r\n"; - wsgi_req->hvec[j].iov_len = 2; - j++; - ah = ah->next; - } - - ah = wsgi_req->additional_headers; - while(ah) { - if (wsgi_req->header_cnt+1 > uwsgi.max_vars) { - uwsgi_log("no more space in iovec. consider increasing max-vars...\n"); - break; - } - wsgi_req->header_cnt++; - wsgi_req->hvec[j].iov_base = ah->value; - wsgi_req->hvec[j].iov_len = ah->len; - j++; - wsgi_req->hvec[j].iov_base = "\r\n"; - wsgi_req->hvec[j].iov_len = 2; - j++; - ah = ah->next; - } - - wsgi_req->hvec[j].iov_base = "\r\n"; wsgi_req->hvec[j].iov_len = 2; - - wsgi_req->headers_size += wsgi_req->socket->proto_writev_header(wsgi_req, wsgi_req->hvec, j+1); - hitem = av_fetch(response, 2, 0); if (!hitem) { diff --git a/plugins/python/wsgi_subhandler.c b/plugins/python/wsgi_subhandler.c index 7c486873..40b414ae 100644 --- a/plugins/python/wsgi_subhandler.c +++ b/plugins/python/wsgi_subhandler.c @@ -160,13 +160,9 @@ int uwsgi_response_subhandler_wsgi(struct wsgi_request *wsgi_req) { // return or yield ? if (PyString_Check((PyObject *)wsgi_req->async_result)) { + char *content = PyString_AsString((PyObject *)wsgi_req->async_result); size_t content_len = PyString_Size((PyObject *)wsgi_req->async_result); - if (content_len > 0 && !wsgi_req->headers_sent) { - if (uwsgi_python_do_send_headers(wsgi_req)) { - goto clear; - } - } - up.hook_write_string(wsgi_req, (PyObject *) wsgi_req->async_result); + uwsgi.nb_write_hook(wsgi_req, content, content_len); uwsgi_py_check_write_errors { uwsgi_py_write_exception(wsgi_req); } diff --git a/proto/fastcgi.c b/proto/fastcgi.c index b6a98d47..4e6feeb0 100644 --- a/proto/fastcgi.c +++ b/proto/fastcgi.c @@ -170,15 +170,12 @@ ssize_t uwsgi_proto_fastcgi_writev(struct wsgi_request * wsgi_req, struct iovec return uwsgi_proto_fastcgi_writev_header(wsgi_req, iovec, iov_len); } -ssize_t uwsgi_proto_fastcgi_write(struct wsgi_request * wsgi_req, char *buf, size_t len) { +int uwsgi_proto_fastcgi_write(struct wsgi_request * wsgi_req, char *buf, size_t len, ssize_t *written) { struct fcgi_record fr; - ssize_t rlen; - size_t chunk_len; - char *ptr = buf; - // in fastcgi we need to not send 0 size frames + // in fastcgi we need to not send 0 sized frames if (!len) - return 0; + return UWSGI_OK; fr.version = 1; fr.type = 6; @@ -188,9 +185,27 @@ ssize_t uwsgi_proto_fastcgi_write(struct wsgi_request * wsgi_req, char *buf, siz fr.reserved = 0; fr.cl = htons(len); - // split response in 64k chunks... + // still trying to send the fcgi header ? + if (*written < 8) { + char *ptr = (char *) &fr; + ssize_t wlen = write(wsgi_req->poll.fd, ptr + *written, 8 - *written); + if (wlen > 0) { + *written += wlen; + if (*written == 8) { + goto body; + } + return UWSGI_AGAIN; + } + if (wlen < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + return UWSGI_AGAIN; + } + } + return -1; + } - if (len <= 65535) { + // split response in 64k chunks... +body: rlen = write(wsgi_req->poll.fd, &fr, 8); if (rlen <= 0) { if (!uwsgi.ignore_write_errors) { @@ -237,10 +252,6 @@ ssize_t uwsgi_proto_fastcgi_write(struct wsgi_request * wsgi_req, char *buf, siz } -ssize_t uwsgi_proto_fastcgi_write_header(struct wsgi_request * wsgi_req, char *buf, size_t len) { - return uwsgi_proto_fastcgi_write(wsgi_req, buf, len); -} - void uwsgi_proto_fastcgi_close(struct wsgi_request *wsgi_req) { if (write(wsgi_req->poll.fd, FCGI_END_REQUEST, 24) <= 0) { diff --git a/proto/uwsgi.c b/proto/uwsgi.c index 129487c5..17047baf 100644 --- a/proto/uwsgi.c +++ b/proto/uwsgi.c @@ -116,44 +116,16 @@ int uwsgi_proto_uwsgi_parser(struct wsgi_request *wsgi_req) { return -1; } -ssize_t uwsgi_proto_uwsgi_writev_header(struct wsgi_request *wsgi_req, struct iovec * iovec, size_t iov_len) { - if (iov_len == 0) return 0; - ssize_t wlen = writev(wsgi_req->poll.fd, iovec, iov_len); - if (wlen < 0) { - if (!uwsgi.ignore_write_errors) { - uwsgi_req_error("writev()"); - } - wsgi_req->write_errors++; - return 0; +size_t uwsgi_proto_uwsgi_write(struct wsgi_request * wsgi_req, char *buf, size_t len) { + ssize_t wlen = write(wsgi_req->poll.fd, ptr, len); + if (wlen > 0) { + return wlen; } - return wlen; -} - -ssize_t uwsgi_proto_uwsgi_writev(struct wsgi_request *wsgi_req, struct iovec * iovec, size_t iov_len) { - return uwsgi_proto_uwsgi_writev_header(wsgi_req, iovec, iov_len); -} - -ssize_t uwsgi_proto_uwsgi_write(struct wsgi_request * wsgi_req, char *buf, size_t len) { - ssize_t wlen; - char *ptr = buf; - if (len == 0) return 0; - - while(len > 0) { - wlen = write(wsgi_req->poll.fd, ptr, len); - if (wlen <= 0) { - if (!uwsgi.ignore_write_errors) { - uwsgi_req_error("write()"); - } - wsgi_req->write_errors++; - return ptr-buf; - } - ptr+=wlen; - len -= wlen; + if (wlen == 0) { + return -1; } - - return ptr-buf; -} - -ssize_t uwsgi_proto_uwsgi_write_header(struct wsgi_request *wsgi_req, char *buf, size_t len) { - return uwsgi_proto_uwsgi_write(wsgi_req, buf, len); + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + return -2; + } + return -1; } diff --git a/uwsgi.h b/uwsgi.h index 3afc1e0b..a619e4dc 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -2073,6 +2073,7 @@ struct uwsgi_server { struct uwsgi_buffer *(*channel_recv_hook)(struct wsgi_request *, int, struct uwsgi_buffer *, int); ssize_t (*buffer_write_hook)(struct wsgi_request *, struct uwsgi_buffer *); + int (*nb_write_hook)(struct wsgi_request *, char *, size_t); }; diff --git a/uwsgiconfig.py b/uwsgiconfig.py index e8e7c971..c06c1ce3 100644 --- a/uwsgiconfig.py +++ b/uwsgiconfig.py @@ -454,7 +454,7 @@ class uConf(object): self.gcc_list = ['core/utils', 'core/protocol', 'core/socket', 'core/logging', 'core/master', 'core/master_utils', 'core/emperor', 'core/notify', 'core/mule', 'core/subscription', 'core/stats', 'core/sendfile', 'core/offload', 'core/io', 'core/static', 'core/websockets', 'core/channels', - 'core/setup_utils', 'core/clock', 'core/init', 'core/buffer', + 'core/setup_utils', 'core/clock', 'core/init', 'core/buffer', 'core/writer', 'core/plugins', 'core/lock', 'core/cache', 'core/daemons', 'core/queue', 'core/event', 'core/signal', 'core/cluster', 'core/rpc', 'core/gateway', 'core/loop', 'lib/rbtree', 'core/rb_timers', 'core/uwsgi']