Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 4 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down Expand Up @@ -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.

Expand All @@ -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.

Expand All @@ -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.

Expand Down Expand Up @@ -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.

Expand Down
4 changes: 3 additions & 1 deletion src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
Expand Down
102 changes: 88 additions & 14 deletions src/proxy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -322,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),
}
Expand All @@ -338,6 +344,7 @@ async fn refresh(
cache_key: PrimaryKey,
entry: &StaleEntry,
response: hyper::Response<Incoming>,
target_uri: Uri,
) {
let (parts, body) = response.into_parts();
let head = Head::from_upstream(&parts);
Expand All @@ -357,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,
Expand Down Expand Up @@ -393,7 +402,7 @@ struct Exchange {

impl Exchange {
/// the client's response once the origin has answered.
fn answered(self, response: hyper::Response<Incoming>) -> HttpResponse {
fn answered(self, response: hyper::Response<Incoming>, target_uri: Uri) -> HttpResponse {
let (parts, body) = response.into_parts();
let head = Head::from_upstream(&parts);
if let Some(stale) = &self.stale {
Expand All @@ -413,6 +422,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)))
}

Expand Down Expand Up @@ -696,7 +706,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<Storing>) -> HttpResponse {
fn client_response<S>(head: Head, body: S, store: Option<Storing>) -> HttpResponse
where
S: Stream<Item = io::Result<Bytes>> + Unpin + 'static,
{
let mut builder = HttpResponse::build(head.status);

let fields = forwardable_fields(&head);
Expand All @@ -706,14 +719,6 @@ fn client_response(head: Head, body: Incoming, store: Option<Storing>) -> 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)
Expand Down Expand Up @@ -770,6 +775,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<impl Stream<Item = io::Result<Bytes>> + 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<S> {
inner: S,
idle: Duration,
sleep: Pin<Box<Sleep>>,
/// 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<S> Stream for IdleTimeout<S>
where
S: Stream<Item = io::Result<Bytes>> + Unpin,
{
type Item = S::Item;

fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
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.
Expand Down Expand Up @@ -805,9 +879,9 @@ impl<S> Tee<S> {
}
}

impl<S> Stream for Tee<S>
impl<S, E> Stream for Tee<S>
where
S: Stream<Item = Result<Bytes, hyper::Error>> + Unpin,
S: Stream<Item = Result<Bytes, E>> + Unpin,
{
type Item = S::Item;

Expand Down
Loading
Loading