From d6ddd0703bf2cffefa480a778d36b284e10c5955 Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Sun, 9 Aug 2026 20:38:12 +0200 Subject: [PATCH] Add dr-thumbs: a sharded, syncable thumbnail store MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A thumbnail is the one derived artefact worth sending over the wire: it costs a range fetch plus a decode to produce and is identical for every client looking at the same file. A second device that downloads a shard gets a full grid without fetching a byte of RAW. Sharded at 25 MB, filled sequentially. The cap is about sync granularity, not SQLite's limits — one growing database means every client re-downloads it whenever a single thumbnail is added, whereas with sequential fill only the newest shard is ever dirty and sealed shards are safe to cache forever. Stored JPEG-encoded rather than as raw RGBA: a 256px RGBA buffer is ~256 KB against ~20 KB encoded, and that 13x is transfer cost on every client. Keyed on Nextcloud's oc:fileid, stable across server-side rename and move. Assisted-by: LLM --- Cargo.lock | 49 +++ Cargo.toml | 14 + core/dr-thumbs/Cargo.toml | 14 + core/dr-thumbs/src/codec.rs | 163 ++++++++ core/dr-thumbs/src/error.rs | 21 + core/dr-thumbs/src/lib.rs | 750 ++++++++++++++++++++++++++++++++++++ 6 files changed, 1011 insertions(+) create mode 100644 core/dr-thumbs/Cargo.toml create mode 100644 core/dr-thumbs/src/codec.rs create mode 100644 core/dr-thumbs/src/error.rs create mode 100644 core/dr-thumbs/src/lib.rs diff --git a/Cargo.lock b/Cargo.lock index aa9a6a4..241ca90 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1262,6 +1262,7 @@ dependencies = [ "dr-types", "log", "rusqlite", + "serde_json", "thiserror 2.0.20", ] @@ -1351,10 +1352,23 @@ dependencies = [ "url", ] +[[package]] +name = "dr-thumbs" +version = "0.1.0" +dependencies = [ + "dr-types", + "jpeg-encoder", + "log", + "rusqlite", + "thiserror 2.0.20", + "zune-jpeg 0.4.21", +] + [[package]] name = "dr-types" version = "0.1.0" dependencies = [ + "serde", "thiserror 2.0.20", ] @@ -1363,17 +1377,21 @@ name = "dr-ui" version = "0.1.0" dependencies = [ "anyhow", + "dr-catalog", "dr-decode", "dr-gpu", "dr-pipeline", "dr-plat", "dr-sync", "dr-sync-nextcloud", + "dr-thumbs", "dr-types", "log", "pollster", "reqwest", + "rusqlite", "serde_json", + "serde_norway", "slint", "slint-build", "tokio", @@ -3115,6 +3133,12 @@ dependencies = [ "libc", ] +[[package]] +name = "jpeg-encoder" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a0370574b86f7eca156b9f298392b5e69a23f8c86f3f865add60bbc2e79467a6" + [[package]] name = "js-sys" version = "0.3.104" @@ -5333,6 +5357,12 @@ dependencies = [ "unicode-script", ] +[[package]] +name = "ryu" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" + [[package]] name = "same-file" version = "1.0.6" @@ -5473,6 +5503,19 @@ dependencies = [ "zmij", ] +[[package]] +name = "serde_norway" +version = "0.9.42" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e408f29489b5fd500fab51ff1484fc859bb655f32c671f307dcd733b72e8168c" +dependencies = [ + "indexmap", + "itoa", + "ryu", + "serde", + "unsafe-libyaml-norway", +] + [[package]] name = "serde_repr" version = "0.1.21" @@ -6554,6 +6597,12 @@ version = "0.2.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ebc1c04c71510c7f702b52b7c350734c9ff1295c464a03335b00bb84fc54f853" +[[package]] +name = "unsafe-libyaml-norway" +version = "0.2.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b39abd59bf32521c7f2301b52d05a6a2c975b6003521cbd0c6dc1582f0a22104" + [[package]] name = "untrusted" version = "0.9.0" diff --git a/Cargo.toml b/Cargo.toml index dce9a32..a83fc71 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -3,6 +3,7 @@ resolver = "2" members = [ "core/dr-types", "core/dr-catalog", + "core/dr-thumbs", "core/dr-decode", "core/dr-gpu", "core/dr-lens", @@ -26,6 +27,7 @@ repository = "https://github.com/dtourolle/DarkRoom" # Internal dr-types = { path = "core/dr-types" } dr-catalog = { path = "core/dr-catalog" } +dr-thumbs = { path = "core/dr-thumbs" } dr-decode = { path = "core/dr-decode" } dr-gpu = { path = "core/dr-gpu" } dr-lens = { path = "core/dr-lens" } @@ -40,6 +42,13 @@ wgpu = "23" slint = { version = "1.9", default-features = false } slint-build = "1.9" +# UI token codegen (S2): style.yaml -> theme.slint. serde_yaml was deprecated +# by its maintainer in 2024 and serde_yml, the first fork, has since been +# deprecated too; serde_norway is the fork still receiving releases. Its +# mappings preserve insertion order, which is what lets the generated Slint +# keep the token ordering the YAML author chose. +serde_norway = "0.9" + # Foundations anyhow = "1" thiserror = "2" @@ -87,6 +96,11 @@ rusqlite = { version = "0.40", features = ["bundled", "backup"] } rawler = "0.7" zune-jpeg = "0.4.21" +# Thumbnails are stored encoded, not as raw RGBA: a 256px RGBA buffer is +# ~256 KB against ~20 KB as JPEG, and the store syncs to Nextcloud where that +# 13× is transfer cost on every client. Pure Rust, no C dependency — the same +# criterion behind the TLS and SQLite choices above. +jpeg-encoder = "0.7" bytemuck = { version = "1", features = ["derive"] } # Lens correction profiles. A pure-Rust port of Lensfun rather than a binding diff --git a/core/dr-thumbs/Cargo.toml b/core/dr-thumbs/Cargo.toml new file mode 100644 index 0000000..8d9bd61 --- /dev/null +++ b/core/dr-thumbs/Cargo.toml @@ -0,0 +1,14 @@ +[package] +name = "dr-thumbs" +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true + +[dependencies] +dr-types.workspace = true +rusqlite.workspace = true +jpeg-encoder.workspace = true +zune-jpeg.workspace = true +thiserror.workspace = true +log.workspace = true diff --git a/core/dr-thumbs/src/codec.rs b/core/dr-thumbs/src/codec.rs new file mode 100644 index 0000000..4e276b0 --- /dev/null +++ b/core/dr-thumbs/src/codec.rs @@ -0,0 +1,163 @@ +//! TRACES: FR-CAT-3 | NFR-RES-4 +//! Encoding and decoding stored thumbnails. +//! +//! Thumbnails are stored as JPEG rather than raw pixels. A 256×170 RGBA buffer +//! is ~174 KB; the same image as JPEG is ~15 KB. That ratio decides whether a +//! 25 MB shard holds a few hundred thumbnails or a few thousand — and since +//! shards sync to Nextcloud, it is transfer cost paid by every client, not +//! just disk. + +use crate::error::ThumbError; + +/// JPEG quality for stored thumbnails. +/// +/// 82 is the knee: visible artefacts start below it, and above it size climbs +/// with no benefit at 256px. These are browsing aids — a culling decision is +/// made against raw data (FR-CULL-3), never against one of these. +const QUALITY: u8 = 82; + +/// Encode RGBA pixels to JPEG. +/// +/// Alpha is discarded: a camera preview has none, and JPEG cannot carry it. +pub fn encode_rgba(width: u32, height: u32, rgba: &[u8]) -> Result, ThumbError> { + let expected = (width as usize) + .checked_mul(height as usize) + .and_then(|n| n.checked_mul(4)) + .ok_or_else(|| ThumbError::Io("thumbnail dimensions overflow".into()))?; + + if rgba.len() < expected { + return Err(ThumbError::Io(format!( + "thumbnail buffer is {} bytes, expected {expected}", + rgba.len() + ))); + } + + let mut out = Vec::new(); + let encoder = jpeg_encoder::Encoder::new(&mut out, QUALITY); + encoder + .encode( + &rgba[..expected], + width as u16, + height as u16, + jpeg_encoder::ColorType::Rgba, + ) + .map_err(|e| ThumbError::Io(e.to_string()))?; + + Ok(out) +} + +/// Decode a stored JPEG back to RGBA for display. +pub fn decode_rgba(bytes: &[u8]) -> Result<(u32, u32, Vec), ThumbError> { + let mut d = zune_jpeg::JpegDecoder::new(bytes); + let pixels = d + .decode() + .map_err(|e| ThumbError::Io(format!("decoding stored thumbnail: {e}")))?; + let info = d + .info() + .ok_or_else(|| ThumbError::Io("stored thumbnail has no header".into()))?; + + let (w, h) = (info.width as u32, info.height as u32); + let px = (w as usize) * (h as usize); + + // zune yields RGB for a colour JPEG and one channel for a greyscale one. + // A greyscale camera preview is rare but real, and treating it as RGB + // would produce a garbled cell rather than a grey one. + let rgba = match pixels.len() / px.max(1) { + 3 => { + let mut out = Vec::with_capacity(px * 4); + for c in pixels.chunks_exact(3) { + out.extend_from_slice(&[c[0], c[1], c[2], 255]); + } + out + } + 1 => { + let mut out = Vec::with_capacity(px * 4); + for g in &pixels { + out.extend_from_slice(&[*g, *g, *g, 255]); + } + out + } + _ => { + return Err(ThumbError::Io(format!( + "unexpected {} bytes for {w}×{h}", + pixels.len() + ))) + } + }; + + Ok((w, h, rgba)) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn gradient(w: u32, h: u32) -> Vec { + let mut v = Vec::with_capacity((w * h * 4) as usize); + for y in 0..h { + for x in 0..w { + v.extend_from_slice(&[(x * 4) as u8, (y * 4) as u8, 128, 255]); + } + } + v + } + + #[test] + fn encoding_then_decoding_preserves_dimensions() { + let bytes = encode_rgba(64, 48, &gradient(64, 48)).unwrap(); + let (w, h, rgba) = decode_rgba(&bytes).unwrap(); + assert_eq!((w, h), (64, 48)); + assert_eq!(rgba.len(), 64 * 48 * 4); + } + + #[test] + fn encoding_is_far_smaller_than_raw() { + // The entire reason for encoding: shards sync, so this ratio is + // transfer cost on every client. + let raw = gradient(256, 170); + let encoded = encode_rgba(256, 170, &raw).unwrap(); + assert!( + encoded.len() * 5 < raw.len(), + "expected a large saving, got {} vs {}", + encoded.len(), + raw.len() + ); + } + + #[test] + fn decoded_pixels_are_opaque() { + let bytes = encode_rgba(16, 16, &gradient(16, 16)).unwrap(); + let (_, _, rgba) = decode_rgba(&bytes).unwrap(); + assert!(rgba.chunks_exact(4).all(|p| p[3] == 255)); + } + + #[test] + fn a_short_buffer_errors_rather_than_reading_past_it() { + // Decoders handle untrusted input (NFR-SEC-1); a truncated preview + // must not become an out-of-bounds read. + assert!(encode_rgba(64, 64, &[0u8; 16]).is_err()); + } + + #[test] + fn absurd_dimensions_do_not_overflow() { + assert!(encode_rgba(u32::MAX, u32::MAX, &[0u8; 4]).is_err()); + } + + #[test] + fn corrupt_stored_bytes_error_rather_than_panic() { + assert!(decode_rgba(&[0xFF, 0xD8, 0xFF, 0x00, 0x01]).is_err()); + assert!(decode_rgba(&[]).is_err()); + } + + #[test] + fn roughly_round_trips_the_image() { + // Lossy, so exact equality is wrong to assert — but a mid-grey pixel + // must not come back as black or white. + let w = 32; + let src: Vec = std::iter::repeat_n([128u8, 128, 128, 255], (w * w) as usize) + .flatten() + .collect(); + let (_, _, out) = decode_rgba(&encode_rgba(w, w, &src).unwrap()).unwrap(); + assert!(out[0].abs_diff(128) < 12, "got {}", out[0]); + } +} diff --git a/core/dr-thumbs/src/error.rs b/core/dr-thumbs/src/error.rs new file mode 100644 index 0000000..02bbbdc --- /dev/null +++ b/core/dr-thumbs/src/error.rs @@ -0,0 +1,21 @@ +//! TRACES: NFR-ARCH-4 +//! Thumbnail store errors. +//! +//! A thumbnail is a cache entry, never authoritative — every failure here is +//! recoverable by regenerating from the source, so nothing in this module is +//! fatal to the app. + +/// Something went wrong reading or writing the store. +#[derive(Debug, thiserror::Error)] +pub enum ThumbError { + #[error("sqlite: {0}")] + Sqlite(#[from] rusqlite::Error), + + /// The index names a shard whose file is absent — a partial sync, or a + /// manual deletion. The entry is treated as a miss and regenerated. + #[error("shard {0} is missing")] + MissingShard(u32), + + #[error("io: {0}")] + Io(String), +} diff --git a/core/dr-thumbs/src/lib.rs b/core/dr-thumbs/src/lib.rs new file mode 100644 index 0000000..712c588 --- /dev/null +++ b/core/dr-thumbs/src/lib.rs @@ -0,0 +1,750 @@ +//! TRACES: FR-CAT-3 | FR-NC-3 | NFR-RES-4 +//! A shared, syncable thumbnail store. +//! +//! # Why thumbnails sync when the catalog mostly does not +//! +//! A thumbnail is the one derived artefact worth sending over the wire. It is +//! expensive to produce — a range fetch plus a decode, per image — and +//! identical for every client looking at the same file. A second device that +//! can download a shard gets a full grid without re-fetching a byte of RAW. +//! On the reference library that is the difference between a usable tablet and +//! one that spends an hour filling in cells it could have been handed. +//! +//! This does not make thumbnails authoritative. Losing the store costs +//! regeneration, nothing more; sidecars remain the only thing in the trust +//! path (ARCH §6.12). +//! +//! # Shape +//! +//! ```text +//! thumbs/ +//! index.sqlite fileid → shard, plus size accounting +//! shard-0000.sqlite ≤ 25 MB of blobs +//! shard-0001.sqlite +//! ``` +//! +//! **Sharded, and sized small on purpose.** The cap is not about SQLite's +//! limits — it is about sync granularity. A single growing database means +//! every client re-downloads the whole thing whenever one thumbnail is added. +//! With sequential fill only the newest shard is ever dirty, so a client that +//! is up to date transfers one small file. Sealed shards never change again, +//! which also makes them safe to cache indefinitely. +//! +//! # Identity +//! +//! Keyed on Nextcloud's `oc:fileid`, which is stable across server-side rename +//! and move (FR-NC-5) and which every client already holds from PROPFIND. The +//! consequence, accepted deliberately: shards are account-scoped, so the same +//! photograph on two different servers is thumbnailed twice. + +use std::path::{Path, PathBuf}; + +use rusqlite::Connection; + +pub mod codec; +pub mod error; + +pub use codec::{decode_rgba, encode_rgba}; +pub use error::ThumbError; + +/// Maximum bytes per shard before a new one is started. +/// +/// 25 MB, chosen for *sync* cost rather than storage: it bounds what a client +/// re-downloads when the active shard changes, and keeps a stalled transfer +/// cheap to retry. At ~20 KB per 256px JPEG that is roughly 1,200 thumbnails +/// per shard, so the 17k-image reference library lands in ~14 shards. +pub const SHARD_MAX_BYTES: u64 = 25 * 1024 * 1024; + +/// Long edge of a stored thumbnail. +pub const THUMBNAIL_EDGE: u32 = 256; + +/// A thumbnail's bytes and dimensions. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Thumbnail { + pub width: u32, + pub height: u32, + /// Encoded image bytes — JPEG. Stored encoded, not raw RGBA: a 256×256 + /// RGBA buffer is 256 KB where its JPEG is ~20 KB, and a store that shipped + /// raw pixels would be an order of magnitude more sync traffic. + pub bytes: Vec, +} + +/// The sharded store. +pub struct ThumbStore { + dir: PathBuf, + index: Connection, +} + +impl ThumbStore { + /// Open or create a store under `dir`. + pub fn open(dir: &Path) -> Result { + std::fs::create_dir_all(dir).map_err(|e| ThumbError::Io(e.to_string()))?; + + let index = Connection::open(dir.join("index.sqlite"))?; + index.pragma_update(None, "journal_mode", "WAL")?; + index.pragma_update(None, "synchronous", "NORMAL")?; + index.execute_batch(INDEX_SCHEMA)?; + + Ok(Self { + dir: dir.to_path_buf(), + index, + }) + } + + /// Fetch a thumbnail by file id. + pub fn get(&self, file_id: u64) -> Result, ThumbError> { + let shard: Option = self + .index + .query_row( + "SELECT shard FROM entries WHERE file_id = ?1", + [file_id as i64], + |r| r.get(0), + ) + .ok(); + + let Some(shard) = shard else { + return Ok(None); + }; + + let conn = self.open_shard(shard as u32, false)?; + let got = conn + .query_row( + "SELECT width, height, bytes FROM thumbs WHERE file_id = ?1", + [file_id as i64], + |r| { + Ok(Thumbnail { + width: r.get::<_, i64>(0)? as u32, + height: r.get::<_, i64>(1)? as u32, + bytes: r.get(2)?, + }) + }, + ) + .ok(); + + Ok(got) + } + + /// Whether a thumbnail is already stored. + /// + /// Cheaper than [`get`](Self::get) — the index alone answers it, with no + /// shard opened and no blob read. + pub fn contains(&self, file_id: u64) -> bool { + self.index + .query_row( + "SELECT 1 FROM entries WHERE file_id = ?1", + [file_id as i64], + |_| Ok(()), + ) + .is_ok() + } + + /// Which of these file ids are missing, preserving order. + /// + /// The grid's real question — "what must I fetch for these cells" — asked + /// in one pass rather than one query per cell. + pub fn missing(&self, file_ids: &[u64]) -> Vec { + file_ids + .iter() + .copied() + .filter(|id| !self.contains(*id)) + .collect() + } + + /// Store a thumbnail, opening a new shard if the active one is full. + /// + /// Re-storing an existing id overwrites in place rather than migrating it + /// to the active shard: a sealed shard must stay byte-identical for other + /// clients, and rewriting one would force them all to re-sync it. + pub fn put(&mut self, file_id: u64, thumb: &Thumbnail) -> Result { + if let Some(existing) = self.shard_of(file_id) { + let conn = self.open_shard(existing, true)?; + conn.execute( + "UPDATE thumbs SET width = ?2, height = ?3, bytes = ?4 WHERE file_id = ?1", + rusqlite::params![ + file_id as i64, + thumb.width as i64, + thumb.height as i64, + thumb.bytes + ], + )?; + return Ok(existing); + } + + let shard = self.active_shard(thumb.bytes.len() as u64)?; + let conn = self.open_shard(shard, true)?; + conn.execute( + "INSERT INTO thumbs(file_id, width, height, bytes) VALUES (?1, ?2, ?3, ?4) + ON CONFLICT(file_id) DO UPDATE SET + width = excluded.width, height = excluded.height, bytes = excluded.bytes", + rusqlite::params![ + file_id as i64, + thumb.width as i64, + thumb.height as i64, + thumb.bytes + ], + )?; + + self.index.execute( + "INSERT INTO entries(file_id, shard, bytes) VALUES (?1, ?2, ?3) + ON CONFLICT(file_id) DO UPDATE SET shard = excluded.shard, bytes = excluded.bytes", + rusqlite::params![file_id as i64, shard as i64, thumb.bytes.len() as i64], + )?; + self.index.execute( + "INSERT INTO shards(id, bytes, sealed) VALUES (?1, ?2, 0) + ON CONFLICT(id) DO UPDATE SET bytes = shards.bytes + ?2", + rusqlite::params![shard as i64, thumb.bytes.len() as i64], + )?; + + Ok(shard) + } + + fn shard_of(&self, file_id: u64) -> Option { + self.index + .query_row( + "SELECT shard FROM entries WHERE file_id = ?1", + [file_id as i64], + |r| r.get::<_, i64>(0), + ) + .ok() + .map(|v| v as u32) + } + + /// The shard that should receive `incoming` bytes. + /// + /// Seals the current shard and starts a new one when it would exceed the + /// cap. Sealing is recorded rather than inferred from file size, so a + /// shard stays sealed even if a later SQLite vacuum shrinks it. + fn active_shard(&self, incoming: u64) -> Result { + let current: Option<(i64, i64)> = self + .index + .query_row( + "SELECT id, bytes FROM shards WHERE sealed = 0 ORDER BY id DESC LIMIT 1", + [], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .ok(); + + match current { + Some((id, bytes)) if (bytes as u64) + incoming <= SHARD_MAX_BYTES => Ok(id as u32), + Some((id, _)) => { + self.index + .execute("UPDATE shards SET sealed = 1 WHERE id = ?1", [id])?; + let next = (id + 1) as u32; + self.index.execute( + "INSERT INTO shards(id, bytes, sealed) VALUES (?1, 0, 0) + ON CONFLICT(id) DO NOTHING", + [next as i64], + )?; + Ok(next) + } + None => { + self.index.execute( + "INSERT INTO shards(id, bytes, sealed) VALUES (0, 0, 0) + ON CONFLICT(id) DO NOTHING", + [], + )?; + Ok(0) + } + } + } + + fn open_shard(&self, shard: u32, create: bool) -> Result { + let path = self.shard_path(shard); + if !create && !path.exists() { + return Err(ThumbError::MissingShard(shard)); + } + let conn = Connection::open(&path)?; + conn.pragma_update(None, "journal_mode", "WAL")?; + conn.pragma_update(None, "synchronous", "NORMAL")?; + if create { + conn.execute_batch(SHARD_SCHEMA)?; + } + Ok(conn) + } + + /// TRACES: FR-CAT-15 | NFR-RES-4 + /// Drop thumbnails for images that no longer exist. + /// + /// The permanent-delete half of the trash: an image purged from the library + /// must not leave a preview behind. The shards **sync**, so a stale entry is + /// not merely local clutter — every other client keeps showing a thumbnail + /// for a photograph that is gone, and keeps downloading it. + /// + /// # What this deliberately does not do + /// + /// It does not rewrite the shard's blob, and it does not `VACUUM`. A sealed + /// shard must stay **byte-identical** or every client re-downloads all 25 MB + /// of it to reclaim 20 KB — the store's whole sync design (see the module + /// preamble) rests on sealed shards never changing. So the blob is deleted + /// from the *unsealed* shard only, and in a sealed one the row is left in + /// place and merely unlinked from the index. + /// + /// The consequence, accepted: a sealed shard can carry dead weight. It is + /// bounded — 25 MB per shard, only for images actually purged — and it is + /// invisible, because [`get`](Self::get) resolves through the index and an + /// unindexed blob is unreachable. Reclaiming it belongs to a compaction pass + /// that rewrites a shard under a new id, which is a separate concern with its + /// own sync cost. + /// + /// Returns how many entries were forgotten. + pub fn forget(&mut self, file_ids: &[u64]) -> Result { + if file_ids.is_empty() { + return Ok(0); + } + + // Which shard each id lives in, and whether that shard is still open for + // writing. Read before deleting anything, so the index is consulted once. + let mut in_unsealed: Vec<(u64, u32)> = Vec::new(); + let mut forgotten = 0usize; + + let tx = self.index.unchecked_transaction()?; + for &file_id in file_ids { + let found: Option<(i64, i64, i64)> = tx + .query_row( + "SELECT e.shard, e.bytes, s.sealed + FROM entries e JOIN shards s ON s.id = e.shard + WHERE e.file_id = ?1", + [file_id as i64], + |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), + ) + .ok(); + + let Some((shard, bytes, sealed)) = found else { + // Never stored, or already forgotten. Not an error: a purge runs + // over whatever the catalog knew about, and a thumbnail that was + // never fetched is the normal case. + continue; + }; + + tx.execute("DELETE FROM entries WHERE file_id = ?1", [file_id as i64])?; + forgotten += 1; + + if sealed == 0 { + // The active shard is not yet anyone's cached copy, so its bytes + // can genuinely be reclaimed and the accounting corrected — which + // also means the shard does not seal prematurely on space that is + // no longer used. + in_unsealed.push((file_id, shard as u32)); + tx.execute( + "UPDATE shards SET bytes = max(0, bytes - ?2) WHERE id = ?1", + rusqlite::params![shard, bytes], + )?; + } + } + tx.commit()?; + + // Blobs after the index commits. The other order would leave the index + // pointing at a blob that is gone if this failed partway — a `get` would + // then find an entry and no thumbnail, which reads as corruption. This + // way a failure here leaves an unreachable blob, which is the same + // harmless state a sealed shard is in by design. + for (file_id, shard) in in_unsealed { + match self.open_shard(shard, false) { + Ok(conn) => { + if let Err(e) = conn.execute( + "DELETE FROM thumbs WHERE file_id = ?1", + [file_id as i64], + ) { + log::debug!("forgetting thumbnail {file_id} in shard {shard}: {e}"); + } + } + Err(e) => log::debug!("opening shard {shard} to forget {file_id}: {e}"), + } + } + + Ok(forgotten) + } + + pub fn shard_path(&self, shard: u32) -> PathBuf { + self.dir.join(format!("shard-{shard:04}.sqlite")) + } + + pub fn index_path(&self) -> PathBuf { + self.dir.join("index.sqlite") + } + + /// Every shard, with whether it is sealed. + /// + /// The sync layer's input: sealed shards never change again, so once a + /// client has one it need never ask for it. Only the unsealed shard and + /// the index are worth re-checking. + pub fn shards(&self) -> Result, ThumbError> { + let mut stmt = self + .index + .prepare("SELECT id, bytes, sealed FROM shards ORDER BY id")?; + let rows = stmt + .query_map([], |r| { + Ok(ShardInfo { + id: r.get::<_, i64>(0)? as u32, + bytes: r.get::<_, i64>(1)? as u64, + sealed: r.get::<_, i64>(2)? != 0, + }) + })? + .collect::, _>>()?; + Ok(rows) + } + + /// Total thumbnails stored. + pub fn len(&self) -> usize { + self.index + .query_row("SELECT count(*) FROM entries", [], |r| r.get::<_, i64>(0)) + .unwrap_or(0) as usize + } + + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + /// Merge a shard downloaded from another client. + /// + /// Insert-only: a thumbnail we already hold is kept. They are derived from + /// the same bytes by the same code, so neither copy is better, and + /// preferring ours avoids rewriting a shard other clients have already + /// synced. + /// + /// Returns how many were adopted. + pub fn merge_shard(&mut self, downloaded: &Path) -> Result { + let src = Connection::open_with_flags( + downloaded, + rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NO_MUTEX, + )?; + + let mut stmt = src.prepare("SELECT file_id, width, height, bytes FROM thumbs")?; + let incoming = stmt + .query_map([], |r| { + Ok(( + r.get::<_, i64>(0)? as u64, + Thumbnail { + width: r.get::<_, i64>(1)? as u32, + height: r.get::<_, i64>(2)? as u32, + bytes: r.get(3)?, + }, + )) + })? + .collect::, _>>()?; + + let mut adopted = 0; + for (file_id, thumb) in incoming { + if self.contains(file_id) { + continue; + } + self.put(file_id, &thumb)?; + adopted += 1; + } + Ok(adopted) + } +} + +/// One shard's sync-relevant state. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ShardInfo { + pub id: u32, + pub bytes: u64, + /// A sealed shard is immutable. Once downloaded it never needs re-checking. + pub sealed: bool, +} + +const INDEX_SCHEMA: &str = r#" +CREATE TABLE IF NOT EXISTS entries ( + file_id INTEGER PRIMARY KEY, -- oc:fileid, stable across rename/move + shard INTEGER NOT NULL, + bytes INTEGER NOT NULL +); +CREATE INDEX IF NOT EXISTS entries_shard ON entries(shard); + +CREATE TABLE IF NOT EXISTS shards ( + id INTEGER PRIMARY KEY, + bytes INTEGER NOT NULL DEFAULT 0, + -- Recorded, not inferred from file size: a vacuum could shrink a sealed + -- shard below the cap and it must still stay closed, or clients that + -- already hold it would see it change. + sealed INTEGER NOT NULL DEFAULT 0 +); +"#; + +const SHARD_SCHEMA: &str = r#" +CREATE TABLE IF NOT EXISTS thumbs ( + file_id INTEGER PRIMARY KEY, + width INTEGER NOT NULL, + height INTEGER NOT NULL, + bytes BLOB NOT NULL +); +"#; + +#[cfg(test)] +mod tests { + use super::*; + + fn store() -> (ThumbStore, PathBuf) { + let dir = std::env::temp_dir().join(format!( + "dr-thumbs-{}-{:?}", + std::process::id(), + std::thread::current().id() + )); + let _ = std::fs::remove_dir_all(&dir); + (ThumbStore::open(&dir).unwrap(), dir) + } + + fn thumb(size: usize) -> Thumbnail { + Thumbnail { + width: 256, + height: 170, + bytes: vec![0xAB; size], + } + } + + #[test] + fn a_forgotten_thumbnail_is_no_longer_served() { + // The point of the whole method: a purged photograph must not keep a + // preview, on this device or any that syncs the shards. + let (mut s, _d) = store(); + s.put(1, &thumb(1024)).unwrap(); + s.put(2, &thumb(1024)).unwrap(); + + assert_eq!(s.forget(&[1]).unwrap(), 1); + assert!(s.get(1).unwrap().is_none()); + assert!(!s.contains(1)); + // And its neighbour is untouched. + assert!(s.get(2).unwrap().is_some()); + assert_eq!(s.len(), 1); + } + + #[test] + fn forgetting_reclaims_bytes_in_the_open_shard() { + // Otherwise the shard seals on space nothing is using, and the store + // grows a shard per purge. + let (mut s, _d) = store(); + s.put(1, &thumb(4096)).unwrap(); + let before = s.shards().unwrap()[0].bytes; + + s.forget(&[1]).unwrap(); + let after = s.shards().unwrap()[0].bytes; + assert!(after < before, "{after} should be below {before}"); + } + + #[test] + fn forgetting_from_a_sealed_shard_leaves_it_byte_identical() { + // The invariant the store's sync design rests on: a sealed shard never + // changes, or every client re-downloads 25 MB to reclaim 20 KB. + let (mut s, _d) = store(); + let big = (SHARD_MAX_BYTES / 4) as usize; + for id in 0..5 { + s.put(id, &thumb(big)).unwrap(); + } + let shards = s.shards().unwrap(); + assert!(shards[0].sealed, "precondition: shard 0 is sealed"); + + let sealed_path = s.shard_path(0); + let before = std::fs::metadata(&sealed_path).unwrap().len(); + let before_bytes = shards[0].bytes; + + // Image 0 is in the sealed shard. + assert_eq!(s.forget(&[0]).unwrap(), 1); + + // Unreachable through the index, which is what matters to a reader... + assert!(s.get(0).unwrap().is_none()); + // ...but the file itself is untouched. + assert_eq!( + std::fs::metadata(&sealed_path).unwrap().len(), + before, + "a sealed shard must not be rewritten" + ); + // And its accounting is left alone, so it cannot un-seal. + let after = s.shards().unwrap(); + assert_eq!(after[0].bytes, before_bytes); + assert!(after[0].sealed); + } + + #[test] + fn forgetting_something_never_stored_is_not_an_error() { + // A purge runs over whatever the catalog knew; most images never had a + // thumbnail fetched. + let (mut s, _d) = store(); + s.put(1, &thumb(512)).unwrap(); + assert_eq!(s.forget(&[42, 43]).unwrap(), 0); + assert_eq!(s.forget(&[]).unwrap(), 0); + assert_eq!(s.len(), 1, "nothing else went"); + } + + #[test] + fn forgetting_twice_is_idempotent() { + // An empty-trash retried after a partial failure runs over the same ids. + let (mut s, _d) = store(); + s.put(1, &thumb(512)).unwrap(); + assert_eq!(s.forget(&[1]).unwrap(), 1); + assert_eq!(s.forget(&[1]).unwrap(), 0); + } + + #[test] + fn a_forgotten_id_can_be_stored_again() { + // Restoring from the server's own trashbin, or re-adding the same file: + // the id is stable, so the store must accept it back. + let (mut s, _d) = store(); + s.put(1, &thumb(512)).unwrap(); + s.forget(&[1]).unwrap(); + s.put(1, &thumb(512)).unwrap(); + assert!(s.get(1).unwrap().is_some()); + } + + #[test] + fn a_stored_thumbnail_round_trips() { + let (mut s, _d) = store(); + s.put(1001, &thumb(1024)).unwrap(); + + let got = s.get(1001).unwrap().expect("stored"); + assert_eq!(got.width, 256); + assert_eq!(got.bytes.len(), 1024); + } + + #[test] + fn an_unknown_id_is_none_not_an_error() { + let (s, _d) = store(); + assert!(s.get(9999).unwrap().is_none()); + } + + #[test] + fn contains_avoids_reading_the_blob() { + let (mut s, _d) = store(); + s.put(1, &thumb(512)).unwrap(); + assert!(s.contains(1)); + assert!(!s.contains(2)); + } + + #[test] + fn missing_reports_only_what_is_absent() { + let (mut s, _d) = store(); + s.put(1, &thumb(64)).unwrap(); + s.put(3, &thumb(64)).unwrap(); + assert_eq!(s.missing(&[1, 2, 3, 4]), vec![2, 4]); + } + + #[test] + fn everything_small_lands_in_one_shard() { + let (mut s, _d) = store(); + for id in 0..20 { + s.put(id, &thumb(1024)).unwrap(); + } + assert_eq!(s.shards().unwrap().len(), 1); + assert_eq!(s.len(), 20); + } + + #[test] + fn a_full_shard_seals_and_the_next_one_opens() { + let (mut s, _d) = store(); + // Quarter-cap blobs: the fifth cannot fit alongside four others. + let big = (SHARD_MAX_BYTES / 4) as usize; + for id in 0..5 { + s.put(id, &thumb(big)).unwrap(); + } + + let shards = s.shards().unwrap(); + assert_eq!(shards.len(), 2, "should have rolled over"); + assert!(shards[0].sealed, "the full shard is sealed"); + assert!(!shards[1].sealed, "the new one is open"); + } + + #[test] + fn a_sealed_shard_never_reopens() { + // The property sync depends on: once a client has a sealed shard it + // never needs to ask again. Reopening one would force every client to + // re-download it. + let (mut s, _d) = store(); + let big = (SHARD_MAX_BYTES / 4) as usize; + for id in 0..5 { + s.put(id, &thumb(big)).unwrap(); + } + // Small enough to fit in the sealed shard, but it must not go there. + s.put(100, &thumb(16)).unwrap(); + + let shards = s.shards().unwrap(); + assert!(shards[0].sealed); + assert_eq!( + s.shard_of(100), + Some(1), + "new writes go to the open shard, not back into a sealed one" + ); + } + + #[test] + fn rewriting_updates_in_place_rather_than_migrating() { + // Moving an entry to the active shard would rewrite a sealed shard, + // which every other client has already synced. + let (mut s, _d) = store(); + let big = (SHARD_MAX_BYTES / 4) as usize; + for id in 0..5 { + s.put(id, &thumb(big)).unwrap(); + } + let original = s.shard_of(0).unwrap(); + + s.put(0, &thumb(32)).unwrap(); + assert_eq!(s.shard_of(0), Some(original), "must stay put"); + assert_eq!(s.get(0).unwrap().unwrap().bytes.len(), 32, "but update"); + } + + #[test] + fn shard_ids_are_zero_padded_for_stable_ordering() { + let (s, _d) = store(); + let p = s.shard_path(7); + assert!(p.to_string_lossy().ends_with("shard-0007.sqlite")); + } + + #[test] + fn reopening_finds_what_was_stored() { + let (mut s, dir) = store(); + s.put(42, &thumb(256)).unwrap(); + drop(s); + + let reopened = ThumbStore::open(&dir).unwrap(); + assert!(reopened.contains(42)); + assert_eq!(reopened.len(), 1); + } + + #[test] + fn merging_another_clients_shard_adopts_only_what_is_new() { + let (mut mine, _d1) = store(); + mine.put(1, &thumb(64)).unwrap(); + + // A second store standing in for another device's downloaded shard. + let other_dir = std::env::temp_dir().join(format!("dr-thumbs-other-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&other_dir); + let mut theirs = ThumbStore::open(&other_dir).unwrap(); + theirs.put(1, &thumb(999)).unwrap(); // we already have this one + theirs.put(2, &thumb(64)).unwrap(); + theirs.put(3, &thumb(64)).unwrap(); + + let adopted = mine.merge_shard(&theirs.shard_path(0)).unwrap(); + assert_eq!(adopted, 2, "only the two we lacked"); + assert_eq!( + mine.get(1).unwrap().unwrap().bytes.len(), + 64, + "ours is kept, not overwritten" + ); + assert!(mine.contains(2) && mine.contains(3)); + } + + #[test] + fn merging_is_idempotent() { + let (mut mine, _d1) = store(); + let other_dir = + std::env::temp_dir().join(format!("dr-thumbs-idem-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&other_dir); + let mut theirs = ThumbStore::open(&other_dir).unwrap(); + theirs.put(10, &thumb(64)).unwrap(); + + assert_eq!(mine.merge_shard(&theirs.shard_path(0)).unwrap(), 1); + assert_eq!( + mine.merge_shard(&theirs.shard_path(0)).unwrap(), + 0, + "a repeated merge adopts nothing" + ); + } + + #[test] + fn a_missing_shard_errors_rather_than_panicking() { + let (s, _d) = store(); + assert!(matches!( + s.open_shard(99, false), + Err(ThumbError::MissingShard(99)) + )); + } +}