diff --git a/Cargo.lock b/Cargo.lock index b055f27..bf9324e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1619,6 +1619,20 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "prometheus" +version = "0.14.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3ca5326d8d0b950a9acd87e6a3f94745394f62e4dae1b1ee22b2bc0c394af43a" +dependencies = [ + "cfg-if", + "fnv", + "lazy_static", + "memchr", + "parking_lot", + "thiserror", +] + [[package]] name = "quote" version = "1.0.40" @@ -1939,6 +1953,7 @@ dependencies = [ "mime_guess", "moka", "parking_lot", + "prometheus", "rcgen", "rustls", "serde_json", diff --git a/Cargo.toml b/Cargo.toml index 5e4051d..2c095c9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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 diff --git a/README.md b/README.md index f914feb..e225058 100644 --- a/README.md +++ b/README.md @@ -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). @@ -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. @@ -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`. @@ -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. @@ -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 @@ -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: diff --git a/src/assets.rs b/src/assets.rs index 33c3cf7..91bd87d 100644 --- a/src/assets.rs +++ b/src/assets.rs @@ -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 @@ -34,14 +35,23 @@ pub async fn serve_asset( state: web::Data, 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); @@ -82,18 +92,16 @@ async fn load(state: &AppState, path: PathBuf) -> std::io::Result<(Arc, & 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, )) } @@ -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() } diff --git a/src/cache.rs b/src/cache.rs index a5a16c9..f850e52 100644 --- a/src/cache.rs +++ b/src/cache.rs @@ -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). diff --git a/src/config.rs b/src/config.rs index 8686b9d..da983eb 100644 --- a/src/config.rs +++ b/src/config.rs @@ -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, } impl Config { diff --git a/src/lib.rs b/src/lib.rs index 0950741..afec100 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -14,6 +14,7 @@ mod assets; mod cache; mod coalesce; mod forwarded; +mod metrics; mod proxy; use actix_web::body::MessageBody; @@ -28,7 +29,6 @@ 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; @@ -36,52 +36,13 @@ 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, http_client: Client, proxy::UpstreamBody>, @@ -92,24 +53,19 @@ pub struct AppState { #[get("/health")] async fn health_check(state: web::Data) -> 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 { @@ -121,6 +77,13 @@ async fn health_check(state: web::Data) -> impl Responder { })) } +#[get("/metrics")] +async fn metrics_endpoint(state: web::Data) -> 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> { @@ -148,7 +111,8 @@ pub fn build_state(config: &Config) -> io::Result> { 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), @@ -161,7 +125,8 @@ pub fn build_state(config: &Config) -> io::Result> { })) } -/// 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, ) -> App< @@ -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)) } @@ -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, listener: TcpListener) -> io::Result { + let server = HttpServer::new(move || { + App::new() + .app_data(state.clone()) + .wrap(Compress::default()) + .service(metrics_endpoint) + }) + .workers(1) + .listen(listener)?; + Ok(server.run()) +} diff --git a/src/main.rs b/src/main.rs index 44798da..a64fd52 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,8 +1,9 @@ use log::info; use shadowstep::config::Config; use shadowstep::tls::bind_listeners; -use shadowstep::{build_state, run}; +use shadowstep::{build_state, run, run_metrics}; use std::io; +use std::net::TcpListener; use std::num::NonZeroUsize; #[actix_web::main] @@ -19,6 +20,13 @@ async fn main() -> io::Result<()> { ); let (http, tls) = bind_listeners(&config)?; + let proxy = run(state.clone(), http, tls, num_workers)?; - run(state, http, tls, num_workers)?.await + let Some(metrics_addr) = &config.metrics_addr else { + return proxy.await; + }; + info!("Serving /metrics on {}", metrics_addr); + let metrics = run_metrics(state, TcpListener::bind(metrics_addr)?)?; + // each server stops on SIGINT or SIGTERM, so both end together + tokio::try_join!(proxy, metrics).map(|_| ()) } diff --git a/src/metrics.rs b/src/metrics.rs new file mode 100644 index 0000000..ab8b365 --- /dev/null +++ b/src/metrics.rs @@ -0,0 +1,345 @@ +use actix_web::http::StatusCode; +use actix_web::HttpResponse; +use prometheus::core::Collector; +use prometheus::{ + Encoder, Histogram, HistogramOpts, IntCounter, IntCounterVec, IntGauge, Opts, Registry, + TextEncoder, +}; +use std::time::Duration; + +use crate::cache::Store; + +/// the `Content-Type` of the Prometheus text exposition format 0.0.4. +pub const CONTENT_TYPE: &str = "text/plain; version=0.0.4; charset=utf-8"; + +/// the `X-Shadowstep-Cache` values, plus `BYPASS` for a request that could +/// neither be answered from the cache nor stored in it. +pub const HIT: &str = "HIT"; +pub const MISS: &str = "MISS"; +pub const REVALIDATED: &str = "REVALIDATED"; +pub const STALE: &str = "STALE"; +pub const COALESCED: &str = "COALESCED"; +pub const BYPASS: &str = "BYPASS"; + +const CACHE_LABELS: [&str; 6] = [HIT, MISS, REVALIDATED, STALE, COALESCED, BYPASS]; +const ASSET_CACHE_LABELS: [&str; 3] = [HIT, MISS, BYPASS]; +const STATUS_CLASSES: [&str; 5] = ["1xx", "2xx", "3xx", "4xx", "5xx"]; + +/// upper bounds of the origin response time buckets, in seconds. 30 s is +/// the default upstream timeout. +const ORIGIN_BUCKETS: [f64; 12] = [ + 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, +]; + +/// the handler a request reached. paths are left out of labels because +/// each one would be a new series. +#[derive(Clone, Copy)] +pub enum Route { + Proxy, + Asset, +} + +impl Route { + fn label(self) -> &'static str { + match self { + Route::Proxy => "proxy", + Route::Asset => "asset", + } + } +} + +/// whether an origin request answers a client or revalidates a stale +/// response in the background. +#[derive(Clone, Copy)] +pub enum OriginKind { + Foreground, + Background, +} + +impl OriginKind { + fn label(self) -> &'static str { + match self { + OriginKind::Foreground => "foreground", + OriginKind::Background => "background", + } + } +} + +/// why an origin request ended without response headers. +#[derive(Clone, Copy)] +pub enum OriginFailure { + /// the connection or the request failed + Error, + /// no response headers within the upstream timeout + Timeout, +} + +/// the counts `/health` reports, read from the same counters as `/metrics`. +pub struct CacheCounts { + pub hits: u64, + pub misses: u64, + pub revalidations: u64, + pub stale: u64, + pub coalesced: u64, + pub background_refreshes: u64, +} + +/// the proxy's Prometheus metrics, in a registry of their own so that each +/// `AppState` counts only its own requests. +pub struct Metrics { + registry: Registry, + requests: IntCounterVec, + responses: IntCounterVec, + origin_requests: IntCounterVec, + origin_response_seconds: Histogram, + origin_body_idle_timeouts: IntCounterVec, + revalidations: IntCounter, + background_refreshes: IntCounter, + cache_bytes: IntGauge, + cache_entries: IntGauge, +} + +/// registers `metric` in `registry`. every name, label name and bucket +/// list here is a constant, so an error means a programming mistake, which +/// any test that builds an `AppState` catches. +fn register( + registry: &Registry, + metric: prometheus::Result, +) -> M { + let metric = metric.expect("invalid metric definition"); + registry + .register(Box::new(metric.clone())) + .expect("duplicate metric name"); + metric +} + +fn counter_vec(name: &str, help: &str, labels: &[&str]) -> prometheus::Result { + IntCounterVec::new(Opts::new(name, help), labels) +} + +const REQUESTS_HELP: &str = "Proxied and asset requests, by route and X-Shadowstep-Cache value, \ + or BYPASS for requests the cache could neither answer nor store."; +const RESPONSES_HELP: &str = "Responses to proxied and asset requests, by route and status class."; +const ORIGIN_REQUESTS_HELP: &str = "Requests to the origin, by kind and by the status class of \ + the response, or error or timeout when no response headers arrived."; +const ORIGIN_RESPONSE_SECONDS_HELP: &str = + "Time from sending an origin request to receiving its response headers."; +const BODY_IDLE_HELP: &str = + "Origin response bodies given up after the origin sent nothing for the upstream timeout."; +const REVALIDATIONS_HELP: &str = "304 responses from the origin that freshened a stored \ + response, in the foreground or the background."; +const BACKGROUND_REFRESHES_HELP: &str = + "Background revalidations started under stale-while-revalidate."; +const CACHE_BYTES_HELP: &str = "Bytes held in the cache, for origin responses and assets."; +const CACHE_ENTRIES_HELP: &str = "Entries in the cache, for origin responses and assets."; + +impl Metrics { + pub fn new() -> Self { + let registry = Registry::new(); + let r = ®istry; + let histogram = HistogramOpts::new( + "shadowstep_origin_response_seconds", + ORIGIN_RESPONSE_SECONDS_HELP, + ) + .buckets(ORIGIN_BUCKETS.to_vec()); + let metrics = Metrics { + requests: register( + r, + counter_vec( + "shadowstep_requests_total", + REQUESTS_HELP, + &["route", "cache"], + ), + ), + responses: register( + r, + counter_vec( + "shadowstep_responses_total", + RESPONSES_HELP, + &["route", "status_class"], + ), + ), + origin_requests: register( + r, + counter_vec( + "shadowstep_origin_requests_total", + ORIGIN_REQUESTS_HELP, + &["kind", "outcome"], + ), + ), + origin_response_seconds: register(r, Histogram::with_opts(histogram)), + origin_body_idle_timeouts: register( + r, + counter_vec( + "shadowstep_origin_body_idle_timeouts_total", + BODY_IDLE_HELP, + &["kind"], + ), + ), + revalidations: register( + r, + IntCounter::new("shadowstep_cache_revalidations_total", REVALIDATIONS_HELP), + ), + background_refreshes: register( + r, + IntCounter::new( + "shadowstep_cache_background_refreshes_total", + BACKGROUND_REFRESHES_HELP, + ), + ), + cache_bytes: register(r, IntGauge::new("shadowstep_cache_bytes", CACHE_BYTES_HELP)), + cache_entries: register( + r, + IntGauge::new("shadowstep_cache_entries", CACHE_ENTRIES_HELP), + ), + registry, + }; + metrics.initialise_series(); + metrics + } + + /// creates each expected series at 0, so that `rate()` sees the first + /// increment of a series. + fn initialise_series(&self) { + for cache in CACHE_LABELS { + self.requests.with_label_values(&["proxy", cache]); + } + for cache in ASSET_CACHE_LABELS { + self.requests.with_label_values(&["asset", cache]); + } + for kind in [OriginKind::Foreground, OriginKind::Background] { + for outcome in STATUS_CLASSES.iter().chain(&["error", "timeout"]) { + self.origin_requests + .with_label_values(&[kind.label(), outcome]); + } + self.origin_body_idle_timeouts + .with_label_values(&[kind.label()]); + } + for route in [Route::Proxy, Route::Asset] { + for class in STATUS_CLASSES { + self.responses.with_label_values(&[route.label(), class]); + } + } + } + + /// counts a response to a proxied or asset request. + pub fn observe(&self, route: Route, cache: &str, status: StatusCode) { + self.requests + .with_label_values(&[route.label(), cache]) + .inc(); + self.responses + .with_label_values(&[route.label(), status_class(status.as_u16())]) + .inc(); + } + + /// counts an origin request whose response headers arrived after + /// `elapsed`. + pub fn origin_responded(&self, kind: OriginKind, status: u16, elapsed: Duration) { + self.origin_requests + .with_label_values(&[kind.label(), status_class(status)]) + .inc(); + self.origin_response_seconds.observe(elapsed.as_secs_f64()); + } + + /// counts an origin request that ended without response headers. + pub fn origin_failed(&self, kind: OriginKind, failure: OriginFailure) { + let outcome = match failure { + OriginFailure::Error => "error", + OriginFailure::Timeout => "timeout", + }; + self.origin_requests + .with_label_values(&[kind.label(), outcome]) + .inc(); + } + + /// the counter of origin bodies of `kind` given up for going idle. + pub fn body_idle_timeouts(&self, kind: OriginKind) -> IntCounter { + self.origin_body_idle_timeouts + .with_label_values(&[kind.label()]) + } + + /// a 304 that freshened a stored response, in the foreground or the + /// background + pub fn revalidation(&self) { + self.revalidations.inc(); + } + + pub fn background_refresh(&self) { + self.background_refreshes.inc(); + } + + fn requests(&self, route: Route, cache: &str) -> u64 { + self.requests + .with_label_values(&[route.label(), cache]) + .get() + } + + fn both_routes(&self, cache: &str) -> u64 { + self.requests(Route::Proxy, cache) + self.requests(Route::Asset, cache) + } + + /// the counts for `/health`. its misses are the responses that came + /// from the origin, bypass or not, and asset misses. an asset request + /// that found no file counts as a bypass here and not in `/health`. + pub fn cache_counts(&self) -> CacheCounts { + CacheCounts { + hits: self.both_routes(HIT), + misses: self.both_routes(MISS) + self.requests(Route::Proxy, BYPASS), + revalidations: self.revalidations.get(), + stale: self.both_routes(STALE), + coalesced: self.both_routes(COALESCED), + background_refreshes: self.background_refreshes.get(), + } + } + + /// the text exposition, with the cache gauges read from `store` now. + pub fn encode(&self, store: &Store) -> Vec { + self.cache_bytes.set(gauge_value(store.weighted_size())); + self.cache_entries.set(gauge_value(store.entry_count())); + let mut buffer = Vec::new(); + if let Err(e) = TextEncoder::new().encode(&self.registry.gather(), &mut buffer) { + log::error!("Failed to encode metrics: {}", e); + } + buffer + } +} + +fn gauge_value(value: u64) -> i64 { + i64::try_from(value).unwrap_or(i64::MAX) +} + +/// `2xx` for 200 to 299 and so on. an origin may send a status up to 999, +/// which counts as `other`. +fn status_class(status: u16) -> &'static str { + match status / 100 { + 1 => "1xx", + 2 => "2xx", + 3 => "3xx", + 4 => "4xx", + 5 => "5xx", + _ => "other", + } +} + +/// the response's `X-Shadowstep-Cache` value, if it is one the proxy sets. +pub fn cache_label(response: &HttpResponse) -> Option<&'static str> { + let value = response.headers().get("x-shadowstep-cache")?; + CACHE_LABELS + .into_iter() + .find(|label| value.as_bytes() == label.as_bytes()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn status_classes_cover_each_hundred() { + assert_eq!(status_class(100), "1xx"); + assert_eq!(status_class(299), "2xx"); + assert_eq!(status_class(304), "3xx"); + assert_eq!(status_class(404), "4xx"); + assert_eq!(status_class(599), "5xx"); + assert_eq!(status_class(600), "other"); + } +} diff --git a/src/proxy.rs b/src/proxy.rs index 3f3b3fd..a48fc7d 100644 --- a/src/proxy.rs +++ b/src/proxy.rs @@ -10,6 +10,7 @@ use http_body_util::{BodyExt, Empty, Limited, StreamBody}; use hyper::body::{Frame, Incoming}; use hyper::{Request as HyperRequest, Uri}; use log::{debug, error, warn}; +use prometheus::IntCounter; use std::convert::TryFrom; use std::future::Future; use std::io; @@ -17,6 +18,7 @@ use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Duration; use tokio::sync::mpsc; +use tokio::time::error::Elapsed; use tokio::time::{Instant, Sleep}; use url::{Position, Url}; @@ -26,7 +28,8 @@ use crate::cache::{ }; use crate::coalesce::{FlightGuard, Join, Waiter}; use crate::forwarded::{ClientInfo, CLIENT_FORWARDING_HEADERS, URL_OVERRIDE_HEADERS}; -use crate::{AppState, CacheStats}; +use crate::metrics::{self, OriginFailure, OriginKind, Route}; +use crate::AppState; const CACHE_STATUS: &str = "x-shadowstep-cache"; @@ -37,10 +40,34 @@ pub(crate) type UpstreamBody = UnsyncBoxBody; type Flight = FlightGuard; +type OriginError = hyper_util::client::legacy::Error; + pub async fn forward_to_upstream( req: HttpRequest, payload: web::Payload, state: web::Data, +) -> HttpResponse { + let request_policy = RequestPolicy::new(req.method(), req.headers()); + let bypass = request_policy.bypasses_cache(); + let response = forward(req, payload, state.clone(), request_policy).await; + // an error response carries no `X-Shadowstep-Cache`, and it came from + // the origin exchange, so it counts as a miss + let cache = if bypass { + metrics::BYPASS + } else { + metrics::cache_label(&response).unwrap_or(metrics::MISS) + }; + state + .metrics + .observe(Route::Proxy, cache, response.status()); + response +} + +async fn forward( + req: HttpRequest, + payload: web::Payload, + state: web::Data, + request_policy: RequestPolicy, ) -> HttpResponse { let client = ClientInfo::from_request(&req); @@ -55,7 +82,6 @@ pub async fn forward_to_upstream( // the key takes scheme and host from `ClientInfo`, the source of the // X-Forwarded-Proto and X-Forwarded-Host that the origin receives let cache_key = PrimaryKey::new(client.scheme, &client.host, path_and_query); - let request_policy = RequestPolicy::new(req.method(), req.headers()); let ToOrigin { stale, flight } = match serve_from_cache(&req, &state, &cache_key, &request_policy).await { Before::Answer(response) => return response, @@ -64,7 +90,6 @@ pub async fn forward_to_upstream( let Some((builder, target_uri)) = upstream_request(&req, &state, &client, path_and_query) else { - state.cache_stats.miss(); return error_response(StatusCode::INTERNAL_SERVER_ERROR); }; // a client's own preconditions go to the origin unchanged, and then the @@ -86,19 +111,10 @@ pub async fn forward_to_upstream( let hyper_req = match origin_request(builder, payload).await { Ok(hyper_req) => hyper_req, - Err(status) => { - exchange.state.cache_stats.miss(); - return error_response(status); - } + Err(status) => return error_response(status), }; - let upstream = tokio::time::timeout( - exchange.state.upstream_timeout, - exchange.state.http_client.request(hyper_req), - ) - .await; - - match upstream { + match send_to_origin(&exchange.state, OriginKind::Foreground, hyper_req).await { Ok(Ok(upstream_response)) => { debug!( "Received response from upstream: {:?}", @@ -206,17 +222,11 @@ enum FreshUse { } impl FreshUse { - /// counts the use and returns its `X-Shadowstep-Cache` value. - fn count(self, stats: &CacheStats) -> &'static str { + /// the use's `X-Shadowstep-Cache` value. + fn cache_status(self) -> &'static str { match self { - FreshUse::Hit => { - stats.hit(); - "HIT" - } - FreshUse::Coalesced => { - stats.coalesced(); - "COALESCED" - } + FreshUse::Hit => metrics::HIT, + FreshUse::Coalesced => metrics::COALESCED, } } } @@ -262,8 +272,7 @@ fn from_cache( { Some(Lookup::Fresh(stored)) => { debug!("Cache hit for {} {}", req.method(), req.uri()); - let cache_status = fresh_use.count(&state.cache_stats); - Cached::Answer(cached_response(&stored, cache_status)) + Cached::Answer(cached_response(&stored, fresh_use.cache_status())) } Some(Lookup::Stale(entry)) if req.method() == Method::GET => { while_revalidating(req, state, cache_key, request_policy, entry) @@ -286,8 +295,7 @@ fn while_revalidating( return Cached::Stale(entry); } debug!("Serving stale {} while revalidating", req.uri()); - state.cache_stats.stale(); - let response = cached_response(&entry.response, "STALE"); + let response = cached_response(&entry.response, metrics::STALE); refresh_in_background(req, state, cache_key.clone(), entry); Cached::Answer(response) } @@ -318,14 +326,11 @@ fn refresh_in_background( let Ok(hyper_req) = with_validators(builder, &entry.response).body(empty_body()) else { return; }; - state.cache_stats.background_refresh(); + state.metrics.background_refresh(); let (req, state) = (req.clone(), state.clone()); actix_web::rt::spawn(async move { let _guard = guard; - let upstream = - tokio::time::timeout(state.upstream_timeout, state.http_client.request(hyper_req)) - .await; - match upstream { + match send_to_origin(&state, OriginKind::Background, hyper_req).await { Ok(Ok(response)) => { refresh(&req, &state, cache_key, &entry, response, target_uri).await } @@ -357,14 +362,15 @@ async fn refresh( &stored_fields(forwardable_fields(&head)), head.map.get(header::AGE), ); - state.cache_stats.revalidation(); + state.metrics.revalidation(); return; } let Some(store) = store_plan(req, state, cache_key, &request_policy, &head) else { return; }; let limit = usize::try_from(store.limit).unwrap_or(usize::MAX); - let body = origin_body(body, state.upstream_timeout, target_uri); + let idle_timeouts = state.metrics.body_idle_timeouts(OriginKind::Background); + let body = origin_body(body, state.upstream_timeout, target_uri, idle_timeouts); let body = StreamBody::new(body.map(|chunk| chunk.map(Frame::data))); match Limited::new(body, limit).collect().await { Ok(collected) => (store.finish)( @@ -413,7 +419,6 @@ impl Exchange { return self.serve_stale(stale); } } - self.state.cache_stats.miss(); let flight = self.flight; let store = store_plan( &self.req, @@ -422,7 +427,11 @@ impl Exchange { &self.request_policy, &head, ); - let body = origin_body(body, self.state.upstream_timeout, target_uri); + let idle_timeouts = self + .state + .metrics + .body_idle_timeouts(OriginKind::Foreground); + let body = origin_body(body, self.state.upstream_timeout, target_uri, idle_timeouts); client_response(head, body, store.map(|store| store.holding(flight))) } @@ -439,7 +448,6 @@ impl Exchange { status = StatusCode::GATEWAY_TIMEOUT; } } - self.state.cache_stats.miss(); error_response(status) } @@ -452,8 +460,7 @@ impl Exchange { "Serving stale {} in place of an origin error", self.req.uri() ); - self.state.cache_stats.stale(); - cached_response(&stale.response, "STALE") + cached_response(&stale.response, metrics::STALE) } /// the stored body with the fields of the origin's 304, which also @@ -466,9 +473,32 @@ impl Exchange { &stored_fields(forwardable_fields(head)), head.map.get(header::AGE), ); - self.state.cache_stats.revalidation(); - cached_response(&fresh, "REVALIDATED") + self.state.metrics.revalidation(); + cached_response(&fresh, metrics::REVALIDATED) + } +} + +/// sends `request` to the origin and waits up to the upstream timeout for +/// its response headers, counting the outcome as `kind`. +async fn send_to_origin( + state: &AppState, + kind: OriginKind, + request: HyperRequest, +) -> Result, OriginError>, Elapsed> { + let started = Instant::now(); + let upstream = + tokio::time::timeout(state.upstream_timeout, state.http_client.request(request)).await; + match &upstream { + Ok(Ok(response)) => { + let status = response.status().as_u16(); + state + .metrics + .origin_responded(kind, status, started.elapsed()); + } + Ok(Err(_)) => state.metrics.origin_failed(kind, OriginFailure::Error), + Err(_) => state.metrics.origin_failed(kind, OriginFailure::Timeout), } + upstream } /// `builder` with the conditional fields from `stored`'s validators (RFC @@ -716,7 +746,7 @@ where for (name, value) in &fields { builder.append_header((name.clone(), value.clone())); } - builder.insert_header((CACHE_STATUS, "MISS")); + builder.insert_header((CACHE_STATUS, metrics::MISS)); let stored_headers = stored_fields(fields); let length = head @@ -781,6 +811,7 @@ fn origin_body( body: Incoming, idle: Duration, target_uri: Uri, + idle_timeouts: IntCounter, ) -> IdleTimeout> + Unpin> { // the data stream drops trailer frames let body = body.into_data_stream().map(|chunk| { @@ -796,6 +827,7 @@ fn origin_body( waiting: false, expired: false, target_uri, + idle_timeouts, } } @@ -812,6 +844,8 @@ struct IdleTimeout { waiting: bool, expired: bool, target_uri: Uri, + /// counts this body if it goes idle + idle_timeouts: IntCounter, } impl Stream for IdleTimeout @@ -835,6 +869,7 @@ where } std::task::ready!(this.sleep.as_mut().poll(cx)); this.expired = true; + this.idle_timeouts.inc(); warn!( "Upstream {} sent no response body data for {:?}", this.target_uri, this.idle diff --git a/tests/integration/body_idle.rs b/tests/integration/body_idle.rs index 79ec61a..8a52355 100644 --- a/tests/integration/body_idle.rs +++ b/tests/integration/body_idle.rs @@ -1,4 +1,4 @@ -use crate::common::{self, Running}; +use crate::common::{self, metric, parse_exposition, Running, Sample}; use bytes::Bytes; use http_body_util::Empty; @@ -197,6 +197,50 @@ async fn sized_body_that_stalls_closes_the_client_connection_short() { proxy.stop().await; } +async fn scrape(proxy: &Running) -> Vec { + let client = common::client::>(); + let resp = timeout(PATIENCE, client.get(proxy.url("/metrics").parse().unwrap())) + .await + .expect("scrape timed out") + .unwrap(); + assert_eq!(resp.status(), 200); + let body = common::body_bytes(resp.into_body()).await; + parse_exposition(std::str::from_utf8(&body).unwrap()) +} + +fn body_idle_timeouts(samples: &[Sample], kind: &str) -> Option { + metric( + samples, + "shadowstep_origin_body_idle_timeouts_total", + &[("kind", kind)], + ) +} + +#[actix_web::test] +async fn stalled_body_counts_one_idle_timeout_and_one_origin_response() { + let mut origin = scripted_origin(vec![stalled_sized()]).await; + let proxy = spawn(&origin); + + raw_get(&proxy, "/page").await; + origin_connection_closed(&mut origin, 0).await; + let samples = scrape(&proxy).await; + + assert_eq!(body_idle_timeouts(&samples, "foreground"), Some(1.0)); + assert_eq!(body_idle_timeouts(&samples, "background"), Some(0.0)); + // the idle timeout comes after the headers, so the origin request is + // still counted once, by its status + let origin_requests = |outcome| { + metric( + &samples, + "shadowstep_origin_requests_total", + &[("kind", "foreground"), ("outcome", outcome)], + ) + }; + assert_eq!(origin_requests("2xx"), Some(1.0)); + assert_eq!(origin_requests("timeout"), Some(0.0)); + 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; @@ -318,6 +362,9 @@ async fn background_revalidation_whose_body_stalls_is_given_up() { raw_get(&proxy, "/page").await; } assert_eq!(background_refreshes(&proxy).await, 2); + let samples = scrape(&proxy).await; + assert_eq!(body_idle_timeouts(&samples, "background"), Some(1.0)); + assert_eq!(body_idle_timeouts(&samples, "foreground"), Some(0.0)); assert!( started.elapsed() < IDLE + MARGIN, "took {:?}", diff --git a/tests/integration/common.rs b/tests/integration/common.rs index 927ed9d..fbc319b 100644 --- a/tests/integration/common.rs +++ b/tests/integration/common.rs @@ -8,7 +8,8 @@ use hyper_util::client::legacy::connect::HttpConnector; use hyper_util::client::legacy::Client; use hyper_util::rt::TokioExecutor; use shadowstep::config::Config; -use shadowstep::{app, build_state, run, AppState}; +use shadowstep::{app, build_state, run, run_metrics, AppState}; +use std::collections::BTreeMap; use std::net::{SocketAddr, TcpListener}; use std::path::Path; use tempfile::TempDir; @@ -26,6 +27,7 @@ pub fn config(origin_url: &str, asset_path: &Path) -> Config { tls_key_path: None, tls_listen_addr: "127.0.0.1:0".to_owned(), upstream_timeout_seconds: 30, + metrics_addr: None, } } @@ -118,6 +120,50 @@ pub fn spawn_workers( } } +/// a proxy with `--metrics-addr` set, and the metrics listener's address. +pub struct WithMetricsListener { + pub proxy: Running, + pub metrics_addr: SocketAddr, + metrics_handle: actix_web::dev::ServerHandle, +} + +impl WithMetricsListener { + pub fn metrics_url(&self, path: &str) -> String { + format!("http://{}{}", self.metrics_addr, path) + } + + pub async fn stop(self) { + self.metrics_handle.stop(true).await; + self.proxy.stop().await; + } +} + +/// `spawn` with a separate metrics listener on an ephemeral port. +pub fn spawn_with_metrics_listener(origin_url: &str) -> WithMetricsListener { + let metrics_listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let metrics_addr = metrics_listener.local_addr().unwrap(); + let (state, assets) = state_with(origin_url, |c| { + c.metrics_addr = Some(metrics_addr.to_string()); + }); + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let server = run(state.clone(), listener, None, 1).unwrap(); + let handle = server.handle(); + actix_web::rt::spawn(server); + let metrics = run_metrics(state, metrics_listener).unwrap(); + let metrics_handle = metrics.handle(); + actix_web::rt::spawn(metrics); + WithMetricsListener { + proxy: Running { + addr, + handle, + _assets: assets, + }, + metrics_addr, + metrics_handle, + } +} + /// an origin URL that refuses connections: bind an ephemeral port, then /// release it. pub fn unreachable_origin() -> String { @@ -148,3 +194,94 @@ where pub async fn body_bytes(body: Incoming) -> Bytes { body.collect().await.unwrap().to_bytes() } + +/// one sample line of a Prometheus text exposition. +#[derive(Debug)] +pub struct Sample { + pub name: String, + pub labels: BTreeMap, + pub value: f64, +} + +/// the samples in a Prometheus text exposition (format 0.0.4). comment +/// lines and blank lines are skipped. +pub fn parse_exposition(text: &str) -> Vec { + text.lines() + .filter(|line| !line.is_empty() && !line.starts_with('#')) + .map(parse_sample) + .collect() +} + +fn parse_sample(line: &str) -> Sample { + let (series, value) = line.rsplit_once(' ').expect("sample without a value"); + let value = value + .parse() + .unwrap_or_else(|_| panic!("bad value in {line}")); + let Some((name, labels)) = series.split_once('{') else { + return Sample { + name: series.to_owned(), + labels: BTreeMap::new(), + value, + }; + }; + let labels = labels.strip_suffix('}').expect("unclosed label set"); + Sample { + name: name.to_owned(), + labels: parse_labels(labels), + value, + } +} + +/// `key="value",...` with the text format's escapes for `\`, `"` and newline. +fn parse_labels(text: &str) -> BTreeMap { + let mut labels = BTreeMap::new(); + let mut chars = text.chars().peekable(); + while chars.peek().is_some() { + let key: String = chars.by_ref().take_while(|&c| c != '=').collect(); + assert_eq!(chars.next(), Some('"'), "label {key} has no opening quote"); + let mut value = String::new(); + while let Some(c) = chars.next() { + match c { + '"' => break, + '\\' => match chars.next() { + Some('n') => value.push('\n'), + Some(other) => value.push(other), + None => panic!("dangling escape in label {key}"), + }, + c => value.push(c), + } + } + labels.insert(key.trim_start_matches(',').to_owned(), value); + if chars.peek() == Some(&',') { + chars.next(); + } + } + labels +} + +/// the value of the series `name` whose labels are exactly `labels`, or +/// `None` when the exposition has no such series. +pub fn metric(samples: &[Sample], name: &str, labels: &[(&str, &str)]) -> Option { + let wanted: BTreeMap = labels + .iter() + .map(|(k, v)| ((*k).to_owned(), (*v).to_owned())) + .collect(); + samples + .iter() + .find(|s| s.name == name && s.labels == wanted) + .map(|s| s.value) +} + +/// the sum over every series of `name` that has all of `labels`. +pub fn metric_sum(samples: &[Sample], name: &str, labels: &[(&str, &str)]) -> f64 { + samples + .iter() + .filter(|s| s.name == name) + .filter(|s| { + labels + .iter() + .all(|(k, v)| s.labels.get(*k).map(String::as_str) == Some(*v)) + }) + .map(|s| s.value) + .sum() +} diff --git a/tests/integration/main.rs b/tests/integration/main.rs index 5e323ff..a9c044d 100644 --- a/tests/integration/main.rs +++ b/tests/integration/main.rs @@ -5,6 +5,7 @@ mod coalescing; mod common; mod forwarded; mod http2; +mod metrics; mod proxy; mod revalidation; mod smoke; diff --git a/tests/integration/metrics.rs b/tests/integration/metrics.rs new file mode 100644 index 0000000..6afabcd --- /dev/null +++ b/tests/integration/metrics.rs @@ -0,0 +1,461 @@ +use crate::common::{self, cache_status, metric, metric_sum, parse_exposition, Sample}; + +use actix_web::body::MessageBody; +use actix_web::dev::{Service, ServiceResponse}; +use actix_web::http::StatusCode; +use actix_web::test; +use bytes::Bytes; +use http_body_util::Empty; +use std::time::{Duration, Instant}; +use wiremock::matchers::{any, method, path}; +use wiremock::{Mock, MockServer, ResponseTemplate}; + +const CONTENT_TYPE: &str = "text/plain; version=0.0.4; charset=utf-8"; + +fn get(uri: &str) -> test::TestRequest { + test::TestRequest::get().uri(uri) +} + +/// sends `req` and reads the whole response, so that a streamed origin body +/// has ended before the next scrape. +async fn call(app: &S, req: test::TestRequest) -> (StatusCode, String) +where + S: Service, Error = actix_web::Error>, + B: MessageBody, +{ + let resp = test::call_service(app, req.to_request()).await; + let status = resp.status(); + let cache = cache_status(&resp); + test::read_body(resp).await; + (status, cache) +} + +async fn scrape(app: &S) -> Vec +where + S: Service, Error = actix_web::Error>, + B: MessageBody, +{ + let resp = test::call_service(app, get("/metrics").to_request()).await; + assert_eq!(resp.status(), StatusCode::OK); + let body = test::read_body(resp).await; + parse_exposition(std::str::from_utf8(&body).unwrap()) +} + +async fn health(app: &S) -> serde_json::Value +where + S: Service, Error = actix_web::Error>, + B: MessageBody, +{ + let resp = test::call_service(app, get("/health").to_request()).await; + test::read_body_json(resp).await +} + +fn requests(samples: &[Sample], route: &str, cache: &str) -> Option { + metric( + samples, + "shadowstep_requests_total", + &[("route", route), ("cache", cache)], + ) +} + +fn responses(samples: &[Sample], route: &str, class: &str) -> Option { + metric( + samples, + "shadowstep_responses_total", + &[("route", route), ("status_class", class)], + ) +} + +fn origin(samples: &[Sample], kind: &str, outcome: &str) -> Option { + metric( + samples, + "shadowstep_origin_requests_total", + &[("kind", kind), ("outcome", outcome)], + ) +} + +fn cacheable(body: &str) -> ResponseTemplate { + ResponseTemplate::new(200) + .insert_header("cache-control", "max-age=60") + .set_body_string(body) +} + +#[actix_web::test] +async fn metrics_has_the_prometheus_text_content_type() { + let (app, _assets) = common::service(&common::unreachable_origin()).await; + + let resp = test::call_service(&app, get("/metrics").to_request()).await; + + assert_eq!(resp.status(), StatusCode::OK); + assert_eq!(resp.headers().get("content-type").unwrap(), CONTENT_TYPE); + let body = test::read_body(resp).await; + let text = std::str::from_utf8(&body).unwrap(); + assert!( + text.contains("# TYPE shadowstep_requests_total counter"), + "{text}" + ); +} + +#[actix_web::test] +async fn miss_then_hit_count_under_their_own_series() { + let origin_server = MockServer::start().await; + Mock::given(path("/page")) + .respond_with(cacheable("page")) + .expect(1) + .mount(&origin_server) + .await; + let (app, _assets) = common::service(&origin_server.uri()).await; + + assert_eq!(call(&app, get("/page")).await.1, "MISS"); + assert_eq!(call(&app, get("/page")).await.1, "HIT"); + let samples = scrape(&app).await; + + assert_eq!(requests(&samples, "proxy", "MISS"), Some(1.0)); + assert_eq!(requests(&samples, "proxy", "HIT"), Some(1.0)); + assert_eq!(metric_sum(&samples, "shadowstep_requests_total", &[]), 2.0); + assert_eq!(responses(&samples, "proxy", "2xx"), Some(2.0)); + assert_eq!(origin(&samples, "foreground", "2xx"), Some(1.0)); + let health = health(&app).await; + assert_eq!(health["cache"]["hits"], 1); + assert_eq!(health["cache"]["misses"], 1); +} + +#[actix_web::test] +async fn post_counts_as_bypass() { + let origin_server = MockServer::start().await; + Mock::given(method("POST")) + .respond_with(ResponseTemplate::new(201)) + .expect(1) + .mount(&origin_server) + .await; + let (app, _assets) = common::service(&origin_server.uri()).await; + + let (status, _) = call(&app, test::TestRequest::post().uri("/form")).await; + assert_eq!(status, StatusCode::CREATED); + let samples = scrape(&app).await; + + assert_eq!(requests(&samples, "proxy", "BYPASS"), Some(1.0)); + assert_eq!(requests(&samples, "proxy", "MISS"), Some(0.0)); + assert_eq!(responses(&samples, "proxy", "2xx"), Some(1.0)); + assert_eq!(origin(&samples, "foreground", "2xx"), Some(1.0)); + // `/health` has always counted a bypass as a miss + assert_eq!(health(&app).await["cache"]["misses"], 1); +} + +#[actix_web::test] +async fn origin_502_and_timeout_count_by_outcome() { + let origin_server = MockServer::start().await; + Mock::given(path("/bad")) + .respond_with(ResponseTemplate::new(502)) + .mount(&origin_server) + .await; + Mock::given(path("/slow")) + .respond_with(ResponseTemplate::new(200).set_delay(Duration::from_secs(3))) + .mount(&origin_server) + .await; + let (app, _assets) = + common::service_with(&origin_server.uri(), |c| c.upstream_timeout_seconds = 1).await; + + assert_eq!(call(&app, get("/bad")).await.0, StatusCode::BAD_GATEWAY); + assert_eq!( + call(&app, get("/slow")).await.0, + StatusCode::GATEWAY_TIMEOUT + ); + let samples = scrape(&app).await; + + assert_eq!(origin(&samples, "foreground", "5xx"), Some(1.0)); + assert_eq!(origin(&samples, "foreground", "timeout"), Some(1.0)); + assert_eq!( + metric_sum(&samples, "shadowstep_origin_requests_total", &[]), + 2.0 + ); + assert_eq!(responses(&samples, "proxy", "5xx"), Some(2.0)); + assert_eq!(requests(&samples, "proxy", "MISS"), Some(2.0)); +} + +#[actix_web::test] +async fn unreachable_origin_counts_as_error() { + let (app, _assets) = common::service(&common::unreachable_origin()).await; + + assert_eq!(call(&app, get("/page")).await.0, StatusCode::BAD_GATEWAY); + let samples = scrape(&app).await; + + assert_eq!(origin(&samples, "foreground", "error"), Some(1.0)); + assert_eq!( + metric_sum(&samples, "shadowstep_origin_requests_total", &[]), + 1.0 + ); + assert_eq!(responses(&samples, "proxy", "5xx"), Some(1.0)); +} + +#[actix_web::test] +async fn histogram_counts_each_origin_response() { + let origin_server = MockServer::start().await; + Mock::given(path("/ok")) + .respond_with(ResponseTemplate::new(200)) + .mount(&origin_server) + .await; + Mock::given(path("/gone")) + .respond_with(ResponseTemplate::new(404)) + .mount(&origin_server) + .await; + Mock::given(path("/bad")) + .respond_with(ResponseTemplate::new(503)) + .mount(&origin_server) + .await; + let (app, _assets) = common::service(&origin_server.uri()).await; + + for uri in ["/ok", "/gone", "/bad", "/ok"] { + call(&app, get(uri)).await; + } + let samples = scrape(&app).await; + + let received = origin_server.received_requests().await.unwrap().len() as f64; + assert_eq!(received, 4.0); + let count = metric(&samples, "shadowstep_origin_response_seconds_count", &[]); + assert_eq!(count, Some(received)); + let infinite = metric( + &samples, + "shadowstep_origin_response_seconds_bucket", + &[("le", "+Inf")], + ); + assert_eq!(infinite, Some(received)); + assert_eq!( + metric_sum(&samples, "shadowstep_origin_requests_total", &[]), + received + ); + let sum = metric(&samples, "shadowstep_origin_response_seconds_sum", &[]).unwrap(); + assert!(sum > 0.0, "{sum}"); +} + +#[actix_web::test] +async fn cache_gauges_match_health() { + let origin_server = MockServer::start().await; + Mock::given(any()) + .respond_with(cacheable("a stored body")) + .mount(&origin_server) + .await; + let (app, assets) = common::service(&origin_server.uri()).await; + std::fs::write(assets.path().join("app.css"), "body{}").unwrap(); + + call(&app, get("/one")).await; + call(&app, get("/two")).await; + call(&app, get("/assets/app.css")).await; + let samples = scrape(&app).await; + let health = health(&app).await; + + let items = health["cache"]["items"].as_f64().unwrap(); + let bytes = health["cache"]["bytes"].as_f64().unwrap(); + assert_eq!(items, 3.0); + assert!(bytes > 0.0); + assert_eq!( + metric(&samples, "shadowstep_cache_entries", &[]), + Some(items) + ); + assert_eq!(metric(&samples, "shadowstep_cache_bytes", &[]), Some(bytes)); +} + +#[actix_web::test] +async fn assets_count_under_their_own_route() { + let (app, assets) = common::service(&common::unreachable_origin()).await; + std::fs::write(assets.path().join("app.css"), "body{}").unwrap(); + + assert_eq!(call(&app, get("/assets/app.css")).await.1, "MISS"); + assert_eq!(call(&app, get("/assets/app.css")).await.1, "HIT"); + assert_eq!( + call(&app, get("/assets/missing.css")).await.0, + StatusCode::NOT_FOUND + ); + let samples = scrape(&app).await; + + assert_eq!(requests(&samples, "asset", "MISS"), Some(1.0)); + assert_eq!(requests(&samples, "asset", "HIT"), Some(1.0)); + assert_eq!(requests(&samples, "asset", "BYPASS"), Some(1.0)); + assert_eq!(responses(&samples, "asset", "2xx"), Some(2.0)); + assert_eq!(responses(&samples, "asset", "4xx"), Some(1.0)); + assert_eq!( + metric_sum(&samples, "shadowstep_requests_total", &[("route", "proxy")]), + 0.0 + ); + // `/health` has never counted an asset that was not found + let health = health(&app).await; + assert_eq!(health["cache"]["hits"], 1); + assert_eq!(health["cache"]["misses"], 1); +} + +#[actix_web::test] +async fn labels_never_hold_paths() { + let origin_server = MockServer::start().await; + Mock::given(any()) + .respond_with(cacheable("page")) + .mount(&origin_server) + .await; + let (app, assets) = common::service(&origin_server.uri()).await; + std::fs::write(assets.path().join("app.css"), "body{}").unwrap(); + + for uri in ["/secret-page?token=abc", "/other", "/assets/app.css"] { + call(&app, get(uri)).await; + } + call(&app, test::TestRequest::post().uri("/secret-form")).await; + let samples = scrape(&app).await; + + assert!(!samples.is_empty()); + let allowed = ["route", "cache", "status_class", "kind", "outcome", "le"]; + for sample in &samples { + for (key, value) in &sample.labels { + assert!(allowed.contains(&key.as_str()), "{sample:?}"); + assert!( + !value.contains('/') && !value.contains("secret") && !value.contains("app"), + "{sample:?}" + ); + } + } +} + +#[actix_web::test] +async fn metrics_and_health_are_not_proxied_or_counted() { + let origin_server = MockServer::start().await; + Mock::given(any()) + .respond_with(ResponseTemplate::new(200).set_body_string("origin metrics")) + .expect(0) + .mount(&origin_server) + .await; + let (app, _assets) = common::service(&origin_server.uri()).await; + + health(&app).await; + scrape(&app).await; + let samples = scrape(&app).await; + + assert_eq!(metric_sum(&samples, "shadowstep_requests_total", &[]), 0.0); + assert_eq!(metric_sum(&samples, "shadowstep_responses_total", &[]), 0.0); + origin_server.verify().await; +} + +#[actix_web::test] +async fn background_revalidation_counts_as_background() { + let origin_server = MockServer::start().await; + Mock::given(any()) + .respond_with( + ResponseTemplate::new(200) + .insert_header("cache-control", "max-age=1, stale-while-revalidate=10") + .insert_header("etag", "\"v1\"") + .set_body_string("old"), + ) + .up_to_n_times(1) + .mount(&origin_server) + .await; + Mock::given(any()) + .respond_with(ResponseTemplate::new(304).insert_header("cache-control", "max-age=60")) + .mount(&origin_server) + .await; + let (app, _assets) = common::service(&origin_server.uri()).await; + + call(&app, get("/page")).await; + // max-age has a granularity of one second + actix_web::rt::time::sleep(Duration::from_millis(1100)).await; + assert_eq!(call(&app, get("/page")).await.1, "STALE"); + let deadline = Instant::now() + Duration::from_secs(5); + let samples = loop { + let samples = scrape(&app).await; + if metric(&samples, "shadowstep_cache_revalidations_total", &[]) == Some(1.0) { + break samples; + } + assert!(Instant::now() < deadline, "the background 304 never landed"); + actix_web::rt::time::sleep(Duration::from_millis(20)).await; + }; + + assert_eq!(origin(&samples, "foreground", "2xx"), Some(1.0)); + assert_eq!(origin(&samples, "background", "3xx"), Some(1.0)); + assert_eq!( + metric(&samples, "shadowstep_cache_background_refreshes_total", &[]), + Some(1.0) + ); + assert_eq!(requests(&samples, "proxy", "STALE"), Some(1.0)); + assert_eq!( + metric(&samples, "shadowstep_origin_response_seconds_count", &[]), + Some(2.0) + ); + let health = health(&app).await; + assert_eq!(health["cache"]["revalidations"], 1); + assert_eq!(health["cache"]["background_refreshes"], 1); + assert_eq!(health["cache"]["stale"], 1); +} + +/// a GET of `url` over a socket: the status, `Content-Type` and body. +async fn socket_get(url: &str) -> (u16, String, String) { + let client = common::client::>(); + let resp = client.get(url.parse().unwrap()).await.unwrap(); + let status = resp.status().as_u16(); + let content_type = resp + .headers() + .get("content-type") + .map(|v| v.to_str().unwrap().to_owned()) + .unwrap_or_default(); + let body = common::body_bytes(resp.into_body()).await; + ( + status, + content_type, + String::from_utf8_lossy(&body).into_owned(), + ) +} + +#[actix_web::test] +async fn metrics_addr_moves_metrics_to_its_own_listener() { + let origin_server = MockServer::start().await; + Mock::given(path("/metrics")) + .respond_with(ResponseTemplate::new(200).set_body_string("origin metrics")) + .expect(1) + .mount(&origin_server) + .await; + Mock::given(path("/page")) + .respond_with(ResponseTemplate::new(200).set_body_string("page")) + .expect(0) + .mount(&origin_server) + .await; + let running = common::spawn_with_metrics_listener(&origin_server.uri()); + + let (status, content_type, body) = socket_get(&running.metrics_url("/metrics")).await; + assert_eq!(status, 200); + assert_eq!(content_type, CONTENT_TYPE); + assert!( + body.contains("# TYPE shadowstep_requests_total counter"), + "{body}" + ); + + // the proxy listener forwards `/metrics` to the origin like any path + let (status, _, body) = socket_get(&running.proxy.url("/metrics")).await; + assert_eq!(status, 200); + assert_eq!(body, "origin metrics"); + + // the metrics listener serves nothing else and proxies nothing + let (status, _, _) = socket_get(&running.metrics_url("/page")).await; + assert_eq!(status, 404); + let (status, _, _) = socket_get(&running.metrics_url("/health")).await; + assert_eq!(status, 404); + + let (_, _, body) = socket_get(&running.metrics_url("/metrics")).await; + let samples = parse_exposition(&body); + assert_eq!(requests(&samples, "proxy", "MISS"), Some(1.0)); + assert_eq!(metric_sum(&samples, "shadowstep_requests_total", &[]), 1.0); + running.stop().await; + origin_server.verify().await; +} + +#[actix_web::test] +async fn metrics_is_served_on_the_proxy_listener_by_default() { + let origin_server = MockServer::start().await; + Mock::given(any()) + .respond_with(ResponseTemplate::new(200).set_body_string("origin metrics")) + .expect(0) + .mount(&origin_server) + .await; + let proxy = common::spawn(&origin_server.uri()); + + let (status, content_type, body) = socket_get(&proxy.url("/metrics")).await; + + assert_eq!(status, 200); + assert_eq!(content_type, CONTENT_TYPE); + assert!(body.contains("shadowstep_cache_entries 0"), "{body}"); + proxy.stop().await; + origin_server.verify().await; +}