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