//! Thumbnail cache for the grid view. //! //! Images decode in-process. Videos go through `ffmpeg`, which is optional. //! Nothing here writes to the database: the cache is disposable. use std::path::{Path, PathBuf}; use std::time::{Duration, SystemTime}; use api_types::FileKind; use fast_image_resize::images::Image as FirImage; use fast_image_resize::{FilterType, IntoImageView, ResizeAlg, ResizeOptions, Resizer}; use image::{DynamicImage, ImageReader, Limits, RgbImage}; use sha2::{Digest, Sha256}; use tokio::sync::Semaphore; /// Longest edge of a thumbnail, in pixels. The grid tile is 150 CSS px, so /// this still looks sharp on a 2x display. const MAX_EDGE: u32 = 256; /// WebP quality. At this size WebP is 18% smaller than JPEG at equal quality. /// AVIF and JXL are smaller still, but not every browser decodes them. const QUALITY: f32 = 80.0; /// An entry unused for this long is swept. Every hit refreshes the mtime, so /// this is "cold for a day", not "made a day ago". const TTL: Duration = Duration::from_secs(24 * 60 * 60); /// Only refresh an entry's mtime when it is at least this stale. Without the /// guard, one grid of 100 tiles is 100 metadata writes on every visit. const REFRESH_AFTER: Duration = Duration::from_secs(60 * 60); /// How often the sweeper runs. Each pass walks the whole cache. const SWEEP_EVERY: Duration = Duration::from_secs(60 * 60); /// Refuse to decode an image with more pixels than this. A 100 MP RGB buffer /// is 300 MB, so a crafted or merely enormous file could exhaust memory. const MAX_PIXELS: u64 = 100_000_000; /// Cap on the memory one decode may allocate. Backs up `MAX_PIXELS` for /// formats whose real cost the header does not reveal. const MAX_DECODE_BYTES: u64 = 512 * 1024 * 1024; /// Seconds to seek into a video before grabbing a frame. The first second is /// often black. const VIDEO_SEEK: &str = "1"; pub struct Thumbs { dir: PathBuf, /// Caps concurrent generation. Measured throughput flattens past 8 /// workers, and the cap keeps one big folder from starving the server. limit: Semaphore, /// Whether `ffmpeg` answered at startup. Only video thumbnails need it. ffmpeg: bool, } impl Thumbs { /// Prepare the cache folder. Also probes for `ffmpeg` once, so that a /// missing binary costs one failed spawn per process, not per video. pub async fn new(dir: PathBuf) -> std::io::Result { let ffmpeg = tokio::process::Command::new("ffmpeg") .arg("-version") .stdout(std::process::Stdio::null()) .stderr(std::process::Stdio::null()) .status() .await .is_ok_and(|s| s.success()); if !ffmpeg { tracing::info!("ffmpeg not found: video thumbnails are off"); } Self::with_ffmpeg(dir, ffmpeg).await } /// [`Thumbs::new`] with the probe result supplied. Lets a test take the /// no-ffmpeg branch on a machine that has ffmpeg. pub async fn with_ffmpeg(dir: PathBuf, ffmpeg: bool) -> std::io::Result { std::fs::create_dir_all(&dir)?; // Fail early and loudly rather than on the first thumbnail. let probe = dir.join(".writable"); std::fs::write(&probe, b"")?; let _ = std::fs::remove_file(&probe); let workers = std::thread::available_parallelism() .map(|n| n.get()) .unwrap_or(1) .min(8); tracing::info!(dir = %dir.display(), workers, ffmpeg, "thumbnail cache ready"); Ok(Self { dir, limit: Semaphore::new(workers), ffmpeg, }) } /// The thumbnail for `src`, making it if the cache does not have it. /// /// `None` means no thumbnail: an unsupported kind, a source too large to /// decode, a broken file, or a video with no ffmpeg. A failure is cached /// as an empty file, so a corrupt image is decoded only once. pub async fn get(&self, src: &Path, kind: FileKind, size: u64, mtime: i64) -> Option> { if !matches!(kind, FileKind::Image | FileKind::Video) { return None; } let path = self.entry_path(src, size, mtime); if let Some(hit) = read_hit(&path).await { return hit; } // The permit is held across generation only, never across the disk // read above: a cache hit must not queue behind a 36 MP decode. let _permit = self.limit.acquire().await.ok()?; // Another request may have finished it while we waited. if let Some(hit) = read_hit(&path).await { return hit; } let made = match kind { // Not a property of this file, so nothing is cached for it. // A marker here would outlive installing ffmpeg: every 404 // refreshes its mtime, so the sweeper would never drop it. FileKind::Video if !self.ffmpeg => return None, FileKind::Video => video_thumb(src).await, _ => { let src = src.to_path_buf(); tokio::task::spawn_blocking(move || image_thumb(&src)) .await .ok() .flatten() } }; if made.is_none() { tracing::debug!(src = %src.display(), "no thumbnail"); } // An empty file is the "we tried, it did not work" marker. store(&path, made.as_deref().unwrap_or(&[])).await; made } /// `//.webp`. /// /// The fan-out keeps one folder from collecting every thumbnail on the /// server, which makes the sweeper's `read_dir` walk cheap. fn entry_path(&self, src: &Path, size: u64, mtime: i64) -> PathBuf { let mut h = Sha256::new(); h.update(src.as_os_str().as_encoded_bytes()); h.update(size.to_le_bytes()); h.update(mtime.to_le_bytes()); // Output settings are part of the identity: changing them must not // serve thumbnails made under the old ones. h.update(MAX_EDGE.to_le_bytes()); h.update(QUALITY.to_le_bytes()); let hex = crate::hex(&h.finalize()[..16]); self.dir.join(&hex[..2]).join(format!("{}.webp", &hex[2..])) } } /// Read a cache entry, refreshing its mtime so the sweeper counts it as used. /// /// `None` = not cached. `Some(None)` = cached failure. async fn read_hit(path: &Path) -> Option>> { let bytes = tokio::fs::read(path).await.ok()?; let path = path.to_path_buf(); tokio::task::spawn_blocking(move || touch(&path)); Some((!bytes.is_empty()).then_some(bytes)) } /// Mark an entry as used. Its mtime is the last-used time the sweeper reads. /// /// Not atime: `relatime` updates it once a day and `noatime` mounts never do, /// so a hot thumbnail would look cold. `File::set_times` accepts a read-only /// handle. fn touch(path: &Path) { let Ok(file) = std::fs::File::open(path) else { return; }; let fresh = file .metadata() .and_then(|m| m.modified()) .is_ok_and(|m| m.elapsed().is_ok_and(|age| age < REFRESH_AFTER)); if fresh { return; } let _ = file.set_times(std::fs::FileTimes::new().set_modified(SystemTime::now())); } static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); /// A scratch name for one `store` call. /// /// Unique per call, not just per key. Two requests can generate the same /// uncached file at once, and a shared name lets one truncate the other's /// bytes before the rename publishes them as an empty failure marker. fn tmp_path(path: &Path) -> PathBuf { let seq = TMP_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed); path.with_extension(format!("{}.{seq}.tmp", std::process::id())) } /// Write an entry. Via a temporary file and a rename, so a reader never sees /// a half-written thumbnail, and a crash leaves no truncated one behind. async fn store(path: &Path, bytes: &[u8]) { let Some(parent) = path.parent() else { return }; if tokio::fs::create_dir_all(parent).await.is_err() { return; } let tmp = tmp_path(path); if tokio::fs::write(&tmp, bytes).await.is_err() { let _ = tokio::fs::remove_file(&tmp).await; return; } if tokio::fs::rename(&tmp, path).await.is_err() { let _ = tokio::fs::remove_file(&tmp).await; } } /// Decode, downscale and encode one image. Blocking; run it off the reactor. fn image_thumb(src: &Path) -> Option> { let reader = ImageReader::open(src).ok()?.with_guessed_format().ok()?; // Dimensions come from the header, so an oversized file is refused before // anything is allocated for it. let (w, h) = reader.into_dimensions().ok()?; if u64::from(w) * u64::from(h) > MAX_PIXELS { return None; } let mut reader = ImageReader::open(src).ok()?.with_guessed_format().ok()?; let mut limits = Limits::default(); limits.max_alloc = Some(MAX_DECODE_BYTES); reader.limits(limits); encode(&flatten(reader.decode().ok()?)) } /// Drop the alpha channel by compositing over white. /// /// `to_rgb8` alone keeps the colour behind the alpha, which for a transparent /// PNG is usually black. That turns a cut-out logo into a black tile. White /// matches an image viewer and reads in both themes. fn flatten(img: DynamicImage) -> DynamicImage { if !img.color().has_alpha() { return DynamicImage::ImageRgb8(img.into_rgb8()); } let src = img.into_rgba8(); let mut out = RgbImage::new(src.width(), src.height()); for (o, p) in out.pixels_mut().zip(src.pixels()) { // Straight (not premultiplied) alpha: out = src*a + white*(1-a). let a = u32::from(p.0[3]); let over = |c: u8| ((u32::from(c) * a + 255 * (255 - a)) / 255) as u8; *o = image::Rgb([over(p.0[0]), over(p.0[1]), over(p.0[2])]); } DynamicImage::ImageRgb8(out) } /// One frame from a video, via ffmpeg. The frame comes back as a PNG on /// stdout, already scaled down, so nothing full-size is decoded in-process. async fn video_thumb(src: &Path) -> Option> { let frame = match grab_frame(src, Some(VIDEO_SEEK)).await { Some(f) => Some(f), // A video shorter than the seek point yields no frame. Retry from the // start before giving up on it. None => grab_frame(src, None).await, }?; let img = image::load_from_memory_with_format(&frame, image::ImageFormat::Png).ok()?; encode(&flatten(img)) } async fn grab_frame(src: &Path, seek: Option<&str>) -> Option> { let mut cmd = tokio::process::Command::new("ffmpeg"); cmd.arg("-nostdin").args(["-v", "error"]); if let Some(s) = seek { // Before `-i`, so ffmpeg seeks on the container index instead of // decoding its way there. cmd.args(["-ss", s]); } let out = cmd .arg("-i") .arg(src) .args(["-frames:v", "1"]) .args([ "-vf", // `min` against the source size so a small clip is not enlarged, // which is what `encode` does for an image. &format!( "scale='min({MAX_EDGE},iw)':'min({MAX_EDGE},ih)'\ :force_original_aspect_ratio=decrease" ), ]) .args(["-f", "image2pipe", "-c:v", "png", "-"]) .stdin(std::process::Stdio::null()) .output() .await .ok()?; if !out.status.success() || out.stdout.is_empty() { // `-v error` writes the reason here. Without this a video that stops // working gives no clue why. tracing::debug!( src = %src.display(), status = %out.status, stderr = %String::from_utf8_lossy(&out.stderr).trim(), "ffmpeg produced no frame" ); return None; } Some(out.stdout) } /// Downscale to fit `MAX_EDGE` and encode as WebP. /// /// Lanczos3 through `fast_image_resize`: 46 ms against 720 ms for `image`'s /// own on a 36 MP source. `DynamicImage::thumbnail` is visibly softer. fn encode(img: &DynamicImage) -> Option> { let (w, h) = (img.width(), img.height()); let scale = (MAX_EDGE as f32 / w.max(h) as f32).min(1.0); let (dw, dh) = ( ((w as f32 * scale) as u32).max(1), ((h as f32 * scale) as u32).max(1), ); let mut dst = FirImage::new(dw, dh, img.pixel_type()?); Resizer::new() .resize( img, &mut dst, &ResizeOptions::new().resize_alg(ResizeAlg::Convolution(FilterType::Lanczos3)), ) .ok()?; let rgb = RgbImage::from_raw(dw, dh, dst.into_vec())?; Some( webp::Encoder::from_rgb(&rgb, dw, dh) .encode(QUALITY) .to_vec(), ) } /// Drop entries nobody has used for [`TTL`]. Runs until the process ends. pub async fn sweep_forever(dir: PathBuf) { loop { tokio::time::sleep(SWEEP_EVERY).await; let dir = dir.clone(); let removed = tokio::task::spawn_blocking(move || sweep(&dir)).await; if let Ok(n) = removed && n > 0 { tracing::debug!(removed = n, "swept cold thumbnails"); } } } /// One pass over the cache. Returns how many entries it removed. fn sweep(dir: &Path) -> usize { let Ok(buckets) = std::fs::read_dir(dir) else { return 0; }; let mut removed = 0; for bucket in buckets.flatten() { let Ok(entries) = std::fs::read_dir(bucket.path()) else { continue; }; for e in entries.flatten() { let cold = e .metadata() .and_then(|m| m.modified()) .is_ok_and(|m| m.elapsed().is_ok_and(|age| age > TTL)); if cold && std::fs::remove_file(e.path()).is_ok() { removed += 1; } } // An empty bucket is left in place: there are at most 256 of them and // the next thumbnail landing there would only recreate it. } removed } #[cfg(test)] mod tests { use super::*; /// Put a file in a bucket with a chosen age. fn aged(dir: &Path, name: &str, age: Duration) -> PathBuf { let bucket = dir.join("ab"); std::fs::create_dir_all(&bucket).unwrap(); let p = bucket.join(name); std::fs::write(&p, b"x").unwrap(); let when = SystemTime::now() - age; std::fs::File::options() .write(true) .open(&p) .unwrap() .set_times(std::fs::FileTimes::new().set_modified(when)) .unwrap(); p } #[test] fn two_writes_of_one_entry_use_different_scratch_names() { // A shared name would let the second `write` truncate the first one's // bytes between its write and its rename, publishing an empty file. // An empty entry reads back as a cached failure. let p = Path::new("/cache/ab/cd.webp"); assert_ne!(tmp_path(p), tmp_path(p)); assert_ne!(tmp_path(p), p.to_path_buf()); } #[test] fn the_sweeper_drops_only_cold_entries() { let dir = tempfile::tempdir().unwrap(); let hot = aged(dir.path(), "hot.webp", Duration::from_secs(60)); let warm = aged(dir.path(), "warm.webp", TTL - Duration::from_secs(600)); let cold = aged(dir.path(), "cold.webp", TTL + Duration::from_secs(600)); assert_eq!(sweep(dir.path()), 1); assert!(hot.exists()); assert!(warm.exists()); assert!(!cold.exists()); } #[test] fn a_hit_refreshes_a_stale_entry_only() { let dir = tempfile::tempdir().unwrap(); let stale = aged( dir.path(), "stale.webp", REFRESH_AFTER + Duration::from_secs(60), ); let fresh = aged(dir.path(), "fresh.webp", Duration::from_secs(30)); let fresh_before = std::fs::metadata(&fresh).unwrap().modified().unwrap(); touch(&stale); touch(&fresh); let age = std::fs::metadata(&stale) .unwrap() .modified() .unwrap() .elapsed() .unwrap(); assert!( age < Duration::from_secs(5), "stale entry was not refreshed" ); assert_eq!( std::fs::metadata(&fresh).unwrap().modified().unwrap(), fresh_before, "a fresh entry costs no write" ); } #[test] fn an_edit_changes_the_cache_key() { let t = Thumbs { dir: PathBuf::from("/cache"), limit: Semaphore::new(1), ffmpeg: false, }; let p = Path::new("/data/a.jpg"); assert_ne!(t.entry_path(p, 10, 1), t.entry_path(p, 10, 2)); assert_ne!(t.entry_path(p, 10, 1), t.entry_path(p, 11, 1)); assert_ne!( t.entry_path(p, 10, 1), t.entry_path(Path::new("/b.jpg"), 10, 1) ); assert_eq!(t.entry_path(p, 10, 1), t.entry_path(p, 10, 1)); } #[test] fn the_key_fans_out_into_a_bucket() { let t = Thumbs { dir: PathBuf::from("/cache"), limit: Semaphore::new(1), ffmpeg: false, }; let p = t.entry_path(Path::new("/data/a.jpg"), 1, 1); let bucket = p.parent().unwrap().file_name().unwrap().to_str().unwrap(); assert_eq!(bucket.len(), 2); assert!(bucket.chars().all(|c| c.is_ascii_hexdigit())); assert!(p.starts_with("/cache")); } }