diff --git a/bin/varnishd/cache/cache.h b/bin/varnishd/cache/cache.h index 7ad710995b3..9b57fc9d890 100644 --- a/bin/varnishd/cache/cache.h +++ b/bin/varnishd/cache/cache.h @@ -485,11 +485,13 @@ struct req { void *transport_priv; VTAILQ_ENTRY(req) w_list; + VTAILQ_ENTRY(req) t_list; struct objcore *body_oc; - /* The busy objhead we sleep on */ + /* The busy objhead we sleep on, referenced up to twice */ struct objhead *hash_objhead; + struct objhead *transport_objhead; /* Built Vary string == workspace reservation */ uint8_t *vary_b; diff --git a/bin/varnishd/cache/cache_esi_deliver.c b/bin/varnishd/cache/cache_esi_deliver.c index 5a7bba66cf7..18f490f9d97 100644 --- a/bin/varnishd/cache/cache_esi_deliver.c +++ b/bin/varnishd/cache/cache_esi_deliver.c @@ -47,7 +47,7 @@ #include "vgz.h" static vtr_deliver_f ved_deliver; -static vtr_reembark_f ved_reembark; +static vtr_waitlist_f ved_reembark; static const uint8_t gzip_hdr[] = { 0x1f, 0x8b, 0x08, @@ -91,7 +91,7 @@ static const struct transport VED_transport = { /*--------------------------------------------------------------------*/ -static void v_matchproto_(vtr_reembark_f) +static void v_matchproto_(vtr_waitlist_f) ved_reembark(struct worker *wrk, struct req *req) { struct ecx *ecx; diff --git a/bin/varnishd/cache/cache_hash.c b/bin/varnishd/cache/cache_hash.c index 337727c996e..e33ee404d1e 100644 --- a/bin/varnishd/cache/cache_hash.c +++ b/bin/varnishd/cache/cache_hash.c @@ -355,6 +355,63 @@ hsh_insert_busyobj(const struct worker *wrk, struct objhead *oh) return (oc); } +/*--------------------------------------------------------------------- + */ + +static void +hsh_waitlist(struct worker *wrk, struct req *req, struct objhead *oh) +{ + + CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC); + CHECK_OBJ_NOTNULL(req, REQ_MAGIC); + CHECK_OBJ_NOTNULL(oh, OBJHEAD_MAGIC); + Lck_AssertHeld(&oh->mtx); + AZ(req->hash_ignore_busy); + AZ(req->hash_objhead); + + DSLb(DBG_WAITINGLIST, "on waiting list <%p>", oh); + + /* + * The objhead reference transfers to the sess, we get it + * back when the sess comes off the waiting list and + * calls us again + */ + req->hash_objhead = oh; + req->wrk = NULL; + req->waitinglist = 1; + + if (req->top->topreq->transport->top_waitlist != NULL) + req->top->topreq->transport->top_waitlist(wrk, req); +} + +static void +hsh_reembark(struct worker *wrk, struct req *req) +{ + + CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC); + CHECK_OBJ_NOTNULL(req, REQ_MAGIC); + + if (DO_DEBUG(DBG_WAITINGLIST)) + VSLb(req->vsl, SLT_Debug, "off waiting list"); + + if (req->top->topreq->transport->top_reembark != NULL) + req->top->topreq->transport->top_reembark(wrk, req); + + if (req->transport->reembark != NULL) { + // For ESI includes + req->transport->reembark(wrk, req); + return; + } + + /* + * We ignore the queue limits which apply to new + * requests because if we fail to reschedule there + * may be vmod_privs to cleanup and we need a proper + * worker thread for that. + */ + AZ(Pool_Task(req->sp->pool, req->task, TASK_QUEUE_RUSH)); +} + /*--------------------------------------------------------------------- */ @@ -407,6 +464,12 @@ HSH_Lookup(struct req *req, struct objcore **ocp, struct objcore **bocp) CHECK_OBJ_NOTNULL(oh, OBJHEAD_MAGIC); Lck_AssertHeld(&oh->mtx); + if (req->walkaway) { + assert(oh->refcnt > 0); + (void)hsh_deref_objhead_unlock(wrk, &oh, 0); + return (HSH_WALKAWAY); + } + if (req->hash_always_miss) { /* XXX: should we do predictive Vary in this case ? */ /* Insert new objcore in objecthead and release mutex */ @@ -577,21 +640,11 @@ HSH_Lookup(struct req *req, struct objcore **ocp, struct objcore **bocp) /* There are one or more busy objects, wait for them */ VTAILQ_INSERT_TAIL(&oh->waitinglist, req, w_list); + hsh_waitlist(wrk, req, oh); - AZ(req->hash_ignore_busy); - - /* - * The objhead reference transfers to the sess, we get it - * back when the sess comes off the waiting list and - * calls us again - */ - req->hash_objhead = oh; - req->wrk = NULL; - req->waitinglist = 1; - - if (DO_DEBUG(DBG_WAITINGLIST)) - VSLb(req->vsl, SLT_Debug, "on waiting list <%p>", oh); - + /* NOTE OBS: We are going on the waitinglist. No changing of + * anything on struct req after releasing the mutex, as we may be + * rescheduled immediately. */ Lck_Unlock(&oh->mtx); wrk->stats->busy_sleep++; @@ -648,19 +701,43 @@ hsh_rush2(struct worker *wrk, struct rush *r) req = VTAILQ_FIRST(&r->reqs); CHECK_OBJ_NOTNULL(req, REQ_MAGIC); VTAILQ_REMOVE(&r->reqs, req, w_list); - DSL(DBG_WAITINGLIST, req->vsl->wid, "off waiting list"); - if (req->transport->reembark != NULL) { - // For ESI includes - req->transport->reembark(wrk, req); - } else { - /* - * We ignore the queue limits which apply to new - * requests because if we fail to reschedule there - * may be vmod_privs to cleanup and we need a proper - * workerthread for that. - */ - AZ(Pool_Task(req->sp->pool, req->task, TASK_QUEUE_RUSH)); - } + hsh_reembark(wrk, req); + } +} + +void +HSH_WalkAway(struct worker *wrk, struct objhead **ohp, struct req *req) +{ + struct objhead *oh; + struct rush r[1]; + + CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC); + CHECK_OBJ_NOTNULL(req, REQ_MAGIC); + TAKE_OBJ_NOTNULL(oh, ohp, OBJHEAD_MAGIC); + + INIT_OBJ(r, RUSH_MAGIC); + VTAILQ_INIT(&r->reqs); + + if (DO_DEBUG(DBG_WAITINGLIST)) + VSLb(req->vsl, SLT_Debug, "walking away <%p>", oh); + + Lck_Lock(&oh->mtx); + if (req->waitinglist) { + assert(oh == req->hash_objhead); + assert(oh->refcnt > 1); + oh->refcnt--; + VTAILQ_REMOVE(&oh->waitinglist, req, w_list); + VTAILQ_INSERT_TAIL(&r->reqs, req, w_list); + Lck_Unlock(&oh->mtx); + + wrk->stats->busy_killed++; + AZ(req->walkaway); + req->walkaway = 1; + hsh_rush2(wrk, r); + } else { + AZ(req->hash_objhead); + assert(oh->refcnt > 0); + (void)hsh_deref_objhead_unlock(wrk, &oh, 0); } } diff --git a/bin/varnishd/cache/cache_objhead.h b/bin/varnishd/cache/cache_objhead.h index bc1782379e3..d455a16da9c 100644 --- a/bin/varnishd/cache/cache_objhead.h +++ b/bin/varnishd/cache/cache_objhead.h @@ -76,3 +76,4 @@ unsigned HSH_Purge(struct worker *, struct objhead *, vtim_real ttl_now, vtim_dur ttl, vtim_dur grace, vtim_dur keep); struct objcore *HSH_Private(const struct worker *wrk); void HSH_Cancel(struct worker *, struct objcore *, struct boc *); +void HSH_WalkAway(struct worker *wrk, struct objhead **ohp, struct req *req); diff --git a/bin/varnishd/cache/cache_req_fsm.c b/bin/varnishd/cache/cache_req_fsm.c index 1a1fd7e2a33..12728ee154b 100644 --- a/bin/varnishd/cache/cache_req_fsm.c +++ b/bin/varnishd/cache/cache_req_fsm.c @@ -562,6 +562,21 @@ cnt_lookup(struct worker *wrk, struct req *req) had_objhead = 1; wrk->strangelove = 0; lr = HSH_Lookup(req, &oc, &busy); + if (lr == HSH_WALKAWAY) { + VRY_Finish(req, DISCARD); + AZ(req->objcore); + AZ(req->stale_oc); + VSLb_ts_req(req, "Waitinglist", W_TIM_real(wrk)); + VSLb(req->vsl, SLT_Error, "The client is going away"); + if (req->esi_level > 0) { + /* NB: Interrupt delivery upwards to avoid engaging + * new sub-request tasks. + */ + req->req_step = R_STP_VCLFAIL; + return (REQ_FSM_MORE); + } + return (REQ_FSM_DONE); + } if (lr == HSH_BUSY) { /* * We lost the session to a busy object, disembark the diff --git a/bin/varnishd/cache/cache_transport.h b/bin/varnishd/cache/cache_transport.h index 0e5f03c65ef..f3651e5ea63 100644 --- a/bin/varnishd/cache/cache_transport.h +++ b/bin/varnishd/cache/cache_transport.h @@ -43,7 +43,7 @@ typedef void vtr_req_body_f (struct req *); typedef void vtr_sess_panic_f (struct vsb *, const struct sess *); typedef void vtr_req_panic_f (struct vsb *, const struct req *); typedef void vtr_req_fail_f (struct req *, stream_close_t); -typedef void vtr_reembark_f (struct worker *, struct req *); +typedef void vtr_waitlist_f (struct worker *, struct req *); typedef int vtr_minimal_response_f (struct req *, uint16_t status); struct transport { @@ -63,7 +63,9 @@ struct transport { vtr_deliver_f *deliver; vtr_sess_panic_f *sess_panic; vtr_req_panic_f *req_panic; - vtr_reembark_f *reembark; + vtr_waitlist_f *top_waitlist; + vtr_waitlist_f *top_reembark; + vtr_waitlist_f *reembark; vtr_minimal_response_f *minimal_response; VTAILQ_ENTRY(transport) list; diff --git a/bin/varnishd/hash/hash_slinger.h b/bin/varnishd/hash/hash_slinger.h index a1a9c0e8214..28899772e31 100644 --- a/bin/varnishd/hash/hash_slinger.h +++ b/bin/varnishd/hash/hash_slinger.h @@ -52,7 +52,7 @@ struct hash_slinger { }; enum lookup_e { - HSH_CONTINUE, + HSH_WALKAWAY, HSH_MISS, HSH_BUSY, HSH_HIT, diff --git a/bin/varnishd/http2/cache_http2.h b/bin/varnishd/http2/cache_http2.h index eb6e8f6a9f3..7069ad053b9 100644 --- a/bin/varnishd/http2/cache_http2.h +++ b/bin/varnishd/http2/cache_http2.h @@ -134,6 +134,7 @@ struct h2_req { int counted; struct h2_sess *h2sess; struct req *req; + VTAILQ_HEAD(, req) waitinglist; double t_send; double t_winupd; pthread_cond_t *cond; diff --git a/bin/varnishd/http2/cache_http2_proto.c b/bin/varnishd/http2/cache_http2_proto.c index 63c491fbd2c..98588a4bf2a 100644 --- a/bin/varnishd/http2/cache_http2_proto.c +++ b/bin/varnishd/http2/cache_http2_proto.c @@ -159,6 +159,7 @@ h2_new_req(struct h2_sess *h2, unsigned stream, struct req *req) r2->r_window = h2->local_settings.initial_window_size; r2->t_window = h2->remote_settings.initial_window_size; req->transport_priv = r2; + VTAILQ_INIT(&r2->waitinglist); Lck_Lock(&h2->sess->mtx); if (stream) h2->open_streams++; @@ -182,6 +183,7 @@ h2_del_req(struct worker *wrk, struct h2_req *r2) ASSERT_RXTHR(h2); sp = h2->sess; Lck_Lock(&sp->mtx); + assert(VTAILQ_EMPTY(&r2->waitinglist)); assert(h2->refcnt > 0); --h2->refcnt; /* XXX: PRIORITY reshuffle */ @@ -207,9 +209,13 @@ void h2_kill_req(struct worker *wrk, struct h2_sess *h2, struct h2_req *r2, h2_error h2e) { + VTAILQ_HEAD(, req) wl; + struct req *req, *t; + struct objhead *oh; ASSERT_RXTHR(h2); AN(h2e); + VTAILQ_INIT(&wl); Lck_Lock(&h2->sess->mtx); VSLb(h2->vsl, SLT_Debug, "KILL st=%u state=%d sched=%d", r2->stream, r2->state, r2->scheduled); @@ -221,16 +227,34 @@ h2_kill_req(struct worker *wrk, struct h2_sess *h2, if (r2->error == NULL) r2->error = h2e; if (r2->scheduled) { - if (r2->cond != NULL) - AZ(pthread_cond_signal(r2->cond)); - r2 = NULL; + VTAILQ_FOREACH_SAFE(req, &r2->waitinglist, t_list, t) { + CHECK_OBJ(req, REQ_MAGIC); + VTAILQ_REMOVE(&r2->waitinglist, req, t_list); + VTAILQ_INSERT_TAIL(&wl, req, t_list); + } + if (VTAILQ_EMPTY(&wl)) { + if (r2->cond != NULL) + AZ(pthread_cond_signal(r2->cond)); + r2 = NULL; + } } else { if (r2->state == H2_S_OPEN && h2->new_req == r2->req) (void)h2h_decode_fini(h2); } Lck_Unlock(&h2->sess->mtx); - if (r2 != NULL) + if (VTAILQ_EMPTY(&wl) && r2 != NULL) { h2_del_req(wrk, r2); + return; + } + VTAILQ_FOREACH_SAFE(req, &wl, t_list, t) { + Lck_Lock(&h2->sess->mtx); + VTAILQ_REMOVE(&wl, req, t_list); + oh = req->transport_objhead; + req->transport_objhead = NULL; + Lck_Unlock(&h2->sess->mtx); + HSH_WalkAway(wrk, &oh, req); + AZ(oh); + } } /**********************************************************************/ diff --git a/bin/varnishd/http2/cache_http2_session.c b/bin/varnishd/http2/cache_http2_session.c index 846b319f183..04236f51438 100644 --- a/bin/varnishd/http2/cache_http2_session.c +++ b/bin/varnishd/http2/cache_http2_session.c @@ -36,6 +36,7 @@ #include #include "cache/cache_transport.h" +#include "cache/cache_objhead.h" #include "http2/cache_http2.h" #include "vend.h" @@ -438,6 +439,72 @@ h2_new_session(struct worker *wrk, void *arg) wrk->vsl = NULL; } +/********************************************************************** + */ + +static void +h2_top_waitlist(struct worker *wrk, struct req *req) +{ + struct objhead *oh; + struct h2_req *r2; + + CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC); + CHECK_OBJ_NOTNULL(req, REQ_MAGIC); + oh = req->hash_objhead; + CHECK_OBJ_NOTNULL(oh, OBJHEAD_MAGIC); + CAST_OBJ_NOTNULL(r2, req->top->topreq->transport_priv, H2_REQ_MAGIC); + + if (DO_DEBUG(DBG_WAITINGLIST)) + VSLb(req->vsl, SLT_Debug, "on h2 waiting list <%p>", r2); + + Lck_AssertHeld(&oh->mtx); + AN(req->waitinglist); + AZ(req->wrk); + + assert(oh->refcnt > 0); + oh->refcnt++; + + Lck_Lock(&r2->h2sess->sess->mtx); + AZ(req->transport_objhead); + req->transport_objhead = oh; + VTAILQ_INSERT_TAIL(&r2->waitinglist, req, t_list); + Lck_Unlock(&r2->h2sess->sess->mtx); +} + +static void +h2_top_reembark(struct worker *wrk, struct req *req) +{ + struct objhead *oh; + struct h2_req *r2; + + CHECK_OBJ_NOTNULL(wrk, WORKER_MAGIC); + CHECK_OBJ_NOTNULL(req, REQ_MAGIC); + CAST_OBJ_NOTNULL(r2, req->top->topreq->transport_priv, H2_REQ_MAGIC); + + if (DO_DEBUG(DBG_WAITINGLIST)) + VSLb(req->vsl, SLT_Debug, "off h2 waiting list <%p>", r2); + + Lck_Lock(&r2->h2sess->sess->mtx); + oh = req->transport_objhead; + CHECK_OBJ_ORNULL(oh, OBJHEAD_MAGIC); + if (oh != NULL) { + VTAILQ_REMOVE(&r2->waitinglist, req, t_list); + req->transport_objhead = NULL; + } + Lck_Unlock(&r2->h2sess->sess->mtx); + + if (oh != NULL) { + Lck_Lock(&oh->mtx); + AN(req->hash_objhead); + AZ(req->waitinglist); + AZ(req->wrk); + AZ(req->walkaway); + assert(oh->refcnt > 1); + oh->refcnt--; + Lck_Unlock(&oh->mtx); + } +} + struct transport HTTP2_transport = { .name = "HTTP/2", .magic = TRANSPORT_MAGIC, @@ -447,4 +514,6 @@ struct transport HTTP2_transport = { .req_body = h2_req_body, .req_fail = h2_req_fail, .sess_panic = h2_sess_panic, + .top_waitlist = h2_top_waitlist, + .top_reembark = h2_top_reembark, }; diff --git a/bin/varnishtest/tests/t02023.vtc b/bin/varnishtest/tests/t02023.vtc new file mode 100644 index 00000000000..b0284e69821 --- /dev/null +++ b/bin/varnishtest/tests/t02023.vtc @@ -0,0 +1,66 @@ +varnishtest "Dying h2 session with stream on waiting list" + +barrier b1 cond 2 +barrier b2 cond 2 +barrier b3 cond 2 + +server s1 { + rxreq + barrier b1 sync + barrier b3 sync + txresp +} -start + +varnish v1 -cliok "param.set thread_pools 1" +varnish v1 -cliok "param.set feature +http2" +varnish v1 -cliok "param.set debug +waitinglist" +varnish v1 -cliok "param.set debug +syncvsl" +varnish v1 -vcl+backend {} -start + +logexpect l1 -v v1 -g raw { + expect * 1003 Debug "on waiting list" +} -start + +logexpect l2 -v v1 { + expect * 1003 Debug "walking away" + expect 0 = Debug "off waiting list" + expect * = Timestamp "^Waitinglist: " + expect 0 = Error "The client is going away" +} -start + +client c1 { + stream 1 { + txreq + barrier b1 sync + } -run + stream 3 { + txreq + } -run + stream 5 { + barrier b2 sync + txsettings + } -run + stream 0 { + rxgoaway + expect goaway.laststream == 3 + expect goaway.err == PROTOCOL_ERROR + } -run +} -start + +logexpect l1 -wait +barrier b2 sync + +logexpect l2 -wait +barrier b3 sync + +server s1 -wait + +varnish v1 -expect busy_killed == 1 + +client c2 { + txreq + rxresp + expect resp.status == 200 +} -run + +varnish v1 -expect cache_hit == 1 diff --git a/bin/varnishtest/tests/t02024.vtc b/bin/varnishtest/tests/t02024.vtc new file mode 100644 index 00000000000..37a3e018abe --- /dev/null +++ b/bin/varnishtest/tests/t02024.vtc @@ -0,0 +1,90 @@ +varnishtest "Dying h2 session with ESI stream on waiting list" + +barrier b1 cond 2 +barrier b2 cond 2 +barrier b3 cond 2 + +server s1 { + rxreq + expect req.url == "/esi" + barrier b1 sync + barrier b3 sync + txresp -hdr "included: esi" -body hello +} -start + +server s2 { + rxreq + expect req.url == "/" + txresp -body {} +} -start + +varnish v1 -cliok "param.set thread_pools 1" +varnish v1 -cliok "param.set feature +http2" +varnish v1 -cliok "param.set debug +waitinglist" +varnish v1 -cliok "param.set debug +syncvsl" +varnish v1 -vcl+backend { + sub vcl_recv { + if (req.url == "/") { + set req.backend_hint = s2; + } + } + sub vcl_backend_response { + set beresp.do_esi = !beresp.http.included; + } +} -start + +logexpect l1 -v v1 -g raw { + expect * 1006 Debug "on waiting list" +} -start + +logexpect l2 -v v1 -g raw { + expect * 1006 Debug "walking away" + expect 0 1006 Debug "off waiting list" +} -start + +# c1 will block during a backend fetch +client c1 { + txreq -url "/esi" + rxresp + expect resp.body == hello +} -start + +barrier b1 sync + +client c2 { + stream 1 { + txreq + } -run + stream 3 { + barrier b2 sync + # bork the session + txsettings + } -run + stream 0 { + rxgoaway + expect goaway.laststream == 1 + expect goaway.err == PROTOCOL_ERROR + } -run +} -start + +logexpect l1 -wait +barrier b2 sync + +logexpect l2 -wait +barrier b3 sync + +client c1 -wait +client c2 -wait +server s1 -wait +server s2 -wait + +varnish v1 -expect busy_killed == 1 +varnish v1 -expect cache_hit == 0 + +client c3 { + txreq + rxresp + expect resp.body == hellohellohello +} -run + +varnish v1 -expect cache_hit == 4 diff --git a/include/tbl/req_flags.h b/include/tbl/req_flags.h index 88eaff6391c..296eac813ca 100644 --- a/include/tbl/req_flags.h +++ b/include/tbl/req_flags.h @@ -40,6 +40,7 @@ REQ_FLAG(is_hit, 0, 0, "") REQ_FLAG(waitinglist, 0, 0, "") REQ_FLAG(want100cont, 0, 0, "") REQ_FLAG(late100cont, 0, 0, "") +REQ_FLAG(walkaway, 0, 0, "") #define REQ_BEREQ_FLAG(lower, vcl_r, vcl_w, doc) \ REQ_FLAG(lower, vcl_r, vcl_w, doc) #include "tbl/req_bereq_flags.h" diff --git a/lib/libvsc/VSC_main.vsc b/lib/libvsc/VSC_main.vsc index 047218ad88a..d3925926eca 100644 --- a/lib/libvsc/VSC_main.vsc +++ b/lib/libvsc/VSC_main.vsc @@ -319,10 +319,11 @@ Number of requests taken off the busy object sleep list and rescheduled. .. varnish_vsc:: busy_killed + :group: wrk :oneliner: Number of requests killed after sleep on busy objhdr Number of requests killed from the busy object sleep list due to - lack of resources. + the client going away. .. varnish_vsc:: sess_queued :oneliner: Sessions queued for thread