mirror of
https://github.com/clearlinux/uwsgi.git
synced 2026-08-27 08:55:48 +00:00
another fix for write queue
This commit is contained in:
@@ -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 *);
|
||||
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
+32
-1
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user