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
15 changes: 15 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ bytes = "1"
httpdate = "1.0.3"
futures-util = { version = "0.3", default-features = false, features = ["std"] }
parking_lot = "0.12"
prometheus = { version = "0.14", default-features = false }

[profile.release]
lto = true
Expand Down
41 changes: 41 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ A small caching reverse proxy written in Rust (actix-web 4, hyper 1, rustls 0.23
- Every other path is forwarded to the origin with its method, path, query, headers and body. See [Proxying](#proxying).
- Cacheable origin responses are stored and served from memory. See [Caching](#caching).
- `GET /health` returns `{"status":"ok","cache":{...}}` with `hits`, `misses`, `hit_ratio`, `items` and `bytes` for the cache.
- `GET /metrics` returns Prometheus metrics. See [Metrics](#metrics).
- Responses are compressed according to the request's `Accept-Encoding`.
- HTTPS is served when both a certificate and a key are given. See [TLS](#tls).

Expand Down Expand Up @@ -114,6 +115,8 @@ Assets share the same cache. A stored asset is read from disk again when the fil

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

`/health` reads the same counters as `/metrics`, so the two always agree. `misses` is the sum of the `MISS` series and the proxy's `BYPASS` series of `shadowstep_requests_total`.

Known limits:

- A response with `Cache-Control: no-cache` is not stored, although RFC 9111 allows storing it and revalidating it on every use.
Expand All @@ -122,6 +125,40 @@ Known limits:
- 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.

## Metrics

`GET /metrics` returns metrics in the Prometheus text format 0.0.4, with `Content-Type: text/plain; version=0.0.4; charset=utf-8`. Each process has its own counters, which start at 0.

| Metric | Type | Labels | Meaning |
|---|---|---|---|
| `shadowstep_requests_total` | counter | `route`, `cache` | Proxied and asset requests. `route` is `proxy` or `asset`. `cache` is the response's `X-Shadowstep-Cache` value, or `BYPASS` for a request that the cache could neither answer nor store |
| `shadowstep_responses_total` | counter | `route`, `status_class` | Responses to proxied and asset requests, by status class: `1xx` to `5xx` |
| `shadowstep_origin_requests_total` | counter | `kind`, `outcome` | Requests to the origin. `kind` is `foreground` for a client's request or `background` for a `stale-while-revalidate` revalidation. `outcome` is the status class of the origin's response, `error` when the connection or request failed, or `timeout` when no response headers arrived within `--upstream-timeout-seconds` |
| `shadowstep_origin_body_idle_timeouts_total` | counter | `kind` | Origin response bodies given up because the origin sent nothing for `--upstream-timeout-seconds` |
| `shadowstep_origin_response_seconds` | histogram | none | Time from sending an origin request to receiving its response headers, for both kinds. Buckets run from 5 ms to 30 s |
| `shadowstep_cache_revalidations_total` | counter | none | 304 responses that freshened a stored response, in the foreground or the background |
| `shadowstep_cache_background_refreshes_total` | counter | none | Background revalidations started |
| `shadowstep_cache_bytes` | gauge | none | Bytes held in the cache, read at scrape time |
| `shadowstep_cache_entries` | gauge | none | Entries in the cache, read at scrape time |

Notes on the counts:

- `BYPASS` covers, among others, methods other than `GET` and `HEAD`, requests with `Cache-Control: no-store`, requests with a method-override header, and asset requests that found no file. The response itself still carries `X-Shadowstep-Cache: MISS` for a proxied bypass.
- A proxied error response, such as a 502 or a 504, counts as `MISS`.
- Stale and coalesced responses are the `STALE` and `COALESCED` series of `shadowstep_requests_total`, so they have no metric of their own.
- An origin request counts once in `shadowstep_origin_requests_total`, when its headers arrive or it fails. A body that later goes idle counts again in `shadowstep_origin_body_idle_timeouts_total` but not in `shadowstep_origin_requests_total`, so the latter's total stays the number of origin requests.
- The histogram observes only origin requests whose headers arrived. Its `_count` equals the sum of the status-class outcomes of `shadowstep_origin_requests_total`.
- No label holds a path, host or query, so the number of series is fixed.
- Requests to `/health` and `/metrics` are not counted.

### Exposure

By default `/metrics` is served on the HTTP and HTTPS listeners, as `/health` is. On a public edge node anyone can then read it. The metrics hold aggregate counts, a latency histogram and the cache size, and `/health` already shows most of these counts to anyone.

To keep `/metrics` off the public listeners, set `--metrics-addr` to an address that only the scraper can reach, for example `127.0.0.1:9090` or a pod IP port that no Service exposes. shadowstep then serves `/metrics` only on that listener, which answers 404 to every other path and proxies nothing.

`/health`, and `/metrics` without `--metrics-addr`, shadow the same paths on the origin: a client cannot reach the origin's own `/health` or `/metrics` through shadowstep. With `--metrics-addr` set, `/metrics` on the proxy listeners goes to the origin like any other path.

## Install

Each release tag `vX.Y.Z` publishes a multi-platform image for `linux/amd64` and `linux/arm64` to `ghcr.io/jamiehdev/shadowstep`, tagged `X.Y.Z`, `X.Y` and `latest`. A pre-release tag such as `v2.1.0-rc.1` publishes only `2.1.0-rc.1`.
Expand Down Expand Up @@ -184,6 +221,7 @@ Each option can be set with a flag or an environment variable. The flag wins if
| `--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, and the longest the origin may send nothing during a response body |
| `--metrics-addr` | `METRICS_ADDR` | none | Address for a listener that serves only `/metrics`. When set, the proxy listeners forward `/metrics` to the origin. When unset, they serve `/metrics` themselves. See [Metrics](#metrics) |

`cargo run -- --help` prints the same list.

Expand Down Expand Up @@ -259,6 +297,8 @@ The Service maps port 80 to 8080 and 443 to 8443. Readiness and liveness probes

`CACHE_SIZE_MB` is set to 100 against a 256Mi memory limit. Change the two together.

The Deployment has no `prometheus.io/scrape` annotations. They are a convention of some Prometheus scrape configs, not a Kubernetes or Prometheus standard, and the Prometheus Operator ignores them in favour of a `ServiceMonitor` or `PodMonitor`. Configure scraping of `/metrics` on port 8080 in whichever way the cluster's Prometheus expects. To keep `/metrics` off the LoadBalancer, set `METRICS_ADDR` to `0.0.0.0:9090`, add a container port for it and leave it out of the Service.

The Deployment runs two replicas, and each has its own cache. `X-Forwarded-For` holds whatever source address reaches the pod. With the Service's default `externalTrafficPolicy: Cluster`, that is often a node address rather than the client's.

## Tests
Expand All @@ -275,6 +315,7 @@ Unit tests in `src/assets.rs` cover asset path traversal, `src/cache.rs` covers
- `forwarded.rs`: removal and replacement of forwarding headers
- `cache.rs`: origin response caching
- `revalidation.rs`: conditional requests, `stale-while-revalidate`, `stale-if-error` and `must-revalidate`
- `metrics.rs`: each `/metrics` series, its agreement with `/health`, and `--metrics-addr`

CI also runs:

Expand Down
21 changes: 15 additions & 6 deletions src/assets.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ use std::path::{Component, Path, PathBuf};
use std::sync::Arc;

use crate::cache::Asset;
use crate::metrics::{self, Route};
use crate::AppState;

/// resolves `requested` to a file inside `root`, or `None` if it is missing
Expand Down Expand Up @@ -34,14 +35,23 @@ pub async fn serve_asset(
state: web::Data<AppState>,
req: HttpRequest,
) -> impl Responder {
let filename = path.into_inner();
let response = asset_response(path.into_inner(), &state, &req).await;
// a response without `X-Shadowstep-Cache` found no file, so the cache
// was never consulted
let cache = metrics::cache_label(&response).unwrap_or(metrics::BYPASS);
state
.metrics
.observe(Route::Asset, cache, response.status());
response
}

async fn asset_response(filename: String, state: &AppState, req: &HttpRequest) -> HttpResponse {
let Some(path) = resolve_asset(&state.asset_path, &filename).await else {
warn!("Asset not found: {}", filename);
return HttpResponse::NotFound().body("not found");
};

let (asset, status) = match load(&state, path).await {
let (asset, status) = match load(state, path).await {
Ok(loaded) => loaded,
Err(e) => {
warn!("Asset not found: {} - Error: {}", filename, e);
Expand Down Expand Up @@ -82,18 +92,16 @@ async fn load(state: &AppState, path: PathBuf) -> std::io::Result<(Arc<Asset>, &

if let Some(asset) = state.cache.asset(&path) {
if asset.modified == modified && asset.len == metadata.len() {
state.cache_stats.hit();
return Ok((asset, "HIT"));
return Ok((asset, metrics::HIT));
}
}

debug!("Reading asset from {:?}", path);
let content = Bytes::from(tokio::fs::read(&path).await?);
let etag = format!("\"{}\"", &hex::encode(Sha256::digest(&content))[..32]);
state.cache_stats.miss();
Ok((
state.cache.insert_asset(path, content, etag, modified),
"MISS",
metrics::MISS,
))
}

Expand Down Expand Up @@ -126,6 +134,7 @@ mod tests {
tls_key_path: None,
tls_listen_addr: "127.0.0.1:0".into(),
upstream_timeout_seconds: 30,
metrics_addr: None,
})
.unwrap()
}
Expand Down
6 changes: 6 additions & 0 deletions src/cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -693,6 +693,12 @@ impl RequestPolicy {
}
}

/// whether the cache can neither answer the request nor store its
/// response, as for a POST or a request with `no-store`.
pub fn bypasses_cache(&self) -> bool {
!self.may_serve && !self.may_store
}

/// whether a stale response may answer the request. a request max-age
/// asks for a response no older than that, so it rules out stale ones
/// (RFC 9111 section 5.2.1.1).
Expand Down
6 changes: 6 additions & 0 deletions src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,12 @@ pub struct Config {
/// body before the proxy closes the client's connection.
#[clap(long, env = "UPSTREAM_TIMEOUT_SECONDS", default_value_t = 30)]
pub upstream_timeout_seconds: u64,

/// address for a listener that serves only `/metrics`. when it is set,
/// the proxy listeners forward `/metrics` to the origin like any other
/// path. when it is unset, they serve `/metrics` as they serve `/health`.
#[clap(long, env = "METRICS_ADDR")]
pub metrics_addr: Option<String>,
}

impl Config {
Expand Down
105 changes: 45 additions & 60 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ mod assets;
mod cache;
mod coalesce;
mod forwarded;
mod metrics;
mod proxy;

use actix_web::body::MessageBody;
Expand All @@ -28,60 +29,20 @@ use log::info;
use std::io;
use std::net::TcpListener;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use url::Url;

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, 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 {
fn hit(&self) {
self.hits.fetch_add(1, Ordering::Relaxed);
}

fn miss(&self) {
self.misses.fetch_add(1, Ordering::Relaxed);
}

/// a 304 that freshened a stored response, in the foreground or the
/// background
fn revalidation(&self) {
self.revalidations.fetch_add(1, Ordering::Relaxed);
}

fn stale(&self) {
self.stale.fetch_add(1, Ordering::Relaxed);
}

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);
}
}
use crate::metrics::Metrics;

/// application state, including cache
pub struct AppState {
cache_stats: CacheStats,
metrics: Metrics,
/// whether `/metrics` is served on the proxy listeners
metrics_on_proxy_listeners: bool,
cache: Store,
flights: Flights<FlightKey>,
http_client: Client<HttpsConnector<HttpConnector>, proxy::UpstreamBody>,
Expand All @@ -92,24 +53,19 @@ pub struct AppState {

#[get("/health")]
async fn health_check(state: web::Data<AppState>) -> impl Responder {
let stats = &state.cache_stats;
let hits = stats.hits.load(Ordering::Relaxed);
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);
let counts = state.metrics.cache_counts();
// the share of responses whose body came from the cache
let from_cache = hits + revalidations + stale + coalesced;
let total = from_cache + misses;
let from_cache = counts.hits + counts.revalidations + counts.stale + counts.coalesced;
let total = from_cache + counts.misses;
HttpResponse::Ok().json(serde_json::json!({
"status": "ok",
"cache": {
"hits": hits,
"misses": misses,
"revalidations": revalidations,
"stale": stale,
"coalesced": coalesced,
"background_refreshes": stats.background_refreshes.load(Ordering::Relaxed),
"hits": counts.hits,
"misses": counts.misses,
"revalidations": counts.revalidations,
"stale": counts.stale,
"coalesced": counts.coalesced,
"background_refreshes": counts.background_refreshes,
"items": state.cache.entry_count(),
"bytes": state.cache.weighted_size(),
"hit_ratio": if total > 0 {
Expand All @@ -121,6 +77,13 @@ async fn health_check(state: web::Data<AppState>) -> impl Responder {
}))
}

#[get("/metrics")]
async fn metrics_endpoint(state: web::Data<AppState>) -> impl Responder {
HttpResponse::Ok()
.content_type(metrics::CONTENT_TYPE)
.body(state.metrics.encode(&state.cache))
}

/// builds the shared application state from `config`, creating the asset
/// directory if it does not exist.
pub fn build_state(config: &Config) -> io::Result<web::Data<AppState>> {
Expand Down Expand Up @@ -148,7 +111,8 @@ pub fn build_state(config: &Config) -> io::Result<web::Data<AppState>> {
info!("Serving assets from: {:?}", config.asset_path);

Ok(web::Data::new(AppState {
cache_stats: CacheStats::default(),
metrics: Metrics::new(),
metrics_on_proxy_listeners: config.metrics_addr.is_none(),
cache: Store::new(
config.cache_size_mb,
Duration::from_secs(config.cache_ttl_seconds),
Expand All @@ -161,7 +125,8 @@ pub fn build_state(config: &Config) -> io::Result<web::Data<AppState>> {
}))
}

/// the actix application: `/health`, `/assets/*` and a catch-all proxy route.
/// the actix application: `/health`, `/metrics` unless it has a listener of
/// its own, `/assets/*` and a catch-all proxy route.
pub fn app(
state: web::Data<AppState>,
) -> App<
Expand All @@ -173,11 +138,17 @@ pub fn app(
InitError = (),
>,
> {
let metrics_here = state.metrics_on_proxy_listeners;
App::new()
.app_data(state)
.wrap(Compress::default())
.wrap(Logger::new("%r %s %b %D ms"))
.service(health_check)
.configure(|cfg| {
if metrics_here {
cfg.service(metrics_endpoint);
}
})
.service(assets::serve_asset)
.route("/{path:.*}", web::to(proxy::forward_to_upstream))
}
Expand All @@ -201,3 +172,17 @@ pub fn run(

Ok(server.run())
}

/// a server that answers `GET /metrics` on `listener` and 404 on every other
/// path, for `--metrics-addr`.
pub fn run_metrics(state: web::Data<AppState>, listener: TcpListener) -> io::Result<Server> {
let server = HttpServer::new(move || {
App::new()
.app_data(state.clone())
.wrap(Compress::default())
.service(metrics_endpoint)
})
.workers(1)
.listen(listener)?;
Ok(server.run())
}
Loading
Loading