From c8a608999453f0bab97bd15545dfdde78ed563eb Mon Sep 17 00:00:00 2001 From: jamie Date: Sun, 27 Sep 2026 17:57:22 +0100 Subject: [PATCH 1/2] fix(proxy): time out an origin body that goes idle --- README.md | 7 +- src/config.rs | 4 +- src/proxy.rs | 95 ++++++++++-- tests/integration/body_idle.rs | 262 ++++++++++++++++++++++++++++++++ tests/integration/coalescing.rs | 31 ++++ tests/integration/main.rs | 1 + 6 files changed, 382 insertions(+), 18 deletions(-) create mode 100644 tests/integration/body_idle.rs diff --git a/README.md b/README.md index 265fe03..8938011 100644 --- a/README.md +++ b/README.md @@ -34,7 +34,7 @@ The origin then receives one of each of these headers: There is no trusted-proxy setting. Behind another load balancer or proxy, `X-Forwarded-For` holds that proxy's address, and any forwarding headers that proxy adds are removed. -`--upstream-timeout-seconds` (default 30) limits how long shadowstep waits for the origin's response headers, including the time to send the request body. The limit does not cover the response body. +`--upstream-timeout-seconds` (default 30) limits how long shadowstep waits for the origin's response headers, including the time to send the request body. It also limits how long the origin may send nothing while the proxy waits for the next part of the response body. The body may take longer in total, as long as data keeps arriving. When the origin stays silent for longer, shadowstep logs a warning and closes the client's connection before the body is complete: a response with `Content-Length` ends short, and a chunked response gets no last chunk. The proxy does not store that response. Error responses have short generic bodies: @@ -98,7 +98,7 @@ When several requests miss on the same key at once, the first one, the leader, g - Only `GET` and `HEAD` requests that may be served from the cache wait. A request with `Authorization`, `Cookie`, a conditional header, request `Cache-Control: no-cache` or `no-store`, or a method-override header goes to the origin as before. - Only a `GET` leads. A `HEAD` request with no `GET` in flight goes to the origin. -- The leader's response streams to its client as usual. The waiting requests look up the cache once the response is stored, or once the proxy knows it will not be stored: it is not storable, its body passes the entry size limit or fails, the leader's client disconnects, or the origin fails or times out. +- The leader's response streams to its client as usual. The waiting requests look up the cache once the response is stored, or once the proxy knows it will not be stored: it is not storable, its body passes the entry size limit or fails, the leader's client disconnects, or the origin fails, times out or stops sending the body. - A waiting request gets the stored response with `X-Shadowstep-Cache: COALESCED`. If the cache has nothing it may use, for example because the response was `private` or had `Vary` values that differ from the waiting request's, the waiting request goes to the origin on its own. One client's uncacheable response never goes to another client. - A request waits for at most `--upstream-timeout-seconds`, then goes to the origin on its own. @@ -119,7 +119,6 @@ Known limits: - A response with `Cache-Control: no-cache` is not stored, although RFC 9111 allows storing it and revalidating it on every use. - A stored response is never used to answer a client's conditional request with a 304. A fresh hit always gets the full response. - Requests with `Cookie` do not coalesce. If browsers send a cookie with every request to the site, only cookieless clients coalesce. -- The proxy notices that a client has disconnected only when it next writes to it. If the leader's client disconnects while the origin has stopped sending the body, the waiting requests go to the origin after `--upstream-timeout-seconds`. - The host is part of the key, so a client that sends many different `Host` values can create many entries. The byte bound on the cache still applies. - Each process has its own cache. Replicas do not share entries or invalidations. @@ -184,7 +183,7 @@ Each option can be set with a flag or an environment variable. The flag wins if | `--tls-cert` | `TLS_CERT_PATH` | none | PEM certificate chain | | `--tls-key` | `TLS_KEY_PATH` | none | PEM private key in PKCS#8 form | | `--tls-listen-addr` | `TLS_LISTEN_ADDR` | `0.0.0.0:8443` | Address for the HTTPS listener, used only when both TLS paths are set | -| `--upstream-timeout-seconds` | `UPSTREAM_TIMEOUT_SECONDS` | `30` | Seconds to wait for the origin's response headers before answering 504 | +| `--upstream-timeout-seconds` | `UPSTREAM_TIMEOUT_SECONDS` | `30` | Seconds to wait for the origin's response headers before answering 504, and the longest the origin may send nothing during a response body | `cargo run -- --help` prints the same list. diff --git a/src/config.rs b/src/config.rs index 0be4187..8686b9d 100644 --- a/src/config.rs +++ b/src/config.rs @@ -36,7 +36,9 @@ pub struct Config { pub tls_listen_addr: String, /// seconds to wait for the origin's response headers, including the time - /// to send the request body. the proxy answers 504 when it runs out. + /// to send the request body. the proxy answers 504 when it runs out. it + /// is also the longest the origin may send nothing during a response + /// body before the proxy closes the client's connection. #[clap(long, env = "UPSTREAM_TIMEOUT_SECONDS", default_value_t = 30)] pub upstream_timeout_seconds: u64, } diff --git a/src/proxy.rs b/src/proxy.rs index 6881d65..f31afd3 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -11,9 +11,13 @@ use hyper::body::{Frame, Incoming}; use hyper::{Request as HyperRequest, Uri}; use log::{debug, error, warn}; use std::convert::TryFrom; +use std::future::Future; +use std::io; use std::pin::Pin; use std::task::{Context, Poll}; +use std::time::Duration; use tokio::sync::mpsc; +use tokio::time::{Instant, Sleep}; use url::{Position, Url}; use crate::cache::{ @@ -100,7 +104,7 @@ pub async fn forward_to_upstream( "Received response from upstream: {:?}", upstream_response.status() ); - exchange.answered(upstream_response) + exchange.answered(upstream_response, target_uri) } Ok(Err(e)) => { error!("Error forwarding request to upstream {}: {}", target_uri, e); @@ -393,7 +397,7 @@ struct Exchange { impl Exchange { /// the client's response once the origin has answered. - fn answered(self, response: hyper::Response) -> HttpResponse { + fn answered(self, response: hyper::Response, target_uri: Uri) -> HttpResponse { let (parts, body) = response.into_parts(); let head = Head::from_upstream(&parts); if let Some(stale) = &self.stale { @@ -413,6 +417,7 @@ impl Exchange { &self.request_policy, &head, ); + let body = origin_body(body, self.state.upstream_timeout, target_uri); client_response(head, body, store.map(|store| store.holding(flight))) } @@ -696,7 +701,10 @@ struct BodyCopy { /// turns the origin's response into the client's, streaming the body. the /// response carries the origin's length when it sent one. with `store`, a /// body no longer than its limit is also copied into the cache as it streams. -fn client_response(head: Head, body: Incoming, store: Option) -> HttpResponse { +fn client_response(head: Head, body: S, store: Option) -> HttpResponse +where + S: Stream> + Unpin + 'static, +{ let mut builder = HttpResponse::build(head.status); let fields = forwardable_fields(&head); @@ -706,14 +714,6 @@ fn client_response(head: Head, body: Incoming, store: Option) -> HttpRe builder.insert_header((CACHE_STATUS, "MISS")); let stored_headers = stored_fields(fields); - // the data stream drops trailer frames - let body = body.into_data_stream().map(|chunk| { - chunk.map_err(|e| { - error!("Error reading upstream response body: {}", e); - e - }) - }); - let length = head .map .get(header::CONTENT_LENGTH) @@ -770,6 +770,75 @@ fn cached_response(stored: &StoredResponse, cache_status: &'static str) -> HttpR .body(stored.body.clone()) } +/// the origin's response body as a stream of data chunks, which ends with an +/// error once the origin has sent nothing for `idle`. +fn origin_body( + body: Incoming, + idle: Duration, + target_uri: Uri, +) -> IdleTimeout> + Unpin> { + // the data stream drops trailer frames + let body = body.into_data_stream().map(|chunk| { + chunk.map_err(|e| { + error!("Error reading upstream response body: {}", e); + io::Error::other(e) + }) + }); + IdleTimeout { + inner: body, + idle, + sleep: Box::pin(tokio::time::sleep(idle)), + waiting: false, + expired: false, + target_uri, + } +} + +/// passes a body stream through, and ends it with an error when the next +/// chunk takes longer than `idle` to arrive. the error makes actix close +/// the client's connection before the body is whole, and makes the `Tee` +/// drop its copy. +struct IdleTimeout { + inner: S, + idle: Duration, + sleep: Pin>, + /// the deadline runs from the first poll after the last chunk, so time + /// spent waiting for a slow client does not count against the origin + waiting: bool, + expired: bool, + target_uri: Uri, +} + +impl Stream for IdleTimeout +where + S: Stream> + Unpin, +{ + type Item = S::Item; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + let this = &mut *self; + if this.expired { + return Poll::Ready(None); + } + if !this.waiting { + this.sleep.as_mut().reset(Instant::now() + this.idle); + this.waiting = true; + } + if let Poll::Ready(item) = this.inner.poll_next_unpin(cx) { + this.waiting = false; + return Poll::Ready(item); + } + std::task::ready!(this.sleep.as_mut().poll(cx)); + this.expired = true; + warn!( + "Upstream {} sent no response body data for {:?}", + this.target_uri, this.idle + ); + let error = io::Error::new(io::ErrorKind::TimedOut, "origin response body went idle"); + Poll::Ready(Some(Err(error))) + } +} + /// passes a body stream through while copying it into a buffer. once the /// body has ended within the limit, the buffer goes to the finish callback. /// a body that passes the limit or fails is not kept. @@ -805,9 +874,9 @@ impl Tee { } } -impl Stream for Tee +impl Stream for Tee where - S: Stream> + Unpin, + S: Stream> + Unpin, { type Item = S::Item; diff --git a/tests/integration/body_idle.rs b/tests/integration/body_idle.rs new file mode 100644 index 0000000..45b674d --- /dev/null +++ b/tests/integration/body_idle.rs @@ -0,0 +1,262 @@ +use crate::common::{self, Running}; + +use bytes::Bytes; +use http_body_util::Empty; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; +use std::time::{Duration, Instant}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::{TcpListener, TcpStream}; +use tokio::sync::mpsc; +use tokio::time::timeout; + +/// the upstream timeout in these tests, which is also the longest the origin +/// may stay silent during a response body. +const IDLE: Duration = Duration::from_secs(1); + +/// how long past `IDLE` the proxy may take to give up on a silent origin. +const MARGIN: Duration = Duration::from_millis(1500); + +/// longer than any step in these tests should take on a loaded CI runner. +const PATIENCE: Duration = Duration::from_secs(10); + +/// what the origin sends on one connection after the request head. +struct Reply { + head: String, + /// sent in order, each after its delay + parts: Vec<(Duration, Vec)>, + /// after the parts, hold the connection open until the proxy closes it + stall: bool, +} + +fn cacheable_head(framing: &str) -> String { + format!("HTTP/1.1 200 OK\r\ncache-control: max-age=60\r\n{framing}\r\n\r\n") +} + +/// 40 of 100 bytes, then silence. +fn stalled_sized() -> Reply { + Reply { + head: cacheable_head("content-length: 100"), + parts: vec![(Duration::ZERO, vec![b'a'; 40])], + stall: true, + } +} + +/// one chunk of 40 bytes, then silence with no last chunk. +fn stalled_chunked() -> Reply { + let mut chunk = b"28\r\n".to_vec(); + chunk.extend_from_slice(&[b'a'; 40]); + chunk.extend_from_slice(b"\r\n"); + Reply { + head: cacheable_head("transfer-encoding: chunked"), + parts: vec![(Duration::ZERO, chunk)], + stall: true, + } +} + +fn whole(body: &[u8]) -> Reply { + Reply { + head: cacheable_head(&format!("content-length: {}", body.len())), + parts: vec![(Duration::ZERO, body.to_vec())], + stall: false, + } +} + +/// an origin that answers its n-th connection with `replies[n]`. it counts +/// connections and reports on `closed` each stalled connection that the +/// proxy closes. +struct Origin { + url: String, + connections: Arc, + closed: mpsc::UnboundedReceiver, +} + +async fn scripted_origin(replies: Vec) -> Origin { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}", listener.local_addr().unwrap()); + let connections = Arc::new(AtomicUsize::new(0)); + let (closed_tx, closed) = mpsc::unbounded_channel(); + let counter = connections.clone(); + actix_web::rt::spawn(async move { + for reply in replies { + let (stream, _) = listener.accept().await.unwrap(); + let n = counter.fetch_add(1, Ordering::SeqCst); + let closed_tx = closed_tx.clone(); + actix_web::rt::spawn(async move { + if answer(stream, reply).await { + let _ = closed_tx.send(n); + } + }); + } + // any further connection counts but gets no answer + while let Ok((stream, _)) = listener.accept().await { + counter.fetch_add(1, Ordering::SeqCst); + drop(stream); + } + }); + Origin { + url, + connections, + closed, + } +} + +/// sends `reply` after the request head. returns true when the reply +/// stalled and the proxy then closed the connection. +async fn answer(mut stream: TcpStream, reply: Reply) -> bool { + let mut head = Vec::new(); + let mut buf = [0; 1024]; + while !head.windows(4).any(|w| w == b"\r\n\r\n") { + match stream.read(&mut buf).await { + Ok(0) | Err(_) => return false, + Ok(n) => head.extend_from_slice(&buf[..n]), + } + } + let _ = stream.write_all(reply.head.as_bytes()).await; + for (delay, part) in reply.parts { + actix_web::rt::time::sleep(delay).await; + if stream.write_all(&part).await.is_err() { + return false; + } + } + if !reply.stall { + return false; + } + // a reset closes the connection as surely as an EOF + while matches!(stream.read(&mut buf).await, Ok(n) if n > 0) {} + true +} + +fn spawn(origin: &Origin) -> Running { + common::spawn_with(&origin.url, |c| { + c.upstream_timeout_seconds = IDLE.as_secs(); + }) +} + +/// the proxy's whole response to a GET of `path` on a raw connection, which +/// ends when the proxy closes it, and how long that took. +async fn raw_get(proxy: &Running, path: &str) -> (Vec, Duration) { + let started = Instant::now(); + let mut client = TcpStream::connect(proxy.addr).await.unwrap(); + let request = format!("GET {path} HTTP/1.1\r\nhost: localhost\r\nconnection: close\r\n\r\n"); + client.write_all(request.as_bytes()).await.unwrap(); + let mut response = Vec::new(); + timeout(PATIENCE, client.read_to_end(&mut response)) + .await + .expect("the proxy kept the client connection open") + // a reset ends the response as surely as an EOF + .ok(); + (response, started.elapsed()) +} + +fn split(response: &[u8]) -> (String, &[u8]) { + let end = response + .windows(4) + .position(|w| w == b"\r\n\r\n") + .expect("no end of response head"); + let head = String::from_utf8_lossy(&response[..end]).to_ascii_lowercase(); + (head, &response[end + 4..]) +} + +fn assert_given_up_in_time(elapsed: Duration) { + assert!(elapsed >= IDLE - Duration::from_millis(100), "{elapsed:?}"); + assert!(elapsed < IDLE + MARGIN, "took {elapsed:?}"); +} + +/// waits until the proxy has closed the origin connection `n`. +async fn origin_connection_closed(origin: &mut Origin, n: usize) { + let closed = timeout(PATIENCE, origin.closed.recv()) + .await + .expect("the proxy kept the origin connection open"); + assert_eq!(closed, Some(n)); +} + +async fn assert_healthy(proxy: &Running) { + let client = common::client::>(); + let resp = timeout(PATIENCE, client.get(proxy.url("/health").parse().unwrap())) + .await + .expect("health check timed out") + .unwrap(); + assert_eq!(resp.status(), 200); +} + +#[actix_web::test] +async fn sized_body_that_stalls_closes_the_client_connection_short() { + let mut origin = scripted_origin(vec![stalled_sized()]).await; + let proxy = spawn(&origin); + + let (response, elapsed) = raw_get(&proxy, "/page").await; + + assert_given_up_in_time(elapsed); + let (head, body) = split(&response); + assert!(head.starts_with("http/1.1 200 ok"), "{head}"); + assert!(head.contains("content-length: 100"), "{head}"); + assert!(body.len() < 100, "got {} bytes", body.len()); + origin_connection_closed(&mut origin, 0).await; + assert_healthy(&proxy).await; + proxy.stop().await; +} + +#[actix_web::test] +async fn chunked_body_that_stalls_gets_no_last_chunk() { + let mut origin = scripted_origin(vec![stalled_chunked()]).await; + let proxy = spawn(&origin); + + let (response, elapsed) = raw_get(&proxy, "/page").await; + + assert_given_up_in_time(elapsed); + let (head, body) = split(&response); + assert!(head.starts_with("http/1.1 200 ok"), "{head}"); + assert!(head.contains("transfer-encoding: chunked"), "{head}"); + // the body holds only letters, so a `0\r\n\r\n` can only be the last chunk + assert!( + !body.windows(5).any(|w| w == b"0\r\n\r\n"), + "{}", + String::from_utf8_lossy(body) + ); + origin_connection_closed(&mut origin, 0).await; + assert_healthy(&proxy).await; + proxy.stop().await; +} + +#[actix_web::test] +async fn slow_but_steady_body_arrives_whole() { + // ten chunks 300 ms apart take 3 seconds, three times the idle limit + let parts = (0..10) + .map(|_| (Duration::from_millis(300), vec![b'a'; 10])) + .collect(); + let steady = Reply { + head: cacheable_head("content-length: 100"), + parts, + stall: false, + }; + let origin = scripted_origin(vec![steady]).await; + let proxy = spawn(&origin); + + let (response, elapsed) = raw_get(&proxy, "/page").await; + + assert!(elapsed >= Duration::from_secs(3), "{elapsed:?}"); + let (head, body) = split(&response); + assert!(head.starts_with("http/1.1 200 ok"), "{head}"); + assert_eq!(body, [b'a'; 100]); + assert_healthy(&proxy).await; + proxy.stop().await; +} + +#[actix_web::test] +async fn stalled_cacheable_body_is_not_stored() { + let mut origin = scripted_origin(vec![stalled_sized(), whole(b"whole")]).await; + let proxy = spawn(&origin); + + let (_, elapsed) = raw_get(&proxy, "/page").await; + assert_given_up_in_time(elapsed); + origin_connection_closed(&mut origin, 0).await; + let (response, _) = raw_get(&proxy, "/page").await; + + let (head, body) = split(&response); + assert!(head.contains("x-shadowstep-cache: miss"), "{head}"); + assert_eq!(body, b"whole"); + assert_eq!(origin.connections.load(Ordering::SeqCst), 2); + assert_healthy(&proxy).await; + proxy.stop().await; +} diff --git a/tests/integration/coalescing.rs b/tests/integration/coalescing.rs index fe59176..36589dd 100644 --- a/tests/integration/coalescing.rs +++ b/tests/integration/coalescing.rs @@ -550,6 +550,37 @@ async fn followers_of_a_stalled_leader_stop_waiting_after_the_upstream_timeout() server.stop().await; } +#[actix_web::test] +async fn followers_of_a_leader_whose_body_stalls_are_released_when_it_is_given_up() { + let (origin_url, connections) = origin_slowing_the_first_body(FirstBody::Stall).await; + let server = common::spawn_workers(&origin_url, WORKERS, |c| { + c.upstream_timeout_seconds = 1; + }); + let client = common::client(); + + let mut leader = raw_leader(&server).await; + let stalled = Instant::now(); + let followers = spawn_followers(&client, &server).await; + let mut rest = Vec::new(); + timeout(PATIENCE, leader.read_to_end(&mut rest)) + .await + .expect("the proxy kept the leader's connection open") + .ok(); + let given_up = stalled.elapsed(); + let followers = followers.await.unwrap(); + + assert!(given_up < Duration::from_millis(2500), "took {given_up:?}"); + assert!( + stalled.elapsed() < given_up + Duration::from_secs(1), + "followers took {:?} after the leader was given up at {given_up:?}", + stalled.elapsed() + ); + assert_own_whole_responses(&followers); + assert_eq!(connections.load(Ordering::SeqCst), 1 + FOLLOWERS); + assert_eq!(health(&client, &server).await["cache"]["coalesced"], 0); + server.stop().await; +} + #[actix_web::test] async fn requests_that_bypass_the_cache_do_not_coalesce() { let origin = origin_responding(cacheable("body")).await; diff --git a/tests/integration/main.rs b/tests/integration/main.rs index e07e50a..5e323ff 100644 --- a/tests/integration/main.rs +++ b/tests/integration/main.rs @@ -1,4 +1,5 @@ mod body_failure; +mod body_idle; mod cache; mod coalescing; mod common; From 1aecae78fc48117d3cbc866c6ed1e38feff175c9 Mon Sep 17 00:00:00 2001 From: jamie Date: Sun, 27 Sep 2026 18:05:26 +0100 Subject: [PATCH 2/2] fix(proxy): time out an idle background revalidation body --- README.md | 2 +- src/proxy.rs | 7 +++- tests/integration/body_idle.rs | 66 ++++++++++++++++++++++++++++++++++ 3 files changed, 73 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 8938011..a70c8ce 100644 --- a/README.md +++ b/README.md @@ -86,7 +86,7 @@ A `GET` that finds a stale response sends the origin `If-None-Match` from the st - On any other response, the proxy forwards it and stores it under the usual rules. - If the client sent its own `If-None-Match`, `If-Modified-Since`, `If-Match`, `If-Unmodified-Since` or `If-Range`, the proxy forwards those unchanged and adds none of its own. The origin's answer, including a `304`, goes to the client as `MISS`, and a 304 leaves the stored response as it was. -`Cache-Control: stale-while-revalidate=N` lets the proxy serve a response for `N` seconds after it goes stale while it revalidates in the background. At most one background revalidation runs for each stored response at a time. It uses the same header fields, forwarding headers and `--upstream-timeout-seconds` as any request to the origin, with the stored validators in place of the client's conditional headers. +`Cache-Control: stale-while-revalidate=N` lets the proxy serve a response for `N` seconds after it goes stale while it revalidates in the background. At most one background revalidation runs for each stored response at a time. It uses the same header fields, forwarding headers and `--upstream-timeout-seconds` as any request to the origin, with the stored validators in place of the client's conditional headers. If the origin sends nothing for that long during the body of its answer, the proxy stores nothing, and a later request can start another background revalidation. `Cache-Control: stale-if-error=N` lets the proxy serve a response for `N` seconds after it goes stale when the origin answers 500, 502, 503 or 504, refuses the connection or times out. diff --git a/src/proxy.rs b/src/proxy.rs index f31afd3..3f3b3fd 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -326,7 +326,9 @@ fn refresh_in_background( tokio::time::timeout(state.upstream_timeout, state.http_client.request(hyper_req)) .await; match upstream { - Ok(Ok(response)) => refresh(&req, &state, cache_key, &entry, response).await, + Ok(Ok(response)) => { + refresh(&req, &state, cache_key, &entry, response, target_uri).await + } Ok(Err(e)) => warn!("Background revalidation of {} failed: {}", target_uri, e), Err(_) => warn!("Background revalidation of {} timed out", target_uri), } @@ -342,6 +344,7 @@ async fn refresh( cache_key: PrimaryKey, entry: &StaleEntry, response: hyper::Response, + target_uri: Uri, ) { let (parts, body) = response.into_parts(); let head = Head::from_upstream(&parts); @@ -361,6 +364,8 @@ async fn refresh( return; }; let limit = usize::try_from(store.limit).unwrap_or(usize::MAX); + let body = origin_body(body, state.upstream_timeout, target_uri); + let body = StreamBody::new(body.map(|chunk| chunk.map(Frame::data))); match Limited::new(body, limit).collect().await { Ok(collected) => (store.finish)( head.status, diff --git a/tests/integration/body_idle.rs b/tests/integration/body_idle.rs index 45b674d..79ec61a 100644 --- a/tests/integration/body_idle.rs +++ b/tests/integration/body_idle.rs @@ -260,3 +260,69 @@ async fn stalled_cacheable_body_is_not_stored() { assert_healthy(&proxy).await; proxy.stop().await; } + +/// a response that goes stale after a second, which the proxy may then +/// serve stale while it revalidates it in the background. +fn stale_head(framing: &str) -> String { + format!( + "HTTP/1.1 200 OK\r\ncache-control: max-age=1, stale-while-revalidate=60\r\n\ + etag: \"v1\"\r\n{framing}\r\n\r\n" + ) +} + +async fn background_refreshes(proxy: &Running) -> u64 { + let client = common::client::>(); + let resp = timeout(PATIENCE, client.get(proxy.url("/health").parse().unwrap())) + .await + .expect("health check timed out") + .unwrap(); + assert_eq!(resp.status(), 200); + let health: serde_json::Value = + serde_json::from_slice(&common::body_bytes(resp.into_body()).await).unwrap(); + health["cache"]["background_refreshes"].as_u64().unwrap() +} + +#[actix_web::test] +async fn background_revalidation_whose_body_stalls_is_given_up() { + let first = Reply { + head: stale_head("content-length: 5"), + parts: vec![(Duration::ZERO, b"first".to_vec())], + stall: false, + }; + let stalled = Reply { + head: stale_head("content-length: 100"), + parts: vec![(Duration::ZERO, vec![b'a'; 40])], + stall: true, + }; + let not_modified = Reply { + head: "HTTP/1.1 304 Not Modified\r\netag: \"v1\"\r\n\r\n".to_owned(), + parts: Vec::new(), + stall: false, + }; + let mut origin = scripted_origin(vec![first, stalled, not_modified]).await; + let proxy = spawn(&origin); + + raw_get(&proxy, "/page").await; + // max-age has a granularity of one second + actix_web::rt::time::sleep(Duration::from_millis(1100)).await; + let started = Instant::now(); + let (response, _) = raw_get(&proxy, "/page").await; + let (head, _) = split(&response); + assert!(head.contains("x-shadowstep-cache: stale"), "{head}"); + origin_connection_closed(&mut origin, 1).await; + assert_given_up_in_time(started.elapsed()); + + // the refresh guard drops just after the origin connection closes + while background_refreshes(&proxy).await < 2 { + assert!(started.elapsed() < PATIENCE, "no second revalidation"); + raw_get(&proxy, "/page").await; + } + assert_eq!(background_refreshes(&proxy).await, 2); + assert!( + started.elapsed() < IDLE + MARGIN, + "took {:?}", + started.elapsed() + ); + assert_eq!(origin.connections.load(Ordering::SeqCst), 3); + proxy.stop().await; +}