added support for better/true non-blocking write in http router

This commit is contained in:
roberto@quantal64
2012-08-10 11:12:59 +02:00
parent 9fabfd77ee
commit 4322c3b538
6 changed files with 275 additions and 114 deletions
+14 -9
View File
@@ -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);
+24
View File
@@ -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);
+173
View File
@@ -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;
+62 -103
View File
@@ -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
-1
View File
@@ -414,7 +414,6 @@ struct uwsgi_gateway_socket {
// this requires UDP
int subscription;
int shared;
int nb;
char *owner;
struct uwsgi_gateway *gateway;
+2 -1
View File
@@ -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')])