diff --git a/plugins/corerouter/corerouter.c b/plugins/corerouter/corerouter.c index 92e39604..cad7fa5d 100644 --- a/plugins/corerouter/corerouter.c +++ b/plugins/corerouter/corerouter.c @@ -297,6 +297,12 @@ end: free(cr_session->buf_file_name); } + if (cr_session->write_queue) + free(cr_session->write_queue); + + if (cr_session->instance_write_queue) + free(cr_session->instance_write_queue); + // could be used to free additional resources if (cr_session->close) cr_session->close(ucr, cr_session); @@ -363,6 +369,11 @@ struct corerouter_session *corerouter_alloc_session(struct uwsgi_corerouter *ucr ucr->cr_table[new_connection]->timeout = cr_add_timeout(ucr, ucr->cr_table[new_connection]); ucr->cr_table[new_connection]->ugs = ugs; + ucr->cr_table[new_connection]->recv = uwsgi_cr_simple_recv; + ucr->cr_table[new_connection]->send = uwsgi_cr_simple_send; + ucr->cr_table[new_connection]->instance_recv = uwsgi_cr_simple_instance_recv; + ucr->cr_table[new_connection]->instance_send = uwsgi_cr_simple_instance_send; + ucr->alloc_session(ucr, ugs, ucr->cr_table[new_connection], cr_addr, cr_addr_len); event_queue_add_fd_read(ucr->queue, new_connection); @@ -534,15 +545,9 @@ void uwsgi_corerouter_loop(int id, void *data) { break; } - // set socket blocking mode, on non-linux platforms, clients get the server mode -#ifndef __linux__ - if (!ugs->nb) { - uwsgi_socket_b(new_connection); - } -#else - if (ugs->nb) { - uwsgi_socket_nb(new_connection); - } + // set socket in non-blocking mode, on non-linux platforms, clients get the server mode +#ifdef __linux__ + uwsgi_socket_nb(new_connection); #endif corerouter_alloc_session(ucr, ugs, new_connection, (struct sockaddr *) &cr_addr, cr_addr_len); diff --git a/plugins/corerouter/cr.h b/plugins/corerouter/cr.h index 069c7519..24337140 100644 --- a/plugins/corerouter/cr.h +++ b/plugins/corerouter/cr.h @@ -108,6 +108,10 @@ struct corerouter_session { int soopt; int timed_out; + // used for tracking required event + int fd_state; + int instance_fd_state; + struct uwsgi_rb_timer *timeout; int instance_failed; @@ -129,7 +133,21 @@ struct corerouter_session { int keepalive; + char *write_queue; + size_t write_queue_len; + off_t write_queue_pos; + + char *instance_write_queue; + size_t instance_write_queue_len; + off_t instance_write_queue_pos; + void (*close)(struct uwsgi_corerouter *, struct corerouter_session *); + + ssize_t (*recv)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); + ssize_t (*send)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); + + ssize_t (*instance_recv)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); + ssize_t (*instance_send)(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); }; void uwsgi_opt_corerouter(char *, char *, void *); @@ -163,3 +181,9 @@ int uwsgi_cr_map_use_to(struct uwsgi_corerouter *, struct corerouter_session *); int uwsgi_cr_map_use_static_nodes(struct uwsgi_corerouter *, struct corerouter_session *); int uwsgi_courerouter_has_has_backends(struct uwsgi_corerouter *); + +ssize_t uwsgi_cr_simple_recv(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); +ssize_t uwsgi_cr_simple_send(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); + +ssize_t uwsgi_cr_simple_instance_recv(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); +ssize_t uwsgi_cr_simple_instance_send(struct uwsgi_corerouter *, struct corerouter_session *, char *, size_t); diff --git a/plugins/corerouter/cr_common.c b/plugins/corerouter/cr_common.c index 9a6fc29b..11b826ad 100644 --- a/plugins/corerouter/cr_common.c +++ b/plugins/corerouter/cr_common.c @@ -10,6 +10,179 @@ extern struct uwsgi_server uwsgi; #include "cr.h" +ssize_t uwsgi_cr_simple_recv(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { + ssize_t ret = recv(cs->fd, buf, len, 0); + if (ret < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + errno = EINPROGRESS; + return -1; + } + uwsgi_error("recv()"); + } + return ret; +} + +ssize_t uwsgi_cr_simple_instance_recv(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { + ssize_t ret = recv(cs->instance_fd, buf, len, 0); + if (ret < 0) { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + errno = EINPROGRESS; + return -1; + } + uwsgi_error("recv()"); + } + return ret; +} + + +ssize_t uwsgi_cr_simple_send(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { + ssize_t ret; + char *tmp_buf; + if (cs->write_queue_len > 0) { + ret = send(cs->fd, cs->write_queue, cs->write_queue_len, 0); + if (ret > 0) { + cs->write_queue_len-=ret; + cs->write_queue_pos+=ret; + if (cs->write_queue_len == 0) { + free(cs->write_queue); + cs->write_queue = NULL; + cs->write_queue_pos = 0; + if (cs->fd_state) { + event_queue_fd_write_to_read(uc->queue, cs->fd); + cs->fd_state = 0; + } + goto next; + } + goto blocking; + } + else if (ret == 0) { + return 0; + } + else { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + goto blocking; + } + uwsgi_error("send()"); + return -1; + } + } + +next: + ret = send(cs->fd, buf, len, 0); + if (ret > 0) { + if ((size_t)ret == len) return len; + goto blocking; + } + + if (ret == 0) { + return 0; + } + + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + goto blocking; + } + uwsgi_error("send()"); + return -1; + +blocking: + // wait for write + if (!cs->fd_state) { + event_queue_fd_read_to_write(uc->queue, cs->fd); + cs->fd_state = 1; + } + // add new datas to the buffer + tmp_buf = malloc(cs->write_queue_len+len); + if (!tmp_buf) { + uwsgi_error("malloc()"); + return -1; + } + if (cs->write_queue_len>0) { + memcpy(tmp_buf, cs->write_queue+cs->write_queue_pos, cs->write_queue_len); + free(cs->write_queue); + } + memcpy(tmp_buf+cs->write_queue_len, buf, len); + cs->write_queue = tmp_buf; + cs->write_queue_pos = 0; + cs->write_queue_len+=len; + errno = EINPROGRESS; + return -1; +} + +ssize_t uwsgi_cr_simple_instance_send(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { + ssize_t ret; + char *tmp_buf; + if (cs->instance_write_queue_len > 0) { + ret = send(cs->instance_fd, cs->instance_write_queue, cs->instance_write_queue_len, 0); + if (ret > 0) { + cs->instance_write_queue_len-=ret; + cs->instance_write_queue_pos+=ret; + if (cs->instance_write_queue_len == 0) { + free(cs->instance_write_queue); + cs->instance_write_queue = NULL; + cs->instance_write_queue_pos = 0; + if (cs->instance_fd_state) { + event_queue_fd_write_to_read(uc->queue, cs->instance_fd); + cs->instance_fd_state = 0; + } + goto next; + } + goto blocking; + } + else if (ret == 0) { + return 0; + } + else { + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + goto blocking; + } + uwsgi_error("send()"); + return -1; + } + } + +next: + ret = send(cs->instance_fd, buf, len, 0); + if (ret > 0) { + if ((size_t)ret == len) return len; + goto blocking; + } + + if (ret == 0) { + return 0; + } + + if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINPROGRESS) { + goto blocking; + } + uwsgi_error("send()"); + return -1; + +blocking: + // wait for write + if (!cs->instance_fd_state) { + event_queue_fd_read_to_write(uc->queue, cs->instance_fd); + cs->instance_fd_state = 1; + } + // add new datas to the buffer + tmp_buf = malloc(cs->instance_write_queue_len+len); + if (!tmp_buf) { + uwsgi_error("malloc()"); + return -1; + } + if (cs->instance_write_queue_len>0) { + memcpy(tmp_buf, cs->instance_write_queue+cs->instance_write_queue_pos, cs->instance_write_queue_len); + free(cs->instance_write_queue); + } + memcpy(tmp_buf+cs->instance_write_queue_len, buf, len); + cs->instance_write_queue = tmp_buf; + cs->instance_write_queue_pos = 0; + cs->instance_write_queue_len+=len; + errno = EINPROGRESS; + return -1; +} + + + void uwsgi_corerouter_setup_sockets(struct uwsgi_corerouter *ucr) { struct uwsgi_gateway_socket *ugs = uwsgi.gateway_sockets; diff --git a/plugins/http/http.c b/plugins/http/http.c index f890208d..44d9d409 100644 --- a/plugins/http/http.c +++ b/plugins/http/http.c @@ -85,8 +85,6 @@ void uwsgi_opt_https(char *opt, char *value, void *cr) { // initialize ssl context ugs->ctx = uwsgi_ssl_new_server_context(uwsgi_concat3(ucr->short_name, "-", ugs->name),crt, key, ciphers, client_ca); - // the clients must be put in non-blocking mode - ugs->nb = 1; // set the ssl mode ugs->mode = UWSGI_HTTP_SSL; @@ -188,10 +186,6 @@ struct http_session { BIO *ssl_bio; char *ssl_cc; #endif - int fd_state; - - ssize_t (*recv)(struct http_session *, char *, size_t); - ssize_t (*send)(struct http_session *, char *, size_t); in_addr_t ip_addr; char ip[INET_ADDRSTRLEN]; @@ -451,16 +445,12 @@ int http_parse(struct http_session *h_session) { continue; } } + + // this is an hack with dumb/wrong/useless error checking if (uhttp.manage_expect) { if (!uwsgi_strncmp("Expect: 100-continue", 20, base, ptr - base)) { - if (send(h_session->crs.fd, protocol, protocol_len, 0) == (ssize_t) protocol_len) { - if (send(h_session->crs.fd, " 100 Continue\r\n\r\n", 17, 0) != 17) { - uwsgi_error("send()"); - } - } - else { - uwsgi_error("send()"); - } + if (h_session->crs.send(&uhttp.cr, &h_session->crs, protocol, protocol_len) == (ssize_t) protocol_len) + h_session->crs.send(&uhttp.cr, &h_session->crs, " 100 Continue\r\n\r\n", 17); } } h_session->uh.pktsize += http_add_uwsgi_header(h_session, h_session->iov, h_session->uss + c, h_session->uss + c + 2, base, ptr - base, &c); @@ -513,13 +503,13 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se goto choose_node; } - len = hs->recv(hs, hs->buffer + cs->h_pos, UMAX16 - cs->h_pos); + len = cs->recv(&uhttp.cr, cs, hs->buffer + cs->h_pos, UMAX16 - cs->h_pos); #ifdef UWSGI_EVENT_USE_PORT event_queue_add_fd_read(ucs->queue, cs->fd); #endif if (len <= 0) { // check for blocking operation on non-blocking socket - if (len < 0 && cs->ugs->nb && errno == EINPROGRESS) break; + if (len < 0 && errno == EINPROGRESS) break; corerouter_close_session(ucr, cs); break; } @@ -738,26 +728,25 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se // data from instance if (interesting_fd == cs->instance_fd) { - len = recv(cs->instance_fd, hs->buffer, UMAX16, 0); + len = cs->instance_recv(&uhttp.cr, cs, hs->buffer, UMAX16); #ifdef UWSGI_EVENT_USE_PORT event_queue_add_fd_read(uhttp_queue, cs->instance_fd); #endif - if (len <= 0) { - if (len < 0) { - uwsgi_error("recv()"); - } + if (len <= 0) { + if (len < 0 && errno == EINPROGRESS) break; /* Keep-Alive implementation. As soon as the backend close the connection, enable it. The server will start waiting for another request (for a maximum of --http-socket seconds timeout) To have a reliable implementation, we need to reset a bunch of values */ - else if (uhttp.keepalive) { + if (len == 0) { + if (uhttp.keepalive) { #ifdef UWSGI_DEBUG - uwsgi_log("Keep-Alive enabled\n"); + uwsgi_log("Keep-Alive enabled\n"); #endif - cs->keepalive = 1; - corerouter_close_session(ucr, cs); + cs->keepalive = 1; + corerouter_close_session(ucr, cs); cs->status = COREROUTER_STATUS_RECV_HDR; hs->ptr = hs->buffer; hs->rnrn = 0; @@ -771,30 +760,30 @@ To have a reliable implementation, we need to reset a bunch of values cs->instance_address_len = 0; cs->hostname_len = 0; break; - } + } #ifdef UWSGI_SSL - if (len == 0 && cs->ugs->mode == UWSGI_HTTP_SSL) { - int ssd_ret = SSL_shutdown(hs->ssl); - // it could fail or success, in both cases close the connection - if (ssd_ret != 0) { - corerouter_close_session(ucr, cs); + if (cs->ugs->mode == UWSGI_HTTP_SSL) { + int ssd_ret = SSL_shutdown(hs->ssl); + // it could fail or success, in both cases close the connection + if (ssd_ret != 0) { + corerouter_close_session(ucr, cs); + break; + } + cs->status = HTTP_SSL_STATUS_SHUTDOWN; + if (uwsgi_http_ssl_shutdown(hs, 1) != 0) { + corerouter_close_session(ucr, cs); + } break; } - cs->status = HTTP_SSL_STATUS_SHUTDOWN; - if (uwsgi_http_ssl_shutdown(hs, 1) != 0) { - corerouter_close_session(ucr, cs); - } - break; - } #endif + } corerouter_close_session(ucr, cs); break; } - len = hs->send(hs, hs->buffer, len); - + len = cs->send(&uhttp.cr, cs, hs->buffer, len); if (len <= 0) { - if (len < 0 && cs->ugs->nb && errno == EINPROGRESS) break; + if (len < 0 && errno == EINPROGRESS) break; // check for blocking operation non non-blocking socket corerouter_close_session(ucr, cs); break; @@ -809,13 +798,13 @@ To have a reliable implementation, we need to reset a bunch of values // body from client else if (interesting_fd == cs->fd) { - len = hs->recv(hs, bbuf, UMAX16); + len = cs->recv(&uhttp.cr, cs, bbuf, UMAX16); #ifdef UWSGI_EVENT_USE_PORT event_queue_add_fd_read(uhttp_queue, cs->fd); #endif if (len <= 0) { - // check for blocking operation non non-blocking socket - if (len < 0 && cs->ugs->nb && errno == EINPROGRESS) break; + // check for blocking operation on non-blocking socket + if (len < 0 && errno == EINPROGRESS) break; corerouter_close_session(ucr, cs); break; } @@ -833,11 +822,10 @@ To have a reliable implementation, we need to reset a bunch of values } raw: - len = send(cs->instance_fd, bbuf, len, 0); + len = cs->instance_send(&uhttp.cr, cs, bbuf, len); if (len <= 0) { - if (len < 0) - uwsgi_error("send()"); + if (len < 0 && errno == EINPROGRESS) break; corerouter_close_session(ucr, cs); break; } @@ -871,35 +859,6 @@ void http_setup() { uhttp.cr.short_name = uwsgi_str("http"); } -ssize_t uwsgi_http_simple_recv(struct http_session *hs, char *buf, size_t len) { - ssize_t ret = recv(hs->crs.fd, buf, len, 0); - if (ret < 0) { - uwsgi_error("recv()"); - } - return ret; -} - -ssize_t uwsgi_http_simple_send(struct http_session *hs, char *buf, size_t len) { - size_t remains = len; - char *ptr = buf; - while(remains > 0) { - ssize_t ret = send(hs->crs.fd, ptr, remains, 0); - if (ret > 0) { - remains -= ret; - ptr+=ret; - } - else if (ret == 0) { - return -1; - } - // error - else { - uwsgi_error("send()"); - return -1; - } - } - return len; -} - #ifdef UWSGI_SSL int uwsgi_http_ssl_shutdown(struct http_session *hs, int state) { int ret = 0; @@ -909,27 +868,28 @@ int uwsgi_http_ssl_shutdown(struct http_session *hs, int state) { if (ret == 1) return 1; int err = SSL_get_error(hs->ssl, ret); if (err == SSL_ERROR_WANT_READ) { - if (hs->fd_state) { + if (hs->crs.fd_state) { event_queue_fd_write_to_read(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 0; + hs->crs.fd_state = 0; } return 0; } else if (err == SSL_ERROR_WANT_WRITE) { - if (!hs->fd_state) { + if (!hs->crs.fd_state) { event_queue_fd_read_to_write(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 1; + hs->crs.fd_state = 1; } return 0; } return -1; } -ssize_t uwsgi_http_ssl_recv(struct http_session *hs, char *buf, size_t len) { +ssize_t uwsgi_http_ssl_recv(struct uwsgi_corerouter *cr, struct corerouter_session *cs, char *buf, size_t len) { + struct http_session *hs = (struct http_session *) cs; int ret = SSL_read(hs->ssl, buf, len); if (ret > 0) { - if (hs->fd_state) { - event_queue_fd_write_to_read(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 0; + if (hs->crs.fd_state) { + event_queue_fd_write_to_read(cr->queue, hs->crs.fd); + hs->crs.fd_state = 0; } return ret; } @@ -937,18 +897,18 @@ ssize_t uwsgi_http_ssl_recv(struct http_session *hs, char *buf, size_t len) { int err = SSL_get_error(hs->ssl, ret); if (err == SSL_ERROR_WANT_READ) { - if (hs->fd_state) { - event_queue_fd_write_to_read(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 0; + if (hs->crs.fd_state) { + event_queue_fd_write_to_read(cr->queue, hs->crs.fd); + hs->crs.fd_state = 0; } errno = EINPROGRESS; return -1; } else if (err == SSL_ERROR_WANT_WRITE) { - if (!hs->fd_state) { - event_queue_fd_read_to_write(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 1; + if (!hs->crs.fd_state) { + event_queue_fd_read_to_write(cr->queue, hs->crs.fd); + hs->crs.fd_state = 1; } errno = EINPROGRESS; return -1; @@ -966,29 +926,30 @@ ssize_t uwsgi_http_ssl_recv(struct http_session *hs, char *buf, size_t len) { } -ssize_t uwsgi_http_ssl_send(struct http_session *hs, char *buf, size_t len) { +ssize_t uwsgi_http_ssl_send(struct uwsgi_corerouter *cr, struct corerouter_session *cs, char *buf, size_t len) { + struct http_session *hs = (struct http_session *) cs; int ret = SSL_write(hs->ssl, buf, len); if (ret > 0) { - if (hs->fd_state) { - event_queue_fd_write_to_read(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 0; + if (hs->crs.fd_state) { + event_queue_fd_write_to_read(cr->queue, hs->crs.fd); + hs->crs.fd_state = 0; } return ret; } if (ret == 0) return 0; int err = SSL_get_error(hs->ssl, ret); if (err == SSL_ERROR_WANT_READ) { - if (hs->fd_state) { - event_queue_fd_write_to_read(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 0; + if (hs->crs.fd_state) { + event_queue_fd_write_to_read(cr->queue, hs->crs.fd); + hs->crs.fd_state = 0; } errno = EINPROGRESS; return -1; } else if (err == SSL_ERROR_WANT_WRITE) { - if (!hs->fd_state) { - event_queue_fd_read_to_write(uhttp.cr.queue, hs->crs.fd); - hs->fd_state = 1; + if (!hs->crs.fd_state) { + event_queue_fd_read_to_write(cr->queue, hs->crs.fd); + cs->fd_state = 1; } errno = EINPROGRESS; return -1; @@ -1039,8 +1000,6 @@ void http_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socke if (sa && sa->sa_family == AF_INET) { hs->ip_addr = ((struct sockaddr_in *) sa)->sin_addr.s_addr; } - hs->recv = uwsgi_http_simple_recv; - hs->send = uwsgi_http_simple_send; if (ugs) { hs->port = ugs->port; hs->port_len = ugs->port_len; @@ -1049,8 +1008,8 @@ void http_alloc_session(struct uwsgi_corerouter *ucr, struct uwsgi_gateway_socke hs->ssl = SSL_new(ugs->ctx); SSL_set_fd(hs->ssl, cs->fd); SSL_set_accept_state(hs->ssl); - hs->recv = uwsgi_http_ssl_recv; - hs->send = uwsgi_http_ssl_send; + cs->recv = uwsgi_http_ssl_recv; + cs->send = uwsgi_http_ssl_send; cs->close = uwsgi_ssl_close; } #endif diff --git a/uwsgi.h b/uwsgi.h index 69ffb5c7..cafb7330 100644 --- a/uwsgi.h +++ b/uwsgi.h @@ -414,7 +414,6 @@ struct uwsgi_gateway_socket { // this requires UDP int subscription; int shared; - int nb; char *owner; struct uwsgi_gateway *gateway; diff --git a/welcome.py b/welcome.py index c7b98e6c..79e675d7 100644 --- a/welcome.py +++ b/welcome.py @@ -44,7 +44,8 @@ def xsendfile(e, sr): def serve_logo(e, sr): # use raw facilities uwsgi.send("%s 200 OK\r\nContent-Type: image/png\r\n\r\n" % e['SERVER_PROTOCOL']) - return uwsgi.sendfile('logo_uWSGI.png') + uwsgi.sendfile('logo_uWSGI.png') + return '' def serve_options(e, sr): sr('200 OK', [('Content-Type', 'text/html')])