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
17 changes: 15 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ Responses served from the cache carry `Age`. Proxied and asset responses carry `
- `MISS`: the origin's response.
- `REVALIDATED`: the origin answered a conditional request with `304 Not Modified`, and the client got the stored body with the 304's header fields.
- `STALE`: a stale stored response, served under `stale-while-revalidate` or `stale-if-error`.
- `COALESCED`: a response that a concurrent request for the same key stored while this request waited for it. See [Request coalescing](#request-coalescing).

`--cache-ttl-seconds` (default 300) caps the freshness lifetime of any entry, whatever the origin sent. `--cache-size-mb` (default 100) bounds origin responses and assets together, measured in bytes. The largest single entry is 8 MiB or the cache size, whichever is smaller. Larger bodies stream to the client without being stored. Setting either option to 0 turns caching off.

Expand All @@ -91,6 +92,16 @@ A `GET` that finds a stale response sends the origin `If-None-Match` from the st

A response with `must-revalidate`, `proxy-revalidate` or `s-maxage` is never served stale, whatever its stale windows. If the origin cannot be reached to revalidate it, the client gets `504 Gateway Timeout`. A request with `Cache-Control: max-age` never gets a stale response. `HEAD` requests are answered only from fresh responses.

### Request coalescing

When several requests miss on the same key at once, the first one, the leader, goes to the origin and the others wait for it. The key is the scheme, host, path and query, plus the request's values for the `Vary` field names of the latest stored response for that URL, if the cache has one. Requests for different known variants therefore do not wait for each other. Requests that find a stale response they must revalidate coalesce the same way, so the origin gets one conditional request.

- 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.
- 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.

Assets share the same cache. A stored asset is read from disk again when the file's size or modified time changes.

`/health` counts responses for proxied requests and assets together. Each response counts once:
Expand All @@ -99,14 +110,16 @@ Assets share the same cache. A stored asset is read from disk again when the fil
- `misses`: responses from the origin, including errors.
- `revalidations`: 304s that freshened a stored response, in the foreground or the background.
- `stale`: stale responses served under `stale-while-revalidate` or `stale-if-error`.
- `coalesced`: responses stored by a concurrent request that this request waited for.

`background_refreshes` counts background revalidations started. `hit_ratio` is the share of responses whose body came from the cache: hits, revalidations and stale responses. `items` and `bytes` describe the whole cache.
`background_refreshes` counts background revalidations started. `hit_ratio` is the share of responses whose body came from the cache: hits, revalidations, stale and coalesced responses. `items` and `bytes` describe the whole cache.

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.
- There is no request coalescing. Concurrent misses for the same URL all go to the origin.
- 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
41 changes: 41 additions & 0 deletions src/cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,15 @@ struct ResponseKey {
vary: Vec<Option<String>>,
}

/// the key that concurrent requests coalesce on: the primary key and the
/// request's values for the `Vary` field names the store knows for it, so
/// that requests for different known variants do not wait for each other.
#[derive(Clone, Debug, Hash, PartialEq, Eq)]
pub struct FlightKey {
primary: PrimaryKey,
vary: Vec<Option<String>>,
}

#[derive(Clone, Debug, Hash, PartialEq, Eq)]
enum Key {
Asset(PathBuf),
Expand Down Expand Up @@ -386,6 +395,23 @@ impl Store {
Some(Lookup::Stale(StaleEntry { response, key }))
}

/// the coalescing key for a request, or `None` when the store keeps
/// nothing, so that no request could be answered from a leader's
/// response.
pub fn flight_key(&self, primary: &PrimaryKey, request: &HeaderMap) -> Option<FlightKey> {
if !self.enabled() {
return None;
}
let vary = self
.index
.get(primary)
.map_or_else(Vec::new, |index| vary_values(&index.vary, request));
Some(FlightKey {
primary: primary.clone(),
vary,
})
}

/// claims the background revalidation of `entry`, or `None` when one is
/// already running.
pub fn start_refresh(&self, entry: &StaleEntry) -> Option<RefreshGuard> {
Expand Down Expand Up @@ -673,6 +699,21 @@ impl RequestPolicy {
pub fn may_serve_stale(&self) -> bool {
self.max_age.is_none()
}

/// whether the request may wait for a concurrent request for the same
/// key and then be answered from the cache. a request with credentials
/// does not, because the origin may answer it for that user alone, and
/// neither does one with its own preconditions, whose answer is for
/// that client.
pub fn may_coalesce(&self) -> bool {
self.may_serve && !self.has_credentials && !self.conditional
}

/// whether other requests may wait for this one's response to be
/// stored. only a GET's response is stored.
pub fn may_lead(&self) -> bool {
self.may_coalesce() && self.may_store
}
}

/// a response that the store may keep, and for how long.
Expand Down
154 changes: 154 additions & 0 deletions src/coalesce.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
//! request coalescing: while one request for a cache key is on its way to the
//! origin, other requests for the key wait for it and then read the cache.

use parking_lot::Mutex;
use std::collections::HashMap;
use std::hash::Hash;
use std::sync::Arc;
use tokio::sync::watch;

type InFlight<K> = Arc<Mutex<HashMap<K, watch::Receiver<()>>>>;

/// the keys with a request in flight, shared by every worker.
pub struct Flights<K> {
in_flight: InFlight<K>,
}

impl<K> Default for Flights<K> {
fn default() -> Self {
Flights {
in_flight: Arc::default(),
}
}
}

/// a request's part in the flight for its key.
pub enum Join<K: Hash + Eq> {
/// the request goes to the origin, and others wait until the guard drops
Lead(FlightGuard<K>),
/// another request is in flight for the key
Follow(Waiter),
/// nothing is in flight and the request may not lead
Alone,
}

impl<K: Hash + Eq + Clone> Flights<K> {
/// joins the flight for `key`, or starts one when there is none and
/// `may_lead` holds.
pub fn join(&self, key: K, may_lead: bool) -> Join<K> {
let mut in_flight = self.in_flight.lock();
if let Some(done) = in_flight.get(&key) {
return Join::Follow(Waiter(done.clone()));
}
if !may_lead {
return Join::Alone;
}
let (sender, done) = watch::channel(());
in_flight.insert(key.clone(), done);
Join::Lead(FlightGuard {
in_flight: self.in_flight.clone(),
key,
_sender: sender,
})
}
}

/// the leader's claim on a key. dropping it ends the flight and wakes the
/// followers.
pub struct FlightGuard<K: Hash + Eq> {
in_flight: InFlight<K>,
key: K,
/// dropped after `drop` has removed the key, so a woken follower never
/// finds the finished flight
_sender: watch::Sender<()>,
}

impl<K: Hash + Eq> Drop for FlightGuard<K> {
fn drop(&mut self) {
self.in_flight.lock().remove(&self.key);
}
}

/// a follower's view of a flight.
pub struct Waiter(watch::Receiver<()>);

impl Waiter {
/// returns once the leader's guard has dropped.
pub async fn wait(mut self) {
// nothing is ever sent, so `changed` returns only when the sender
// drops, or at once if it already has
while self.0.changed().await.is_ok() {}
}
}

#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;

fn lead(flights: &Flights<&'static str>, key: &'static str) -> FlightGuard<&'static str> {
match flights.join(key, true) {
Join::Lead(guard) => guard,
_ => panic!("{key} did not lead"),
}
}

fn follow(flights: &Flights<&'static str>, key: &'static str) -> Waiter {
match flights.join(key, true) {
Join::Follow(waiter) => waiter,
_ => panic!("{key} did not follow"),
}
}

async fn released(waiter: Waiter) -> bool {
tokio::time::timeout(Duration::from_millis(100), waiter.wait())
.await
.is_ok()
}

#[tokio::test]
async fn followers_wait_until_the_guard_drops() {
let flights = Flights::default();
let guard = lead(&flights, "a");

assert!(!released(follow(&flights, "a")).await);
let waiter = follow(&flights, "a");
drop(guard);
assert!(released(waiter).await);
}

#[tokio::test]
async fn a_follower_that_waits_after_the_guard_dropped_is_released() {
let flights = Flights::default();
let guard = lead(&flights, "a");
let waiter = follow(&flights, "a");
drop(guard);

assert!(released(waiter).await);
}

#[test]
fn the_next_request_leads_once_the_flight_ends() {
let flights = Flights::default();
drop(lead(&flights, "a"));

lead(&flights, "a");
}

#[test]
fn keys_have_separate_flights() {
let flights = Flights::default();
let _a = lead(&flights, "a");

lead(&flights, "b");
}

#[test]
fn a_request_that_may_not_lead_goes_alone() {
let flights = Flights::default();

assert!(matches!(flights.join("a", false), Join::Alone));
let _guard = lead(&flights, "a");
assert!(matches!(flights.join("a", false), Join::Follow(_)));
}
}
20 changes: 16 additions & 4 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ pub mod tls;

mod assets;
mod cache;
mod coalesce;
mod forwarded;
mod proxy;

Expand All @@ -32,19 +33,21 @@ use std::sync::Arc;
use std::time::Duration;
use url::Url;

use crate::cache::Store;
use crate::cache::{FlightKey, Store};
use crate::coalesce::Flights;
use crate::config::Config;

/// how the cache answered requests. each proxied or asset response counts
/// once, as a hit, a miss, a revalidation or a stale serve. background
/// refreshes count the background revalidations started.
/// once, as a hit, a miss, a revalidation, a stale serve or a coalesced
/// serve. background refreshes count the background revalidations started.
#[derive(Default)]
pub struct CacheStats {
hits: AtomicU64,
misses: AtomicU64,
revalidations: AtomicU64,
stale: AtomicU64,
background_refreshes: AtomicU64,
coalesced: AtomicU64,
}

impl CacheStats {
Expand All @@ -69,12 +72,18 @@ impl CacheStats {
fn background_refresh(&self) {
self.background_refreshes.fetch_add(1, Ordering::Relaxed);
}

/// a response stored by a concurrent request that this one waited for
fn coalesced(&self) {
self.coalesced.fetch_add(1, Ordering::Relaxed);
}
}

/// application state, including cache
pub struct AppState {
cache_stats: CacheStats,
cache: Store,
flights: Flights<FlightKey>,
http_client: Client<HttpsConnector<HttpConnector>, proxy::UpstreamBody>,
upstream_base_url: Url,
asset_path: PathBuf,
Expand All @@ -88,8 +97,9 @@ async fn health_check(state: web::Data<AppState>) -> impl Responder {
let misses = stats.misses.load(Ordering::Relaxed);
let revalidations = stats.revalidations.load(Ordering::Relaxed);
let stale = stats.stale.load(Ordering::Relaxed);
let coalesced = stats.coalesced.load(Ordering::Relaxed);
// the share of responses whose body came from the cache
let from_cache = hits + revalidations + stale;
let from_cache = hits + revalidations + stale + coalesced;
let total = from_cache + misses;
HttpResponse::Ok().json(serde_json::json!({
"status": "ok",
Expand All @@ -98,6 +108,7 @@ async fn health_check(state: web::Data<AppState>) -> impl Responder {
"misses": misses,
"revalidations": revalidations,
"stale": stale,
"coalesced": coalesced,
"background_refreshes": stats.background_refreshes.load(Ordering::Relaxed),
"items": state.cache.entry_count(),
"bytes": state.cache.weighted_size(),
Expand Down Expand Up @@ -142,6 +153,7 @@ pub fn build_state(config: &Config) -> io::Result<web::Data<AppState>> {
config.cache_size_mb,
Duration::from_secs(config.cache_ttl_seconds),
),
flights: Flights::default(),
http_client,
upstream_base_url,
asset_path: config.asset_path.clone(),
Expand Down
Loading
Loading