Add dr-thumbs: a sharded, syncable thumbnail store
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
This commit is contained in:
@@ -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
|
||||
@@ -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<Vec<u8>, 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<u8>), 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<u8> {
|
||||
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<u8> = 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]);
|
||||
}
|
||||
}
|
||||
@@ -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),
|
||||
}
|
||||
@@ -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<u8>,
|
||||
}
|
||||
|
||||
/// 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<Self, ThumbError> {
|
||||
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<Option<Thumbnail>, ThumbError> {
|
||||
let shard: Option<i64> = 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<u64> {
|
||||
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<u32, ThumbError> {
|
||||
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<u32> {
|
||||
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<u32, ThumbError> {
|
||||
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<Connection, ThumbError> {
|
||||
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<usize, ThumbError> {
|
||||
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<Vec<ShardInfo>, 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::<Result<Vec<_>, _>>()?;
|
||||
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<usize, ThumbError> {
|
||||
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::<Result<Vec<_>, _>>()?;
|
||||
|
||||
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))
|
||||
));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user