From bab89d11524c2ee35eaffcf087a3c81a1f48a684 Mon Sep 17 00:00:00 2001 From: "roberto@quantal64" Date: Fri, 10 Aug 2012 12:49:23 +0200 Subject: [PATCH] another fix for write queue --- plugins/corerouter/cr.h | 3 +-- plugins/corerouter/cr_common.c | 48 +++++++++++++++++++++++----------- plugins/http/http.c | 33 ++++++++++++++++++++++- 3 files changed, 66 insertions(+), 18 deletions(-) diff --git a/plugins/corerouter/cr.h b/plugins/corerouter/cr.h index 24337140..f4548242 100644 --- a/plugins/corerouter/cr.h +++ b/plugins/corerouter/cr.h @@ -135,11 +135,10 @@ struct corerouter_session { char *write_queue; size_t write_queue_len; - off_t write_queue_pos; + int write_queue_close; char *instance_write_queue; size_t instance_write_queue_len; - off_t instance_write_queue_pos; void (*close)(struct uwsgi_corerouter *, struct corerouter_session *); diff --git a/plugins/corerouter/cr_common.c b/plugins/corerouter/cr_common.c index 11b826ad..c2943a2d 100644 --- a/plugins/corerouter/cr_common.c +++ b/plugins/corerouter/cr_common.c @@ -36,17 +36,21 @@ ssize_t uwsgi_cr_simple_instance_recv(struct uwsgi_corerouter *uc, struct corero ssize_t uwsgi_cr_simple_send(struct uwsgi_corerouter *uc, struct corerouter_session *cs, char *buf, size_t len) { - ssize_t ret; + ssize_t ret = -1; char *tmp_buf; + off_t pos = 0; + + size_t partial_len = len; + off_t partial_pos = 0; + 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; + 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; @@ -68,9 +72,13 @@ ssize_t uwsgi_cr_simple_send(struct uwsgi_corerouter *uc, struct corerouter_sess } next: + if (len == 0) goto end; + ret = send(cs->fd, buf, len, 0); if (ret > 0) { if ((size_t)ret == len) return len; + partial_len-=ret; + partial_pos = ret; goto blocking; } @@ -84,6 +92,12 @@ next: uwsgi_error("send()"); return -1; +end: + if (cs->write_queue_close) { + return 0; + } + return ret; + blocking: // wait for write if (!cs->fd_state) { @@ -91,19 +105,18 @@ blocking: cs->fd_state = 1; } // add new datas to the buffer - tmp_buf = malloc(cs->write_queue_len+len); + tmp_buf = malloc(cs->write_queue_len+partial_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); + memcpy(tmp_buf, cs->write_queue+pos, cs->write_queue_len); free(cs->write_queue); } - memcpy(tmp_buf+cs->write_queue_len, buf, len); + memcpy(tmp_buf+cs->write_queue_len, buf+partial_pos, partial_len); cs->write_queue = tmp_buf; - cs->write_queue_pos = 0; - cs->write_queue_len+=len; + cs->write_queue_len+=partial_len; errno = EINPROGRESS; return -1; } @@ -111,15 +124,19 @@ blocking: 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; + off_t pos = 0; + + size_t partial_len = len; + off_t partial_pos = 0; + 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; + 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; @@ -144,6 +161,8 @@ next: ret = send(cs->instance_fd, buf, len, 0); if (ret > 0) { if ((size_t)ret == len) return len; + partial_len-=ret; + partial_pos = ret; goto blocking; } @@ -164,19 +183,18 @@ blocking: cs->instance_fd_state = 1; } // add new datas to the buffer - tmp_buf = malloc(cs->instance_write_queue_len+len); + tmp_buf = malloc(cs->instance_write_queue_len+partial_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); + memcpy(tmp_buf, 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); + memcpy(tmp_buf+cs->instance_write_queue_len, buf+partial_pos, partial_len); cs->instance_write_queue = tmp_buf; - cs->instance_write_queue_pos = 0; - cs->instance_write_queue_len+=len; + cs->instance_write_queue_len+=partial_len; errno = EINPROGRESS; return -1; } diff --git a/plugins/http/http.c b/plugins/http/http.c index 44d9d409..8b37e7d9 100644 --- a/plugins/http/http.c +++ b/plugins/http/http.c @@ -728,6 +728,19 @@ void uwsgi_http_switch_events(struct uwsgi_corerouter *ucr, struct corerouter_se // data from instance if (interesting_fd == cs->instance_fd) { + // writable ? + if (cs->instance_fd_state) { + len = cs->instance_send(&uhttp.cr, cs, NULL, 0); +#ifdef UWSGI_EVENT_USE_PORT + event_queue_add_fd_write(uhttp_queue, cs->instance_fd); +#endif + if (len <= 0) { + if (len < 0 && errno == EINPROGRESS) break; + corerouter_close_session(ucr, cs); + } + break; + } + 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); @@ -775,6 +788,11 @@ To have a reliable implementation, we need to reset a bunch of values } break; } + // something to send in the queue ? + if (cs->write_queue) { + cs->write_queue_close = 1; + break; + } #endif } corerouter_close_session(ucr, cs); @@ -795,9 +813,22 @@ To have a reliable implementation, we need to reset a bunch of values cs->un->transferred += len; } - // body from client + // body from client or client ready to receive else if (interesting_fd == cs->fd) { + // writable ? + if (cs->fd_state) { + len = cs->send(&uhttp.cr, cs, NULL, 0); +#ifdef UWSGI_EVENT_USE_PORT + event_queue_add_fd_write(uhttp_queue, cs->fd); +#endif + if (len <= 0) { + if (len < 0 && errno == EINPROGRESS) break; + corerouter_close_session(ucr, cs); + } + break; + } + len = cs->recv(&uhttp.cr, cs, bbuf, UMAX16); #ifdef UWSGI_EVENT_USE_PORT event_queue_add_fd_read(uhttp_queue, cs->fd);