//! An OSM tile caching proxy. //! //! The web UI never talks to `tile.openstreetmap.org` directly. Two reasons, and //! only the second is about performance: the browser would leak every viewer's IP //! and viewport to a third party, and a small deployment re-requesting the same //! city block all day is exactly the traffic the OSM tile usage policy asks //! proxies to absorb. //! //! Bytes live at `{cache_dir}/tiles/{z}/{x}/{y}.png`. The `tiles` table is the //! index and the byte accounting; the filesystem alone could not answer "what is //! the least recently used tile" without walking it. use std::path::{Path as FsPath, PathBuf}; use std::sync::OnceLock; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; use anyhow::{Context, Result}; use axum::Router; use axum::extract::{Path, State}; use axum::http::{StatusCode, header}; use axum::response::{IntoResponse, Response}; use axum::routing::get; use sqlx::SqlitePool; use tokio::io::AsyncWriteExt; use tower_sessions::Session; use tracing::{debug, warn}; use crate::api::{ApiError, Shared, current_user}; use crate::db::now; /// Deepest zoom the proxy will fetch. OSM itself stops at 19, and accepting more /// only lets a caller mint cache entries upstream will never satisfy. const MAX_ZOOM: u8 = 19; /// Upper bound on an accepted response body. A PNG tile is a few tens of KiB; a /// hostile or misconfigured upstream must not be able to fill the disk with one /// response. const MAX_TILE_BYTES: usize = 256 * 1024; /// Used when upstream sends no `Cache-Control: max-age`. /// /// This deliberately overrides a bare `no-cache`, which is what /// `tile.openstreetmap.org` answers with. Honouring it would send every single /// tile request upstream and make the proxy worse than useless to the very /// servers it exists to spare. A rendered tile changes on the order of days. const DEFAULT_TTL_S: i64 = 7 * 86_400; /// How long the *browser* may reuse a tile without asking us again. Tiles are /// immutable in practice, so this is the cheapest hit of all: no request at all. const BROWSER_CACHE_CONTROL: &str = "public, max-age=86400"; const UPSTREAM_TIMEOUT: Duration = Duration::from_secs(10); /// Merged into the main router by [`crate::api::router`]. /// /// No `.png` suffix on the route: axum allows one parameter per path segment, so /// `{y}.png` is rejected at startup. Leaflet does not care about the extension, /// and the `Content-Type` header is what a browser actually reads. pub fn router() -> Router { Router::new().route("/tiles/{z}/{x}/{y}", get(tile)) } /// One client for the whole process, so connections to the tile server are kept /// alive across requests instead of paying a TLS handshake per tile. fn client() -> &'static reqwest::Client { static CLIENT: OnceLock = OnceLock::new(); CLIENT.get_or_init(|| { reqwest::Client::builder() .timeout(UPSTREAM_TIMEOUT) .build() // The builder only fails on a broken TLS backend, and a default // client still works — refusing to serve any map at all would be a // worse answer than one without keep-alive tuning. .unwrap_or_default() }) } #[derive(sqlx::FromRow)] struct TileRow { etag: Option, last_modified: Option, expires_at: i64, } /// Serve one tile, from cache when possible. /// /// Requires a signed-in session, like every other read path. An unauthenticated /// `/tiles` endpoint is an open proxy: strangers would burn this deployment's OSM /// quota and get its IP blocked. Leaflet loads tiles as same-origin `` /// requests, so the session cookie rides along without any JavaScript help. async fn tile( State(state): State, session: Session, Path((z, x, y)): Path<(u8, u32, u32)>, ) -> Result { current_user(&state, &session).await?; // Before anything touches the filesystem: the extractor guarantees these are // integers, not that they name a tile that can exist. Without this a garbage // request creates a directory tree. if !in_range(z, x, y) { return Err(ApiError::BadRequest(format!( "no such tile: z must be 0..={MAX_ZOOM} and x, y must be below 2^z" ))); } let path = tile_path(&state.cfg.cache_dir, z, x, y); let row: Option = sqlx::query_as( "SELECT etag, last_modified, expires_at FROM tiles WHERE z = ? AND x = ? AND y = ?", ) .bind(i64::from(z)) .bind(i64::from(x)) .bind(i64::from(y)) .fetch_optional(&state.db.read) .await?; // The row and the file can disagree — a manually cleared cache directory, a // half-restored backup. The bytes are the truth; a row without them is a miss. let cached = match &row { Some(_) => tokio::fs::read(&path).await.ok(), None => None, }; if row.is_some() && cached.is_none() { delete_row(&state.db.write, z, x, y).await?; } let at = now(); if let (Some(meta), Some(bytes)) = (&row, &cached) && meta.expires_at > at { touch(&state.db.write, z, x, y, at, None).await?; return Ok(png(bytes.clone())); } // Only when we hold the bytes a 304 would refer to. Sending a validator we // cannot honour would turn every request into an empty response. let conditional = row .as_ref() .filter(|_| cached.is_some()) .map(|m| (m.etag.clone(), m.last_modified.clone())); match fetch(&state.cfg, z, x, y, conditional).await { // The case that actually keeps the OSM quota happy: a revalidation costs // a few hundred bytes and refreshes a tile we already hold. Ok(Fetched::NotModified { expires_at }) => match cached { Some(bytes) => { touch(&state.db.write, z, x, y, at, Some(expires_at)).await?; Ok(png(bytes)) } // Upstream answered 304 to a request that carried no validator. // Serving the zero bytes we hold would render a broken image. None => Err(ApiError::Internal(anyhow::anyhow!( "tile upstream sent 304 for a tile we do not have" ))), }, Ok(Fetched::Body { bytes, etag, last_modified, expires_at, }) => { write_tile(&path, &bytes).await?; upsert( &state.db.write, z, x, y, &etag, &last_modified, at, expires_at, bytes.len() as i64, ) .await?; Ok(png(bytes)) } Err(e) => match cached { // A stale tile is a correct-looking map. An error is a grey square in // the middle of one, which reads as a broken deployment. Some(bytes) => { debug!(error = %e, z, x, y, "serving a stale tile; upstream is unavailable"); Ok(png(bytes)) } None => Err(ApiError::Internal(e)), }, } } fn in_range(z: u8, x: u32, y: u32) -> bool { z <= MAX_ZOOM && u64::from(x) < 1u64 << z && u64::from(y) < 1u64 << z } fn tile_path(cache_dir: &FsPath, z: u8, x: u32, y: u32) -> PathBuf { cache_dir.join(format!("tiles/{z}/{x}/{y}.png")) } fn upstream_url(template: &str, z: u8, x: u32, y: u32) -> String { template .replace("{z}", &z.to_string()) .replace("{x}", &x.to_string()) .replace("{y}", &y.to_string()) } fn png(bytes: Vec) -> Response { ( StatusCode::OK, [ (header::CONTENT_TYPE, "image/png"), (header::CACHE_CONTROL, BROWSER_CACHE_CONTROL), ], bytes, ) .into_response() } enum Fetched { NotModified { expires_at: i64, }, Body { bytes: Vec, etag: Option, last_modified: Option, expires_at: i64, }, } /// One upstream GET, conditional when we already hold a copy. /// /// The `User-Agent` is not decoration: the OSM tile usage policy prohibits /// library defaults and blocks unidentified proxies without notice, which is why /// [`crate::config::Config::validate`] refuses to start without a contact address. async fn fetch( cfg: &crate::config::Config, z: u8, x: u32, y: u32, conditional: Option<(Option, Option)>, ) -> Result { let url = upstream_url(&cfg.tile_upstream_url, z, x, y); let mut req = client() .get(&url) .header(header::USER_AGENT, cfg.tile_user_agent()); if let Some((etag, last_modified)) = conditional { if let Some(etag) = etag { req = req.header(header::IF_NONE_MATCH, etag); } if let Some(lm) = last_modified { req = req.header(header::IF_MODIFIED_SINCE, lm); } } let resp = req.send().await.with_context(|| format!("GET {url}"))?; let status = resp.status(); let expires_at = now() + max_age_of(resp.headers()).unwrap_or(DEFAULT_TTL_S); if status == reqwest::StatusCode::NOT_MODIFIED { return Ok(Fetched::NotModified { expires_at }); } if !status.is_success() { anyhow::bail!("tile upstream answered {status} for {url}"); } let etag = header_string(resp.headers(), header::ETAG); let last_modified = header_string(resp.headers(), header::LAST_MODIFIED); let bytes = read_capped(resp) .await .with_context(|| format!("body of {url}"))?; Ok(Fetched::Body { bytes, etag, last_modified, expires_at, }) } /// Read the body chunk by chunk, refusing to buffer more than /// [`MAX_TILE_BYTES`]. `Response::bytes` would happily allocate whatever the /// server sends, and `Content-Length` is the sender's claim, not a bound. async fn read_capped(mut resp: reqwest::Response) -> Result> { let mut out = Vec::new(); while let Some(chunk) = resp.chunk().await? { if out.len() + chunk.len() > MAX_TILE_BYTES { anyhow::bail!("tile body exceeds {MAX_TILE_BYTES} bytes"); } out.extend_from_slice(&chunk); } Ok(out) } fn header_string(headers: &reqwest::header::HeaderMap, name: header::HeaderName) -> Option { headers .get(name) .and_then(|v| v.to_str().ok()) .map(str::to_string) } /// `max-age` from a `Cache-Control` header, in seconds. fn max_age_of(headers: &reqwest::header::HeaderMap) -> Option { let value = headers.get(header::CACHE_CONTROL)?.to_str().ok()?; value .split(',') .filter_map(|part| part.trim().strip_prefix("max-age=")) .find_map(|n| n.trim().parse::().ok()) .filter(|n| *n > 0) } /// Write the bytes, then rename into place. /// /// The rename is the point: a concurrent reader either sees the previous file or /// the complete new one, never a half-written PNG. async fn write_tile(path: &FsPath, bytes: &[u8]) -> Result<()> { let dir = path.parent().context("tile path has no parent")?; tokio::fs::create_dir_all(dir) .await .with_context(|| format!("creating {}", dir.display()))?; // In the same directory, so the rename stays within one filesystem. static SEQ: AtomicU64 = AtomicU64::new(0); let temp = dir.join(format!( ".{}.{}.tmp", std::process::id(), SEQ.fetch_add(1, Ordering::Relaxed) )); let mut file = tokio::fs::File::create(&temp) .await .with_context(|| format!("creating {}", temp.display()))?; let written = async { file.write_all(bytes).await?; file.sync_all().await } .await; if let Err(e) = written { let _ = tokio::fs::remove_file(&temp).await; return Err(anyhow::Error::new(e).context("writing a tile")); } tokio::fs::rename(&temp, path) .await .with_context(|| format!("renaming into {}", path.display()))?; Ok(()) } async fn touch( pool: &SqlitePool, z: u8, x: u32, y: u32, at: i64, expires_at: Option, ) -> Result<(), sqlx::Error> { sqlx::query( "UPDATE tiles SET last_access = ?, expires_at = COALESCE(?, expires_at) \ WHERE z = ? AND x = ? AND y = ?", ) .bind(at) .bind(expires_at) .bind(i64::from(z)) .bind(i64::from(x)) .bind(i64::from(y)) .execute(pool) .await .map(|_| ()) } async fn delete_row(pool: &SqlitePool, z: u8, x: u32, y: u32) -> Result<(), sqlx::Error> { sqlx::query("DELETE FROM tiles WHERE z = ? AND x = ? AND y = ?") .bind(i64::from(z)) .bind(i64::from(x)) .bind(i64::from(y)) .execute(pool) .await .map(|_| ()) } #[allow(clippy::too_many_arguments)] async fn upsert( pool: &SqlitePool, z: u8, x: u32, y: u32, etag: &Option, last_modified: &Option, at: i64, expires_at: i64, bytes: i64, ) -> Result<(), sqlx::Error> { sqlx::query( "INSERT INTO tiles (z, x, y, etag, last_modified, fetched_at, expires_at, bytes, last_access) \ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) \ ON CONFLICT (z, x, y) DO UPDATE SET \ etag = excluded.etag, last_modified = excluded.last_modified, \ fetched_at = excluded.fetched_at, expires_at = excluded.expires_at, \ bytes = excluded.bytes, last_access = excluded.last_access", ) .bind(i64::from(z)) .bind(i64::from(x)) .bind(i64::from(y)) .bind(etag) .bind(last_modified) .bind(at) .bind(expires_at) .bind(bytes) .bind(at) .execute(pool) .await .map(|_| ()) } // --------------------------------------------------------------------------- // Eviction // --------------------------------------------------------------------------- /// How often the cache is measured against its ceiling. const EVICT_INTERVAL: Duration = Duration::from_secs(10 * 60); /// Rows deleted per pass, so the write lock is never held for long — the same /// reasoning as [`crate::retention`]'s batches. const EVICT_BATCH: i64 = 500; /// Eviction stops here rather than at the ceiling, so a cache sitting exactly at /// the limit does not evict a tile on every single fetch. fn low_water(max_bytes: u64) -> i64 { (max_bytes / 10 * 9) as i64 } /// Enforce `max_cache_bytes`, oldest access first. /// /// On a timer rather than on the request path: eviction is a whole-table sum and /// a batch of deletes, and making a map pan pay for that would be felt. pub fn spawn_eviction( pool: SqlitePool, cache_dir: PathBuf, max_bytes: u64, ) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { let mut ticker = tokio::time::interval(EVICT_INTERVAL); // Sleep first: startup already has enough to do. ticker.tick().await; loop { ticker.tick().await; if let Err(e) = evict_once(&pool, &cache_dir, max_bytes).await { warn!(error = %e, "tile eviction failed; will retry next tick"); } } }) } /// Delete least-recently-used tiles until the cache is under [`low_water`]. /// /// Returns the number of tiles removed. async fn evict_once(pool: &SqlitePool, cache_dir: &FsPath, max_bytes: u64) -> Result { let total: i64 = sqlx::query_scalar("SELECT COALESCE(SUM(bytes), 0) FROM tiles") .fetch_one(pool) .await .context("summing the tile cache")?; if total <= max_bytes as i64 { return Ok(0); } let mut remaining = total; let target = low_water(max_bytes); let mut removed = 0; // ponytail: one row deleted per statement, driven by the tiles_last_access // index. Fine up to the ~100k rows a 1 GiB cache holds; if a deployment runs // a much larger ceiling, delete by a last_access cutoff in one statement. while remaining > target { let batch: Vec<(i64, i64, i64, i64)> = sqlx::query_as("SELECT z, x, y, bytes FROM tiles ORDER BY last_access LIMIT ?") .bind(EVICT_BATCH) .fetch_all(pool) .await .context("listing the least recently used tiles")?; if batch.is_empty() { break; } for (z, x, y, bytes) in batch { sqlx::query("DELETE FROM tiles WHERE z = ? AND x = ? AND y = ?") .bind(z) .bind(x) .bind(y) .execute(pool) .await .context("evicting a tile row")?; // A leftover file is only wasted space, and the next fetch of that // tile overwrites it — so a failed unlink must not abort the sweep. let path = tile_path(cache_dir, z as u8, x as u32, y as u32); let _ = tokio::fs::remove_file(&path).await; remaining -= bytes; removed += 1; if remaining <= target { break; } } } debug!(removed, total, target, "tile cache eviction"); Ok(removed) } #[cfg(test)] mod tests { use super::*; /// A throwaway database, migrated and ready. Deliberately a local copy of /// `db::tests::test_db`: that module is private to `db.rs`, so it is not /// reachable from here even under `cfg(test)`. async fn test_db() -> (crate::db::Db, tempfile::TempDir) { let dir = tempfile::tempdir().expect("temp dir"); let db = crate::db::Db::open(&dir.path().join("test.db")) .await .expect("open"); (db, dir) } #[test] fn tiles_outside_the_pyramid_are_rejected() { assert!(in_range(0, 0, 0)); assert!(in_range(1, 1, 1)); assert!(in_range(19, (1 << 19) - 1, (1 << 19) - 1)); assert!(!in_range(20, 0, 0), "zoom beyond what OSM serves"); assert!(!in_range(1, 2, 0), "x must be below 2^z"); assert!(!in_range(1, 0, 2), "y must be below 2^z"); assert!(!in_range(0, 1, 0)); assert!(!in_range(19, 1 << 19, 0)); assert!(!in_range(3, u32::MAX, u32::MAX)); } #[test] fn the_upstream_url_is_built_from_the_template() { assert_eq!( upstream_url("https://tile.openstreetmap.org/{z}/{x}/{y}.png", 7, 66, 44), "https://tile.openstreetmap.org/7/66/44.png" ); // Subdomain-style templates put the placeholders elsewhere; the // substitution must not care where they are. assert_eq!( upstream_url("https://t.example/{x}-{y}@{z}", 3, 1, 2), "https://t.example/1-2@3" ); } #[test] fn the_cache_path_mirrors_the_request_path() { assert_eq!( tile_path(FsPath::new("/var/cache/ot"), 7, 66, 44), PathBuf::from("/var/cache/ot/tiles/7/66/44.png") ); } #[test] fn an_upstream_max_age_sets_the_expiry() { let mut headers = reqwest::header::HeaderMap::new(); assert_eq!( max_age_of(&headers), None, "no header means the default TTL" ); headers.insert( header::CACHE_CONTROL, "public, max-age=604800".parse().expect("literal"), ); assert_eq!(max_age_of(&headers), Some(604_800)); // `no-cache` carries no max-age, so the default applies rather than a // zero TTL that would revalidate on every single request. headers.insert(header::CACHE_CONTROL, "no-cache".parse().expect("literal")); assert_eq!(max_age_of(&headers), None); } #[tokio::test] async fn a_tile_is_renamed_into_place_rather_than_written_in_pieces() { let dir = tempfile::tempdir().expect("temp dir"); let path = tile_path(dir.path(), 4, 1, 2); write_tile(&path, b"\x89PNG").await.expect("write"); assert_eq!( tokio::fs::read(&path).await.expect("read"), b"\x89PNG", "the file must exist with its full contents" ); // No temp file survives a successful write. let leftovers: Vec<_> = std::fs::read_dir(path.parent().expect("parent")) .expect("read_dir") .filter_map(|e| e.ok()) .filter(|e| e.file_name().to_string_lossy().ends_with(".tmp")) .collect(); assert!(leftovers.is_empty(), "a temp file was left behind"); } /// Inserts `n` tiles of `bytes` each, oldest access first. async fn seed(pool: &SqlitePool, dir: &FsPath, n: u32, bytes: i64) { for i in 0..n { let path = tile_path(dir, 1, i, 0); write_tile(&path, &vec![0u8; bytes as usize]) .await .expect("tile file"); upsert(pool, 1, i, 0, &None, &None, i64::from(i), 0, bytes) .await .expect("row"); } } #[tokio::test] async fn eviction_removes_the_least_recently_used_tiles_down_to_the_low_water_mark() { let (db, _db_dir) = test_db().await; let cache = tempfile::tempdir().expect("temp dir"); // 10 tiles of 100 bytes against a 500-byte ceiling: 1000 bytes cached, // and eviction must stop at 450, not at 500. seed(&db.write, cache.path(), 10, 100).await; let removed = evict_once(&db.write, cache.path(), 500) .await .expect("evict"); assert_eq!(removed, 6, "1000 bytes down to 450 needs six tiles gone"); let survivors: Vec = sqlx::query_scalar("SELECT x FROM tiles ORDER BY x") .fetch_all(&db.read) .await .expect("rows"); assert_eq!( survivors, vec![6, 7, 8, 9], "the oldest accesses must go first" ); assert!( !tile_path(cache.path(), 1, 0, 0).exists(), "an evicted row must take its file with it, or the accounting lies" ); assert!(tile_path(cache.path(), 1, 9, 0).exists()); } #[tokio::test] async fn a_cache_under_its_ceiling_is_left_alone() { let (db, _db_dir) = test_db().await; let cache = tempfile::tempdir().expect("temp dir"); seed(&db.write, cache.path(), 4, 100).await; assert_eq!( evict_once(&db.write, cache.path(), 1_000) .await .expect("evict"), 0 ); let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM tiles") .fetch_one(&db.read) .await .expect("count"); assert_eq!(count, 4); } /// A fresh access must move a tile to the back of the eviction queue, which /// is the entire reason `last_access` is written on a cache hit. #[tokio::test] async fn a_recently_served_tile_outlives_an_older_one() { let (db, _db_dir) = test_db().await; let cache = tempfile::tempdir().expect("temp dir"); seed(&db.write, cache.path(), 4, 100).await; touch(&db.write, 1, 0, 0, 9_999, None).await.expect("touch"); evict_once(&db.write, cache.path(), 200) .await .expect("evict"); let survivors: Vec = sqlx::query_scalar("SELECT x FROM tiles ORDER BY x") .fetch_all(&db.read) .await .expect("rows"); assert!( survivors.contains(&0), "tile 0 was just served and must not be the first evicted, got {survivors:?}" ); } #[tokio::test] async fn a_refreshed_tile_keeps_one_row_rather_than_accumulating() { let (db, _db_dir) = test_db().await; upsert(&db.write, 5, 1, 2, &None, &None, 1, 100, 10) .await .expect("insert"); upsert( &db.write, 5, 1, 2, &Some("\"abc\"".into()), &None, 2, 200, 20, ) .await .expect("update"); let rows: Vec<(i64, Option, i64)> = sqlx::query_as("SELECT bytes, etag, expires_at FROM tiles") .fetch_all(&db.read) .await .expect("rows"); assert_eq!(rows.len(), 1); assert_eq!(rows[0].0, 20, "byte accounting must follow the new body"); assert_eq!(rows[0].1.as_deref(), Some("\"abc\"")); assert_eq!(rows[0].2, 200); } } // ponytail: two simultaneous requests for the same missing tile both fetch it. // A duplicated upstream GET is cheap and the atomic rename keeps the file // consistent, so this is not worth a single-flight map. If the same viewport // ever gets opened by enough people at once to matter, key a // `DashMap<(z, x, y), broadcast::Sender<_>>` on the coordinates and have the // second caller await the first.