thumb.rs
⎇
Raw
1//! Thumbnail cache for the grid view.
2//!
3//! Images decode in-process. Videos go through `ffmpeg`, which is optional.
4//! Nothing here writes to the database: the cache is disposable.
5
6use std::path::{Path, PathBuf};
7use std::time::{Duration, SystemTime};
8
9use api_types::FileKind;
10use fast_image_resize::images::Image as FirImage;
11use fast_image_resize::{FilterType, IntoImageView, ResizeAlg, ResizeOptions, Resizer};
12use image::{DynamicImage, ImageReader, Limits, RgbImage};
13use sha2::{Digest, Sha256};
14use tokio::sync::Semaphore;
15
16/// Longest edge of a thumbnail, in pixels. The grid tile is 150 CSS px, so
17/// this still looks sharp on a 2x display.
18const MAX_EDGE: u32 = 256;
19
20/// WebP quality. At this size WebP is 18% smaller than JPEG at equal quality.
21/// AVIF and JXL are smaller still, but not every browser decodes them.
22const QUALITY: f32 = 80.0;
23
24/// An entry unused for this long is swept. Every hit refreshes the mtime, so
25/// this is "cold for a day", not "made a day ago".
26const TTL: Duration = Duration::from_secs(24 * 60 * 60);
27
28/// Only refresh an entry's mtime when it is at least this stale. Without the
29/// guard, one grid of 100 tiles is 100 metadata writes on every visit.
30const REFRESH_AFTER: Duration = Duration::from_secs(60 * 60);
31
32/// How often the sweeper runs. Each pass walks the whole cache.
33const SWEEP_EVERY: Duration = Duration::from_secs(60 * 60);
34
35/// Refuse to decode an image with more pixels than this. A 100 MP RGB buffer
36/// is 300 MB, so a crafted or merely enormous file could exhaust memory.
37const MAX_PIXELS: u64 = 100_000_000;
38
39/// Cap on the memory one decode may allocate. Backs up `MAX_PIXELS` for
40/// formats whose real cost the header does not reveal.
41const MAX_DECODE_BYTES: u64 = 512 * 1024 * 1024;
42
43/// Seconds to seek into a video before grabbing a frame. The first second is
44/// often black.
45const VIDEO_SEEK: &str = "1";
46
47pub struct Thumbs {
48 dir: PathBuf,
49 /// Caps concurrent generation. Measured throughput flattens past 8
50 /// workers, and the cap keeps one big folder from starving the server.
51 limit: Semaphore,
52 /// Whether `ffmpeg` answered at startup. Only video thumbnails need it.
53 ffmpeg: bool,
54}
55
56impl Thumbs {
57 /// Prepare the cache folder. Also probes for `ffmpeg` once, so that a
58 /// missing binary costs one failed spawn per process, not per video.
59 pub async fn new(dir: PathBuf) -> std::io::Result<Self> {
60 let ffmpeg = tokio::process::Command::new("ffmpeg")
61 .arg("-version")
62 .stdout(std::process::Stdio::null())
63 .stderr(std::process::Stdio::null())
64 .status()
65 .await
66 .is_ok_and(|s| s.success());
67 if !ffmpeg {
68 tracing::info!("ffmpeg not found: video thumbnails are off");
69 }
70 Self::with_ffmpeg(dir, ffmpeg).await
71 }
72
73 /// [`Thumbs::new`] with the probe result supplied. Lets a test take the
74 /// no-ffmpeg branch on a machine that has ffmpeg.
75 pub async fn with_ffmpeg(dir: PathBuf, ffmpeg: bool) -> std::io::Result<Self> {
76 std::fs::create_dir_all(&dir)?;
77 // Fail early and loudly rather than on the first thumbnail.
78 let probe = dir.join(".writable");
79 std::fs::write(&probe, b"")?;
80 let _ = std::fs::remove_file(&probe);
81
82 let workers = std::thread::available_parallelism()
83 .map(|n| n.get())
84 .unwrap_or(1)
85 .min(8);
86 tracing::info!(dir = %dir.display(), workers, ffmpeg, "thumbnail cache ready");
87 Ok(Self {
88 dir,
89 limit: Semaphore::new(workers),
90 ffmpeg,
91 })
92 }
93
94 /// The thumbnail for `src`, making it if the cache does not have it.
95 ///
96 /// `None` means no thumbnail: an unsupported kind, a source too large to
97 /// decode, a broken file, or a video with no ffmpeg. A failure is cached
98 /// as an empty file, so a corrupt image is decoded only once.
99 pub async fn get(&self, src: &Path, kind: FileKind, size: u64, mtime: i64) -> Option<Vec<u8>> {
100 if !matches!(kind, FileKind::Image | FileKind::Video) {
101 return None;
102 }
103 let path = self.entry_path(src, size, mtime);
104
105 if let Some(hit) = read_hit(&path).await {
106 return hit;
107 }
108
109 // The permit is held across generation only, never across the disk
110 // read above: a cache hit must not queue behind a 36 MP decode.
111 let _permit = self.limit.acquire().await.ok()?;
112 // Another request may have finished it while we waited.
113 if let Some(hit) = read_hit(&path).await {
114 return hit;
115 }
116
117 let made = match kind {
118 // Not a property of this file, so nothing is cached for it.
119 // A marker here would outlive installing ffmpeg: every 404
120 // refreshes its mtime, so the sweeper would never drop it.
121 FileKind::Video if !self.ffmpeg => return None,
122 FileKind::Video => video_thumb(src).await,
123 _ => {
124 let src = src.to_path_buf();
125 tokio::task::spawn_blocking(move || image_thumb(&src))
126 .await
127 .ok()
128 .flatten()
129 }
130 };
131 if made.is_none() {
132 tracing::debug!(src = %src.display(), "no thumbnail");
133 }
134 // An empty file is the "we tried, it did not work" marker.
135 store(&path, made.as_deref().unwrap_or(&[])).await;
136 made
137 }
138
139 /// `<cache>/<first 2 hex>/<rest>.webp`.
140 ///
141 /// The fan-out keeps one folder from collecting every thumbnail on the
142 /// server, which makes the sweeper's `read_dir` walk cheap.
143 fn entry_path(&self, src: &Path, size: u64, mtime: i64) -> PathBuf {
144 let mut h = Sha256::new();
145 h.update(src.as_os_str().as_encoded_bytes());
146 h.update(size.to_le_bytes());
147 h.update(mtime.to_le_bytes());
148 // Output settings are part of the identity: changing them must not
149 // serve thumbnails made under the old ones.
150 h.update(MAX_EDGE.to_le_bytes());
151 h.update(QUALITY.to_le_bytes());
152 let hex = crate::hex(&h.finalize()[..16]);
153 self.dir.join(&hex[..2]).join(format!("{}.webp", &hex[2..]))
154 }
155}
156
157/// Read a cache entry, refreshing its mtime so the sweeper counts it as used.
158///
159/// `None` = not cached. `Some(None)` = cached failure.
160async fn read_hit(path: &Path) -> Option<Option<Vec<u8>>> {
161 let bytes = tokio::fs::read(path).await.ok()?;
162 let path = path.to_path_buf();
163 tokio::task::spawn_blocking(move || touch(&path));
164 Some((!bytes.is_empty()).then_some(bytes))
165}
166
167/// Mark an entry as used. Its mtime is the last-used time the sweeper reads.
168///
169/// Not atime: `relatime` updates it once a day and `noatime` mounts never do,
170/// so a hot thumbnail would look cold. `File::set_times` accepts a read-only
171/// handle.
172fn touch(path: &Path) {
173 let Ok(file) = std::fs::File::open(path) else {
174 return;
175 };
176 let fresh = file
177 .metadata()
178 .and_then(|m| m.modified())
179 .is_ok_and(|m| m.elapsed().is_ok_and(|age| age < REFRESH_AFTER));
180 if fresh {
181 return;
182 }
183 let _ = file.set_times(std::fs::FileTimes::new().set_modified(SystemTime::now()));
184}
185
186static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
187
188/// A scratch name for one `store` call.
189///
190/// Unique per call, not just per key. Two requests can generate the same
191/// uncached file at once, and a shared name lets one truncate the other's
192/// bytes before the rename publishes them as an empty failure marker.
193fn tmp_path(path: &Path) -> PathBuf {
194 let seq = TMP_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
195 path.with_extension(format!("{}.{seq}.tmp", std::process::id()))
196}
197
198/// Write an entry. Via a temporary file and a rename, so a reader never sees
199/// a half-written thumbnail, and a crash leaves no truncated one behind.
200async fn store(path: &Path, bytes: &[u8]) {
201 let Some(parent) = path.parent() else { return };
202 if tokio::fs::create_dir_all(parent).await.is_err() {
203 return;
204 }
205 let tmp = tmp_path(path);
206 if tokio::fs::write(&tmp, bytes).await.is_err() {
207 let _ = tokio::fs::remove_file(&tmp).await;
208 return;
209 }
210 if tokio::fs::rename(&tmp, path).await.is_err() {
211 let _ = tokio::fs::remove_file(&tmp).await;
212 }
213}
214
215/// Decode, downscale and encode one image. Blocking; run it off the reactor.
216fn image_thumb(src: &Path) -> Option<Vec<u8>> {
217 let reader = ImageReader::open(src).ok()?.with_guessed_format().ok()?;
218 // Dimensions come from the header, so an oversized file is refused before
219 // anything is allocated for it.
220 let (w, h) = reader.into_dimensions().ok()?;
221 if u64::from(w) * u64::from(h) > MAX_PIXELS {
222 return None;
223 }
224
225 let mut reader = ImageReader::open(src).ok()?.with_guessed_format().ok()?;
226 let mut limits = Limits::default();
227 limits.max_alloc = Some(MAX_DECODE_BYTES);
228 reader.limits(limits);
229 encode(&flatten(reader.decode().ok()?))
230}
231
232/// Drop the alpha channel by compositing over white.
233///
234/// `to_rgb8` alone keeps the colour behind the alpha, which for a transparent
235/// PNG is usually black. That turns a cut-out logo into a black tile. White
236/// matches an image viewer and reads in both themes.
237fn flatten(img: DynamicImage) -> DynamicImage {
238 if !img.color().has_alpha() {
239 return DynamicImage::ImageRgb8(img.into_rgb8());
240 }
241 let src = img.into_rgba8();
242 let mut out = RgbImage::new(src.width(), src.height());
243 for (o, p) in out.pixels_mut().zip(src.pixels()) {
244 // Straight (not premultiplied) alpha: out = src*a + white*(1-a).
245 let a = u32::from(p.0[3]);
246 let over = |c: u8| ((u32::from(c) * a + 255 * (255 - a)) / 255) as u8;
247 *o = image::Rgb([over(p.0[0]), over(p.0[1]), over(p.0[2])]);
248 }
249 DynamicImage::ImageRgb8(out)
250}
251
252/// One frame from a video, via ffmpeg. The frame comes back as a PNG on
253/// stdout, already scaled down, so nothing full-size is decoded in-process.
254async fn video_thumb(src: &Path) -> Option<Vec<u8>> {
255 let frame = match grab_frame(src, Some(VIDEO_SEEK)).await {
256 Some(f) => Some(f),
257 // A video shorter than the seek point yields no frame. Retry from the
258 // start before giving up on it.
259 None => grab_frame(src, None).await,
260 }?;
261 let img = image::load_from_memory_with_format(&frame, image::ImageFormat::Png).ok()?;
262 encode(&flatten(img))
263}
264
265async fn grab_frame(src: &Path, seek: Option<&str>) -> Option<Vec<u8>> {
266 let mut cmd = tokio::process::Command::new("ffmpeg");
267 cmd.arg("-nostdin").args(["-v", "error"]);
268 if let Some(s) = seek {
269 // Before `-i`, so ffmpeg seeks on the container index instead of
270 // decoding its way there.
271 cmd.args(["-ss", s]);
272 }
273 let out = cmd
274 .arg("-i")
275 .arg(src)
276 .args(["-frames:v", "1"])
277 .args([
278 "-vf",
279 // `min` against the source size so a small clip is not enlarged,
280 // which is what `encode` does for an image.
281 &format!(
282 "scale='min({MAX_EDGE},iw)':'min({MAX_EDGE},ih)'\
283 :force_original_aspect_ratio=decrease"
284 ),
285 ])
286 .args(["-f", "image2pipe", "-c:v", "png", "-"])
287 .stdin(std::process::Stdio::null())
288 .output()
289 .await
290 .ok()?;
291 if !out.status.success() || out.stdout.is_empty() {
292 // `-v error` writes the reason here. Without this a video that stops
293 // working gives no clue why.
294 tracing::debug!(
295 src = %src.display(),
296 status = %out.status,
297 stderr = %String::from_utf8_lossy(&out.stderr).trim(),
298 "ffmpeg produced no frame"
299 );
300 return None;
301 }
302 Some(out.stdout)
303}
304
305/// Downscale to fit `MAX_EDGE` and encode as WebP.
306///
307/// Lanczos3 through `fast_image_resize`: 46 ms against 720 ms for `image`'s
308/// own on a 36 MP source. `DynamicImage::thumbnail` is visibly softer.
309fn encode(img: &DynamicImage) -> Option<Vec<u8>> {
310 let (w, h) = (img.width(), img.height());
311 let scale = (MAX_EDGE as f32 / w.max(h) as f32).min(1.0);
312 let (dw, dh) = (
313 ((w as f32 * scale) as u32).max(1),
314 ((h as f32 * scale) as u32).max(1),
315 );
316
317 let mut dst = FirImage::new(dw, dh, img.pixel_type()?);
318 Resizer::new()
319 .resize(
320 img,
321 &mut dst,
322 &ResizeOptions::new().resize_alg(ResizeAlg::Convolution(FilterType::Lanczos3)),
323 )
324 .ok()?;
325 let rgb = RgbImage::from_raw(dw, dh, dst.into_vec())?;
326 Some(
327 webp::Encoder::from_rgb(&rgb, dw, dh)
328 .encode(QUALITY)
329 .to_vec(),
330 )
331}
332
333/// Drop entries nobody has used for [`TTL`]. Runs until the process ends.
334pub async fn sweep_forever(dir: PathBuf) {
335 loop {
336 tokio::time::sleep(SWEEP_EVERY).await;
337 let dir = dir.clone();
338 let removed = tokio::task::spawn_blocking(move || sweep(&dir)).await;
339 if let Ok(n) = removed
340 && n > 0
341 {
342 tracing::debug!(removed = n, "swept cold thumbnails");
343 }
344 }
345}
346
347/// One pass over the cache. Returns how many entries it removed.
348fn sweep(dir: &Path) -> usize {
349 let Ok(buckets) = std::fs::read_dir(dir) else {
350 return 0;
351 };
352 let mut removed = 0;
353 for bucket in buckets.flatten() {
354 let Ok(entries) = std::fs::read_dir(bucket.path()) else {
355 continue;
356 };
357 for e in entries.flatten() {
358 let cold = e
359 .metadata()
360 .and_then(|m| m.modified())
361 .is_ok_and(|m| m.elapsed().is_ok_and(|age| age > TTL));
362 if cold && std::fs::remove_file(e.path()).is_ok() {
363 removed += 1;
364 }
365 }
366 // An empty bucket is left in place: there are at most 256 of them and
367 // the next thumbnail landing there would only recreate it.
368 }
369 removed
370}
371
372#[cfg(test)]
373mod tests {
374 use super::*;
375
376 /// Put a file in a bucket with a chosen age.
377 fn aged(dir: &Path, name: &str, age: Duration) -> PathBuf {
378 let bucket = dir.join("ab");
379 std::fs::create_dir_all(&bucket).unwrap();
380 let p = bucket.join(name);
381 std::fs::write(&p, b"x").unwrap();
382 let when = SystemTime::now() - age;
383 std::fs::File::options()
384 .write(true)
385 .open(&p)
386 .unwrap()
387 .set_times(std::fs::FileTimes::new().set_modified(when))
388 .unwrap();
389 p
390 }
391
392 #[test]
393 fn two_writes_of_one_entry_use_different_scratch_names() {
394 // A shared name would let the second `write` truncate the first one's
395 // bytes between its write and its rename, publishing an empty file.
396 // An empty entry reads back as a cached failure.
397 let p = Path::new("/cache/ab/cd.webp");
398 assert_ne!(tmp_path(p), tmp_path(p));
399 assert_ne!(tmp_path(p), p.to_path_buf());
400 }
401
402 #[test]
403 fn the_sweeper_drops_only_cold_entries() {
404 let dir = tempfile::tempdir().unwrap();
405 let hot = aged(dir.path(), "hot.webp", Duration::from_secs(60));
406 let warm = aged(dir.path(), "warm.webp", TTL - Duration::from_secs(600));
407 let cold = aged(dir.path(), "cold.webp", TTL + Duration::from_secs(600));
408
409 assert_eq!(sweep(dir.path()), 1);
410 assert!(hot.exists());
411 assert!(warm.exists());
412 assert!(!cold.exists());
413 }
414
415 #[test]
416 fn a_hit_refreshes_a_stale_entry_only() {
417 let dir = tempfile::tempdir().unwrap();
418 let stale = aged(
419 dir.path(),
420 "stale.webp",
421 REFRESH_AFTER + Duration::from_secs(60),
422 );
423 let fresh = aged(dir.path(), "fresh.webp", Duration::from_secs(30));
424 let fresh_before = std::fs::metadata(&fresh).unwrap().modified().unwrap();
425
426 touch(&stale);
427 touch(&fresh);
428
429 let age = std::fs::metadata(&stale)
430 .unwrap()
431 .modified()
432 .unwrap()
433 .elapsed()
434 .unwrap();
435 assert!(
436 age < Duration::from_secs(5),
437 "stale entry was not refreshed"
438 );
439 assert_eq!(
440 std::fs::metadata(&fresh).unwrap().modified().unwrap(),
441 fresh_before,
442 "a fresh entry costs no write"
443 );
444 }
445
446 #[test]
447 fn an_edit_changes_the_cache_key() {
448 let t = Thumbs {
449 dir: PathBuf::from("/cache"),
450 limit: Semaphore::new(1),
451 ffmpeg: false,
452 };
453 let p = Path::new("/data/a.jpg");
454 assert_ne!(t.entry_path(p, 10, 1), t.entry_path(p, 10, 2));
455 assert_ne!(t.entry_path(p, 10, 1), t.entry_path(p, 11, 1));
456 assert_ne!(
457 t.entry_path(p, 10, 1),
458 t.entry_path(Path::new("/b.jpg"), 10, 1)
459 );
460 assert_eq!(t.entry_path(p, 10, 1), t.entry_path(p, 10, 1));
461 }
462
463 #[test]
464 fn the_key_fans_out_into_a_bucket() {
465 let t = Thumbs {
466 dir: PathBuf::from("/cache"),
467 limit: Semaphore::new(1),
468 ffmpeg: false,
469 };
470 let p = t.entry_path(Path::new("/data/a.jpg"), 1, 1);
471 let bucket = p.parent().unwrap().file_name().unwrap().to_str().unwrap();
472 assert_eq!(bucket.len(), 2);
473 assert!(bucket.chars().all(|c| c.is_ascii_hexdigit()));
474 assert!(p.starts_with("/cache"));
475 }
476}
477