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