diff --git a/plugins/corerouter/corerouter.c b/plugins/corerouter/corerouter.c index a59d5e7c..ef4ec875 100644 --- a/plugins/corerouter/corerouter.c +++ b/plugins/corerouter/corerouter.c @@ -57,8 +57,6 @@ struct corerouter_peer *uwsgi_cr_peer_add(struct corerouter_session *cs) { cs->peers = peers; } - cs->refcnt++; - return peers; } @@ -279,8 +277,6 @@ void corerouter_manage_subscription(char *key, uint16_t keylen, char *val, uint1 } } -static struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *, struct corerouter_peer *); - void corerouter_close_peer(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) { struct corerouter_session *cs = peer->session; @@ -380,9 +376,7 @@ end: corerouter_close_session(ucr, cs); } else { - if (cs->refcnt > 0) - cs->refcnt--; - if (cs->refcnt == 0) { + if (cs->can_keepalive == 0) { corerouter_close_session(ucr, cs); } } @@ -411,7 +405,7 @@ void corerouter_close_session(struct uwsgi_corerouter *ucr, struct corerouter_se free(cr_session); } -static struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) { +struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *ucr, struct corerouter_peer *peer) { cr_del_timeout(ucr, peer); return cr_add_timeout(ucr, peer); } @@ -815,6 +809,8 @@ void uwsgi_corerouter_loop(int id, void *data) { } else if (ret < 0) { if (errno == EINPROGRESS) continue; + // remove keepalive on error + peer->session->can_keepalive = 0; corerouter_close_peer(ucr, peer); continue; } diff --git a/plugins/corerouter/cr.h b/plugins/corerouter/cr.h index e130a6a8..80a93d1d 100644 --- a/plugins/corerouter/cr.h +++ b/plugins/corerouter/cr.h @@ -270,13 +270,12 @@ struct corerouter_session { void (*close)(struct corerouter_session *); int (*retry)(struct corerouter_peer *); + int can_keepalive; + // this is the peer of the client struct corerouter_peer *main_peer; // this is the linked list of backends struct corerouter_peer *peers; - - // when it reaches 0 the session can be destroyed - uint64_t refcnt; }; void uwsgi_opt_corerouter(char *, char *, void *); @@ -317,3 +316,4 @@ int uwsgi_cr_set_hooks(struct corerouter_peer *, ssize_t (*)(struct corerouter_p struct corerouter_peer *uwsgi_cr_peer_add(struct corerouter_session *); struct corerouter_peer *uwsgi_cr_peer_find_by_sid(struct corerouter_session *, uint32_t); void corerouter_close_peer(struct uwsgi_corerouter *, struct corerouter_peer *); +struct uwsgi_rb_timer *corerouter_reset_timeout(struct uwsgi_corerouter *, struct corerouter_peer *); diff --git a/plugins/http/common.h b/plugins/http/common.h index 19b58764..1ff12119 100644 --- a/plugins/http/common.h +++ b/plugins/http/common.h @@ -60,7 +60,6 @@ struct http_session { char *path_info; uint16_t path_info_len; - int can_keepalive; int force_chunked; #ifdef UWSGI_SSL diff --git a/plugins/http/http.c b/plugins/http/http.c index 23dfb6da..9f9cf3fe 100644 --- a/plugins/http/http.c +++ b/plugins/http/http.c @@ -32,7 +32,7 @@ struct uwsgi_option http_options[] = { {"http-subscription-server", required_argument, 0, "enable the subscription server", uwsgi_opt_corerouter_ss, &uhttp, 0}, {"http-timeout", required_argument, 0, "set internal http socket timeout", uwsgi_opt_set_int, &uhttp.cr.socket_timeout, 0}, {"http-manage-expect", no_argument, 0, "manage the Expect HTTP request header", uwsgi_opt_true, &uhttp.manage_expect, 0}, - {"http-keepalive", no_argument, 0, "HTTP 1.1 keepalive support (non-pipelined) requests", uwsgi_opt_true, &uhttp.keepalive, 0}, + {"http-keepalive", optional_argument, 0, "HTTP 1.1 keepalive support (non-pipelined) requests", uwsgi_opt_set_int, &uhttp.keepalive, 0}, {"http-auto-chunked", no_argument, 0, "automatically transform output to chunked encoding during HTTP 1.1 keepalive (if needed)", uwsgi_opt_true, &uhttp.auto_chunked, 0}, {"http-raw-body", no_argument, 0, "blindly send HTTP body to backends (required for WebSockets and Icecast support in backends)", uwsgi_opt_true, &uhttp.raw_body, 0}, @@ -122,12 +122,12 @@ int http_add_uwsgi_header(struct corerouter_peer *peer, char *hh, uint16_t hhlen if (!uwsgi_strncmp("CONTENT_LENGTH", 14, hh, keylen)) { hr->content_length = uwsgi_str_num(val, vallen); - hr->can_keepalive = 0; + hr->session.can_keepalive = 0; } // in the future we could support chunked requests... if (!uwsgi_strncmp("TRANSFER_ENCODING", 17, hh, keylen)) { - hr->can_keepalive = 0; + hr->session.can_keepalive = 0; } if (uwsgi_strncmp("CONTENT_TYPE", 12, hh, keylen) && uwsgi_strncmp("CONTENT_LENGTH", 14, hh, keylen)) { @@ -137,7 +137,7 @@ int http_add_uwsgi_header(struct corerouter_peer *peer, char *hh, uint16_t hhlen if (!uwsgi_strncmp("CONNECTION", 10, hh, keylen)) { if (!uwsgi_strnicmp(val, vallen, "close", 5)) { - hr->can_keepalive = 0; + hr->session.can_keepalive = 0; } } @@ -227,7 +227,7 @@ int http_headers_parse(struct corerouter_peer *peer) { return 0; if (uwsgi_buffer_append_keyval(out, "SERVER_PROTOCOL", 15, base, ptr - base)) return -1; if (uhttp.keepalive && !uwsgi_strncmp("HTTP/1.1", 8, base, ptr-base)) { - hr->can_keepalive = 1; + hr->session.can_keepalive = 1; } ptr += 2; break; @@ -343,7 +343,7 @@ ssize_t hr_send_expect_continue(struct corerouter_peer *main_peer) { ssize_t hr_instance_write(struct corerouter_peer *peer) { ssize_t len = cr_write(peer, "hr_instance_write()"); // end on empty write - if (!len) return 0; + if (!len) { peer->session->can_keepalive = 0; return 0; } // the chunk has been sent, start (again) reading from client and instances if (cr_write_complete(peer)) { @@ -492,10 +492,15 @@ ssize_t hr_instance_read(struct corerouter_peer *peer) { struct http_session *hr = (struct http_session *) peer->session; ssize_t len = cr_read(peer, "hr_instance_read()"); if (!len) { - if (hr->can_keepalive) { + if (hr->session.can_keepalive) { peer->session->main_peer->disabled = 0; - peer->session->refcnt++; hr->rnrn = 0; + if (uhttp.keepalive > 1) { + int orig_timeout = peer->session->corerouter->socket_timeout; + peer->session->corerouter->socket_timeout = uhttp.keepalive; + peer->session->main_peer->timeout = corerouter_reset_timeout(peer->session->corerouter, peer->session->main_peer); + peer->session->corerouter->socket_timeout = orig_timeout; + } if (hr->force_chunked) { if (!hr->last_chunked) { hr->last_chunked = uwsgi_buffer_new(3); @@ -513,7 +518,7 @@ ssize_t hr_instance_read(struct corerouter_peer *peer) { } // need to parse response headers - if (hr->can_keepalive) { + if (hr->session.can_keepalive) { if (peer->r_parser_status != 4) { int ret = hr_check_response_keepalive(peer); if (ret < 0) return -1; @@ -626,7 +631,7 @@ ssize_t http_parse(struct corerouter_peer *main_peer) { new_peer->out->buf[2] = (uint8_t) ((pktsize >> 8) & 0xff); if (hr->remains > 0) { - hr->can_keepalive = 0; + hr->session.can_keepalive = 0; if (hr->content_length < hr->remains) { hr->remains = hr->content_length; hr->content_length = 0; @@ -637,7 +642,7 @@ ssize_t http_parse(struct corerouter_peer *main_peer) { if (uwsgi_buffer_append(new_peer->out, main_peer->in->buf + hr->headers_size + 1, hr->remains)) return -1; } - if (hr->can_keepalive) { + if (hr->session.can_keepalive) { main_peer->disabled = 1; // stop reading from the client if (uwsgi_cr_set_hooks(main_peer, NULL, NULL)) return -1; diff --git a/plugins/http/keepalive.c b/plugins/http/keepalive.c index 6ec3683f..0f781e21 100644 --- a/plugins/http/keepalive.c +++ b/plugins/http/keepalive.c @@ -13,7 +13,7 @@ int http_response_parse(struct http_session *hr, struct uwsgi_buffer *ub, size_t // protocol for(i=0;ican_keepalive && uwsgi_strncmp("HTTP/1.1", 8, buf, i)) { + if (hr->session.can_keepalive && uwsgi_strncmp("HTTP/1.1", 8, buf, i)) { goto end; } if (i+1 >= len) return -1;; @@ -63,7 +63,7 @@ int http_response_parse(struct http_session *hr, struct uwsgi_buffer *ub, size_t if (!colon) return -1; // security check if (colon+2 >= buf+len) return -1; - if (hr->can_keepalive) { + if (hr->session.can_keepalive) { if (!uwsgi_strnicmp(key, colon-key, "Connection", 10)) { if (!uwsgi_strnicmp(colon+2, h_len-((colon-key)+2), "close", 5)) { goto end; @@ -94,7 +94,7 @@ int http_response_parse(struct http_session *hr, struct uwsgi_buffer *ub, size_t } } - if (hr->can_keepalive && !has_size) { + if (hr->session.can_keepalive && !has_size) { if (uhttp.auto_chunked) { char cr = buf[len-2]; char nl = buf[len-1]; @@ -109,12 +109,12 @@ int http_response_parse(struct http_session *hr, struct uwsgi_buffer *ub, size_t return 0; } } - hr->can_keepalive = 0; + hr->session.can_keepalive = 0; } return 0; end: - hr->can_keepalive = 0; + hr->session.can_keepalive = 0; return 0; } diff --git a/plugins/http/spdy3.c b/plugins/http/spdy3.c index 7fb9455c..1f2100a1 100644 --- a/plugins/http/spdy3.c +++ b/plugins/http/spdy3.c @@ -575,7 +575,7 @@ ssize_t spdy_parse(struct corerouter_peer *main_peer) { if (deflateSetDictionary(&hr->spdy_z_out, (Bytef *) SPDY_dictionary_txt, sizeof(SPDY_dictionary_txt)) != Z_OK) { return -1; } - hr->session.refcnt++; + cs->can_keepalive = 1; hr->spdy_initialized = 1; hr->spdy_phase = UWSGI_SPDY_PHASE_HEADER;