From 96d07da15f618ff5575e44bdc6b99e9bd42643aa Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Wed, 26 Aug 2026 22:41:50 +0200 Subject: [PATCH] Sync face data as sealed shards, so a second device does not re-index Indexing 23,500 images is about two hours of CPU, and the result is byte-identical on every device: the same model over the same proxy produces the same embedding. Paying for it once per account rather than once per device is the point. Shards rather than the catalog snapshot, because the snapshot goes up whole on every sync and a fully indexed library carries roughly 30 MB of embeddings. That is exactly the cost the thumbnail store's 25 MB cap exists to bound, so face shards use the same cap -- imported from dr_thumbs rather than restated, since the number is a statement about sync cost and the two must not drift apart. The split follows the one already there: bulk immutable data in sealed shards, small mutable data in the catalog snapshot. Faces, landmarks, embeddings and run markers shard; people, names and assignments ride the catalog and merge by uuid. Keyed on oc:fileid throughout, never on image_id, because a row id means nothing on another device. The run marker travels with the faces it describes. Without it a receiving device cannot tell an image with no faces from one never examined, and would re-detect every landscape it had just adopted. Co-Authored-By: Claude Opus 5 (1M context) --- Cargo.lock | 1 + core/dr-catalog/Cargo.toml | 4 + core/dr-catalog/src/face_shard.rs | 958 ++++++++++++++++++++++++++++++ core/dr-catalog/src/lib.rs | 2 + docs/traceability.md | 10 +- ui/dr-ui/examples/face_index.rs | 47 +- ui/dr-ui/src/derived_sync.rs | 163 +++++ 7 files changed, 1179 insertions(+), 6 deletions(-) create mode 100644 core/dr-catalog/src/face_shard.rs diff --git a/Cargo.lock b/Cargo.lock index 3ea2a48..1617f45 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1408,6 +1408,7 @@ version = "0.7.0" dependencies = [ "dr-face", "dr-plat", + "dr-thumbs", "dr-types", "env_logger", "log", diff --git a/core/dr-catalog/Cargo.toml b/core/dr-catalog/Cargo.toml index 984aa5c..651453c 100644 --- a/core/dr-catalog/Cargo.toml +++ b/core/dr-catalog/Cargo.toml @@ -12,6 +12,10 @@ dr-types.workspace = true # features are off, so this brings in no ONNX runtime and no weights: only the # model-free half compiles here. dr-face.workspace = true +# For `SHARD_MAX_BYTES` alone. The face shards are capped at the same 25 MB the +# thumbnail shards are, and sharing the constant is what keeps them from +# drifting apart — the cap is a statement about sync cost, not about thumbnails. +dr-thumbs.workspace = true # The `Storage` trait, and nothing else from it. A scan has to read a real # directory, and this is how `core/` reaches the platform without a # `#[cfg(target_os)]` of its own (ARCH §4.1: calls go downward). diff --git a/core/dr-catalog/src/face_shard.rs b/core/dr-catalog/src/face_shard.rs new file mode 100644 index 0000000..bf413be --- /dev/null +++ b/core/dr-catalog/src/face_shard.rs @@ -0,0 +1,958 @@ +//! TRACES: FR-CULL-8 | FR-NC-7 | NFR-RES-4 +//! Face data as sealed shards, so a second device does not re-index the library. +//! +//! Indexing a 23,500-image library is on the order of two hours of CPU +//! (docs/faces.md §12.2). It is also **byte-identical on every device**: the +//! same model over the same proxy produces the same embedding. Paying for it +//! once per account rather than once per device is the whole point of this +//! module, and it is the same bargain the thumbnail store already makes. +//! +//! # Why shards, and not the catalog snapshot +//! +//! The catalog is uploaded whole on every sync. Embeddings are 1 KB each, so a +//! fully indexed 23,500-image library carries roughly 30 MB of them — which +//! would be re-uploaded in its entirety every time anything in the catalog +//! changed, and re-downloaded by every device. That is precisely the cost +//! `dr_thumbs`'s 25 MB shard cap exists to bound. +//! +//! So face data splits the way thumbnails and collections already split: +//! +//! - **Here, in sealed shards:** the faces themselves — box, landmarks, +//! embedding, and the run marker. Bulk, immutable once written, and identical +//! everywhere. A sealed shard is never rewritten, so a client downloads each +//! one exactly once and never asks about it again. +//! - **In the catalog snapshot, which already syncs:** people, their names, and +//! which face belongs to whom. Small, changes constantly, and merges by uuid. +//! +//! # Keyed on `file_id`, never on `image_id` +//! +//! A row id is local and means nothing on another device. `oc:fileid` is the +//! identity every client agrees on (FR-NC-5), and it survives a server-side +//! move, so a shard written before a reorganisation still applies after it. + +use std::path::{Path, PathBuf}; + +use rusqlite::{Connection, OptionalExtension}; + +use crate::error::CatalogError; + +/// Maximum bytes in a shard before a new one is started. +/// +/// The same 25 MB `dr_thumbs` uses, for the same reason and deliberately not a +/// separate number: it bounds what a client re-downloads when the open shard +/// changes, and keeps a stalled transfer cheap to retry. At ~1.1 KB per face +/// that is roughly 22,000 faces per shard, so a large personal library lands in +/// one or two. +pub const SHARD_MAX_BYTES: u64 = dr_thumbs::SHARD_MAX_BYTES; + +/// Bytes one stored face occupies, near enough to bound a shard by. +/// +/// Counted rather than measured: the embedding is fixed at 512 × f16, the +/// landmarks at 5 × 2 × f32, and the rest is a handful of numbers. Measuring +/// the file after each insert would mean a `VACUUM` to get an honest answer. +const BYTES_PER_FACE: u64 = 1024 + 40 + 64; + +/// One face as it travels between devices. +#[derive(Debug, Clone, PartialEq)] +pub struct SharedFace { + /// `oc:fileid` of the photograph — the cross-device image identity. + pub file_id: u64, + pub model_id: String, + pub x: f32, + pub y: f32, + pub w: f32, + pub h: f32, + pub landmarks: Vec, + pub confidence: f32, + pub embedding: Vec, + pub crop_px: f32, +} + +/// A store of face shards, beside the thumbnail store. +pub struct FaceShardStore { + dir: PathBuf, + index: Connection, + client: String, +} + +impl FaceShardStore { + pub fn open(dir: &Path) -> Result { + std::fs::create_dir_all(dir).map_err(|e| CatalogError::Io(e.to_string()))?; + let index = Connection::open(dir.join("index.sqlite"))?; + index.execute_batch(INDEX_SCHEMA)?; + let client = mint_client_id(&index)?; + Ok(Self { + dir: dir.to_path_buf(), + index, + client, + }) + } + + /// This device's id, which qualifies its shard numbering. + /// + /// Two devices both writing `shard-0000` would collide on the server; the + /// name carries the writer so they cannot. + pub fn client_id(&self) -> &str { + &self.client + } + + pub fn contains(&self, file_id: u64, model_id: &str) -> bool { + self.index + .query_row( + "SELECT 1 FROM entries WHERE file_id = ?1 AND model_id = ?2", + rusqlite::params![file_id as i64, model_id], + |_| Ok(()), + ) + .optional() + .ok() + .flatten() + .is_some() + } + + /// Number of faces held. + pub fn len(&self) -> u64 { + self.index + .query_row("SELECT COUNT(*) FROM faces_meta", [], |r| { + r.get::<_, i64>(0) + }) + .map(|n| n as u64) + .unwrap_or(0) + } + + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + /// Write a batch of faces for one image into the open shard. + /// + /// Per image rather than per face, because the run marker is per image and + /// the two must land together: a shard holding faces but no marker would + /// make the receiving device re-detect the image it just adopted. + pub fn put_image( + &mut self, + file_id: u64, + model_id: &str, + source_edge: u32, + faces: &[SharedFace], + ) -> Result { + let incoming = faces.len() as u64 * BYTES_PER_FACE + 64; + let shard = self.active_shard(incoming)?; + let conn = self.open_shard(shard, true)?; + + let tx = conn.unchecked_transaction()?; + // Replace rather than append: re-indexing an image must not double its + // faces in the shard, exactly as it must not in the catalog. + tx.execute( + "DELETE FROM faces WHERE file_id = ?1 AND model_id = ?2", + rusqlite::params![file_id as i64, model_id], + )?; + for f in faces { + tx.execute( + "INSERT INTO faces + (file_id, model_id, x, y, w, h, landmarks, confidence, + embedding, crop_px) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", + rusqlite::params![ + f.file_id as i64, + f.model_id, + f.x as f64, + f.y as f64, + f.w as f64, + f.h as f64, + f.landmarks, + f.confidence as f64, + f.embedding, + f.crop_px as f64, + ], + )?; + } + tx.execute( + "INSERT INTO indexed (file_id, model_id, faces_found, source_edge) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT(file_id, model_id) DO UPDATE SET + faces_found = excluded.faces_found, + source_edge = excluded.source_edge", + rusqlite::params![ + file_id as i64, + model_id, + faces.len() as i64, + source_edge as i64 + ], + )?; + tx.commit()?; + + self.index.execute( + "INSERT INTO entries (file_id, model_id, shard, bytes) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT(file_id, model_id) DO UPDATE SET + shard = excluded.shard, bytes = excluded.bytes", + rusqlite::params![file_id as i64, model_id, shard as i64, incoming as i64], + )?; + self.index.execute( + "INSERT INTO faces_meta (file_id, model_id) VALUES (?1, ?2) + ON CONFLICT(file_id, model_id) DO NOTHING", + rusqlite::params![file_id as i64, model_id], + )?; + self.index.execute( + "UPDATE shards SET bytes = bytes + ?2 WHERE id = ?1", + rusqlite::params![shard as i64, incoming as i64], + )?; + Ok(shard) + } + + /// The shard currently open for writing, sealing and starting a new one + /// where this batch would overflow the cap. + fn active_shard(&self, incoming: u64) -> Result { + let open: 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)?)), + ) + .optional()?; + + match open { + Some((id, bytes)) if (bytes as u64) + incoming <= SHARD_MAX_BYTES => Ok(id as u32), + Some((id, _)) => { + // Sealed, and never written to again. That immutability is what + // lets a peer download it once and stop asking. + 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)", + [next as i64], + )?; + Ok(next) + } + None => { + self.index.execute( + "INSERT INTO shards (id, bytes, sealed) VALUES (0, 0, 0)", + [], + )?; + Ok(0) + } + } + } + + fn open_shard(&self, shard: u32, create: bool) -> Result { + let path = self.shard_path(shard); + if !create && !path.is_file() { + return Err(CatalogError::Io(format!("face shard {shard} is missing"))); + } + let conn = Connection::open(&path)?; + conn.execute_batch(SHARD_SCHEMA)?; + Ok(conn) + } + + pub fn shard_path(&self, shard: u32) -> PathBuf { + self.dir.join(format!("shard-{}-{shard:04}.sqlite", self.client)) + } + + /// Every shard, with whether it is sealed. + pub fn shards(&self) -> Result, CatalogError> { + let mut q = self + .index + .prepare("SELECT id, bytes, sealed FROM shards ORDER BY id")?; + let rows = q.query_map([], |r| { + Ok(ShardInfo { + id: r.get::<_, i64>(0)? as u32, + bytes: r.get::<_, i64>(1)? as u64, + sealed: r.get(2)?, + }) + })?; + rows.collect::>().map_err(Into::into) + } + + /// Whether a peer's shard has already been taken in. + /// + /// Nothing in an adopted face records where it came from, so without this + /// ledger a client re-downloads every peer's shard on every sync. + pub fn has_adopted(&self, name: &str, size: u64) -> bool { + self.index + .query_row( + "SELECT 1 FROM adopted WHERE name = ?1 AND size = ?2", + rusqlite::params![name, size as i64], + |_| Ok(()), + ) + .optional() + .ok() + .flatten() + .is_some() + } + + pub fn mark_adopted(&self, name: &str, size: u64) -> Result<(), CatalogError> { + self.index.execute( + "INSERT INTO adopted (name, size) VALUES (?1, ?2) + ON CONFLICT(name) DO UPDATE SET size = excluded.size", + rusqlite::params![name, size as i64], + )?; + Ok(()) + } + + /// Take in a peer's shard, skipping images this store already holds. + /// + /// Returns how many images 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 q = src.prepare("SELECT file_id, model_id, faces_found, source_edge FROM indexed")?; + let images: Vec<(i64, String, i64, i64)> = q + .query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)))? + .collect::>()?; + + let mut adopted = 0; + for (file_id, model_id, _found, edge) in images { + if self.contains(file_id as u64, &model_id) { + continue; + } + let mut fq = src.prepare( + "SELECT file_id, model_id, x, y, w, h, landmarks, confidence, + embedding, crop_px + FROM faces WHERE file_id = ?1 AND model_id = ?2", + )?; + let faces: Vec = fq + .query_map(rusqlite::params![file_id, &model_id], read_shared_face)? + .collect::>()?; + self.put_image(file_id as u64, &model_id, edge as u32, &faces)?; + adopted += 1; + } + Ok(adopted) + } + + /// Read back everything held for one image. + pub fn get_image( + &self, + file_id: u64, + model_id: &str, + ) -> Result, u32)>, CatalogError> { + let Some(shard): Option = self + .index + .query_row( + "SELECT shard FROM entries WHERE file_id = ?1 AND model_id = ?2", + rusqlite::params![file_id as i64, model_id], + |r| r.get(0), + ) + .optional()? + else { + return Ok(None); + }; + + let conn = self.open_shard(shard as u32, false)?; + let edge: Option = conn + .query_row( + "SELECT source_edge FROM indexed WHERE file_id = ?1 AND model_id = ?2", + rusqlite::params![file_id as i64, model_id], + |r| r.get(0), + ) + .optional()?; + let Some(edge) = edge else { return Ok(None) }; + + let mut q = conn.prepare( + "SELECT file_id, model_id, x, y, w, h, landmarks, confidence, embedding, crop_px + FROM faces WHERE file_id = ?1 AND model_id = ?2", + )?; + let faces: Vec = q + .query_map(rusqlite::params![file_id as i64, model_id], read_shared_face)? + .collect::>()?; + Ok(Some((faces, edge as u32))) + } +} + +/// Copy this device's indexed faces into the shard store, ready to upload. +/// +/// Only what the store does not already hold, so a call after a partial sweep +/// writes only the new images. Returns how many images were exported. +/// +/// Images with no `oc:fileid` are skipped rather than failing the pass: a +/// purely local file has no identity another device could match it by, so there +/// is nothing useful to send. +pub fn export_to_shards( + conn: &Connection, + store: &mut FaceShardStore, + model_id: &str, +) -> Result { + let mut q = conn.prepare( + "SELECT r.file_id, fi.image_id, fi.source_edge + FROM face_index fi + JOIN remote r ON r.image_id = fi.image_id + WHERE fi.model_id = ?1 + ORDER BY fi.image_id", + )?; + let rows: Vec<(i64, i64, i64)> = q + .query_map([model_id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))? + .collect::>()?; + + let mut exported = 0; + for (file_id, image_id, edge) in rows { + if store.contains(file_id as u64, model_id) { + continue; + } + let mut fq = conn.prepare( + "SELECT x, y, w, h, landmarks, detector_confidence, embedding, crop_px + FROM faces WHERE image_id = ?1 AND model_id = ?2", + )?; + let faces: Vec = fq + .query_map(rusqlite::params![image_id, model_id], |r| { + Ok(SharedFace { + file_id: file_id as u64, + model_id: model_id.to_string(), + x: r.get::<_, f64>(0)? as f32, + y: r.get::<_, f64>(1)? as f32, + w: r.get::<_, f64>(2)? as f32, + h: r.get::<_, f64>(3)? as f32, + landmarks: r.get(4)?, + confidence: r.get::<_, f64>(5)? as f32, + embedding: r.get(6)?, + crop_px: r.get::<_, f64>(7)? as f32, + }) + })? + .collect::>()?; + + store.put_image(file_id as u64, model_id, edge as u32, &faces)?; + exported += 1; + } + Ok(exported) +} + +/// Take faces out of the shard store and into this device's catalog. +/// +/// **The point of the whole module.** An image a peer already examined is +/// adopted rather than re-detected, which is the difference between a new +/// device being useful in a minute and in two hours. +/// +/// Skips any image this device has already indexed itself. Local work is not +/// second-guessed by a peer's — the two should agree, since the same model over +/// the same proxy is deterministic, but where they do not, the copy this device +/// computed is the one it can vouch for. +/// +/// Returns how many images were adopted. +pub fn import_from_shards( + conn: &Connection, + store: &FaceShardStore, + model_id: &str, +) -> Result { + // Only images this device actually has. A shard covers the whole account, + // and a device holding a subset of the library should take only its own + // part rather than accumulating faces for photographs it cannot show. + let mut q = conn.prepare( + "SELECT r.file_id, r.image_id + FROM remote r + JOIN images i ON i.id = r.image_id + WHERE i.trashed_at IS NULL + AND NOT EXISTS ( + SELECT 1 FROM face_index fi + WHERE fi.image_id = r.image_id AND fi.model_id = ?1 + )", + )?; + let candidates: Vec<(i64, i64)> = q + .query_map([model_id], |r| Ok((r.get(0)?, r.get(1)?)))? + .collect::>()?; + + let mut adopted = 0; + for (file_id, image_id) in candidates { + let Some((faces, edge)) = store.get_image(file_id as u64, model_id)? else { + continue; + }; + let local: Vec = faces + .into_iter() + .map(|f| crate::faces::DetectedFace { + x: f.x, + y: f.y, + w: f.w, + h: f.h, + landmarks: blob_to_landmarks(&f.landmarks), + confidence: f.confidence, + embedding: f.embedding, + crop_px: f.crop_px, + model_id: f.model_id, + }) + .collect(); + + crate::faces::record_detections( + conn, + dr_types::ImageId(image_id as u64), + model_id, + edge, + &local, + )?; + adopted += 1; + } + Ok(adopted) +} + +fn blob_to_landmarks(b: &[u8]) -> [(f32, f32); 5] { + let mut out = [(0.0_f32, 0.0_f32); 5]; + for (i, o) in out.iter_mut().enumerate() { + let at = i * 8; + if at + 8 <= b.len() { + o.0 = f32::from_le_bytes([b[at], b[at + 1], b[at + 2], b[at + 3]]); + o.1 = f32::from_le_bytes([b[at + 4], b[at + 5], b[at + 6], b[at + 7]]); + } + } + out +} + +fn read_shared_face(r: &rusqlite::Row<'_>) -> rusqlite::Result { + Ok(SharedFace { + file_id: r.get::<_, i64>(0)? as u64, + model_id: r.get(1)?, + x: r.get::<_, f64>(2)? as f32, + y: r.get::<_, f64>(3)? as f32, + w: r.get::<_, f64>(4)? as f32, + h: r.get::<_, f64>(5)? as f32, + landmarks: r.get(6)?, + confidence: r.get::<_, f64>(7)? as f32, + embedding: r.get(8)?, + crop_px: r.get::<_, f64>(9)? as f32, + }) +} + +/// One shard's sync-relevant state. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ShardInfo { + pub id: u32, + pub bytes: u64, + /// Sealed shards are immutable; once downloaded they never need re-checking. + pub sealed: bool, +} + +fn mint_client_id(conn: &Connection) -> Result { + conn.execute( + "INSERT INTO meta(key, value) VALUES('client_id', lower(hex(randomblob(6)))) + ON CONFLICT(key) DO NOTHING", + [], + )?; + Ok( + conn.query_row("SELECT value FROM meta WHERE key = 'client_id'", [], |r| { + r.get(0) + })?, + ) +} + +const INDEX_SCHEMA: &str = r#" +CREATE TABLE IF NOT EXISTS entries ( + file_id INTEGER NOT NULL, + model_id TEXT NOT NULL, + shard INTEGER NOT NULL, + bytes INTEGER NOT NULL, + PRIMARY KEY(file_id, model_id) +); +CREATE INDEX IF NOT EXISTS entries_shard ON entries(shard); + +-- One row per (image, model) held, so a count is a count rather than a scan +-- across every shard file. +CREATE TABLE IF NOT EXISTS faces_meta ( + file_id INTEGER NOT NULL, + model_id TEXT NOT NULL, + PRIMARY KEY(file_id, model_id) +); + +CREATE TABLE IF NOT EXISTS shards ( + id INTEGER PRIMARY KEY, + bytes INTEGER NOT NULL DEFAULT 0, + -- Recorded rather than inferred from file size: a vacuum could shrink a + -- sealed shard below the cap and it must still stay closed, or clients + -- holding it would see it change. + sealed INTEGER NOT NULL DEFAULT 0 +); + +CREATE TABLE IF NOT EXISTS meta ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL +); + +CREATE TABLE IF NOT EXISTS adopted ( + name TEXT PRIMARY KEY, + size INTEGER NOT NULL +); +"#; + +const SHARD_SCHEMA: &str = r#" +CREATE TABLE IF NOT EXISTS faces ( + file_id INTEGER NOT NULL, + model_id TEXT NOT NULL, + x REAL NOT NULL, y REAL NOT NULL, w REAL NOT NULL, h REAL NOT NULL, + landmarks BLOB NOT NULL, + confidence REAL NOT NULL, + embedding BLOB NOT NULL, + crop_px REAL NOT NULL +); +CREATE INDEX IF NOT EXISTS faces_file ON faces(file_id, model_id); + +-- The run marker, travelling with the faces it describes. Without it a +-- receiving device cannot tell an image with no faces from one never examined, +-- and would re-detect every landscape it just adopted. +CREATE TABLE IF NOT EXISTS indexed ( + file_id INTEGER NOT NULL, + model_id TEXT NOT NULL, + faces_found INTEGER NOT NULL, + source_edge INTEGER NOT NULL, + PRIMARY KEY(file_id, model_id) +); +"#; + +#[cfg(test)] +mod tests { + use super::*; + + /// A fresh directory per call. + /// + /// Per *call*, not per test: several of these tests stand up two stores to + /// play the two devices against each other, and a helper keyed on the + /// thread would have the second call delete the first store's files. + fn tempdir() -> PathBuf { + use std::sync::atomic::{AtomicU64, Ordering}; + static N: AtomicU64 = AtomicU64::new(0); + let d = std::env::temp_dir().join(format!( + "dr-face-shard-{}-{}", + std::process::id(), + N.fetch_add(1, Ordering::Relaxed) + )); + let _ = std::fs::remove_dir_all(&d); + std::fs::create_dir_all(&d).unwrap(); + d + } + + fn face(file_id: u64, seed: u8) -> SharedFace { + SharedFace { + file_id, + model_id: "w600k_mbf".into(), + x: 0.1, + y: 0.2, + w: 0.15, + h: 0.2, + landmarks: vec![seed; 40], + confidence: 0.87, + embedding: vec![seed; 1024], + crop_px: 180.0, + } + } + + #[test] + fn a_stored_image_round_trips() { + let dir = tempdir(); + let mut s = FaceShardStore::open(&dir).unwrap(); + s.put_image(42, "w600k_mbf", 1024, &[face(42, 1), face(42, 2)]) + .unwrap(); + + assert!(s.contains(42, "w600k_mbf")); + let (faces, edge) = s.get_image(42, "w600k_mbf").unwrap().unwrap(); + assert_eq!(faces.len(), 2); + assert_eq!(edge, 1024); + assert_eq!(faces[0].embedding.len(), 1024); + assert!((faces[0].crop_px - 180.0).abs() < 1e-3); + } + + /// The case the run marker exists for, carried across the wire: an image + /// with no faces must still read as *indexed* on the receiving device. + #[test] + fn an_image_with_no_faces_is_still_recorded() { + let dir = tempdir(); + let mut s = FaceShardStore::open(&dir).unwrap(); + s.put_image(7, "w600k_mbf", 1024, &[]).unwrap(); + + assert!(s.contains(7, "w600k_mbf")); + let (faces, _) = s.get_image(7, "w600k_mbf").unwrap().unwrap(); + assert!(faces.is_empty()); + } + + #[test] + fn a_different_model_is_a_different_entry() { + let dir = tempdir(); + let mut s = FaceShardStore::open(&dir).unwrap(); + s.put_image(1, "w600k_mbf", 1024, &[face(1, 1)]).unwrap(); + assert!(s.contains(1, "w600k_mbf")); + assert!(!s.contains(1, "lvface")); + } + + #[test] + fn re_storing_an_image_replaces_rather_than_doubling_it() { + let dir = tempdir(); + let mut s = FaceShardStore::open(&dir).unwrap(); + s.put_image(1, "w600k_mbf", 1024, &[face(1, 1), face(1, 2)]) + .unwrap(); + s.put_image(1, "w600k_mbf", 2048, &[face(1, 3)]).unwrap(); + + let (faces, edge) = s.get_image(1, "w600k_mbf").unwrap().unwrap(); + assert_eq!(faces.len(), 1); + assert_eq!(edge, 2048); + } + + /// The shard cap is the whole reason this module exists rather than the + /// data riding in the catalog snapshot. + #[test] + fn shards_seal_at_the_cap_and_a_new_one_opens() { + let dir = tempdir(); + let mut s = FaceShardStore::open(&dir).unwrap(); + + // Enough faces to overflow one shard several times over. + let per_image = 64; + let images = (SHARD_MAX_BYTES / (per_image as u64 * BYTES_PER_FACE)) + 2; + for i in 0..images { + let faces: Vec = (0..per_image).map(|k| face(i, k as u8)).collect(); + s.put_image(i, "w600k_mbf", 1024, &faces).unwrap(); + } + + let shards = s.shards().unwrap(); + assert!(shards.len() >= 2, "expected a seal, got {} shard(s)", shards.len()); + assert!(shards[0].sealed, "the first shard should be sealed"); + assert!( + shards[0].bytes <= SHARD_MAX_BYTES, + "sealed shard is {} bytes, over the {SHARD_MAX_BYTES} cap", + shards[0].bytes + ); + assert!(!shards.last().unwrap().sealed, "the last shard stays open"); + } + + #[test] + fn merging_a_peers_shard_adopts_what_this_device_lacks() { + let peer_dir = tempdir(); + let mut peer = FaceShardStore::open(&peer_dir).unwrap(); + peer.put_image(1, "w600k_mbf", 1024, &[face(1, 1)]).unwrap(); + peer.put_image(2, "w600k_mbf", 1024, &[]).unwrap(); + peer.put_image(3, "w600k_mbf", 1024, &[face(3, 3)]).unwrap(); + let peer_shard = peer.shard_path(0); + + let mine_dir = tempdir(); + let mut mine = FaceShardStore::open(&mine_dir).unwrap(); + // Already have image 1, by our own indexing. + mine.put_image(1, "w600k_mbf", 1024, &[face(1, 9)]).unwrap(); + + assert_eq!(mine.merge_shard(&peer_shard).unwrap(), 2); + assert!(mine.contains(2, "w600k_mbf")); + assert!(mine.contains(3, "w600k_mbf")); + + // Ours was not overwritten by theirs. + let (faces, _) = mine.get_image(1, "w600k_mbf").unwrap().unwrap(); + assert_eq!(faces[0].landmarks[0], 9); + } + + #[test] + fn merging_the_same_shard_twice_adopts_nothing_the_second_time() { + let peer_dir = tempdir(); + let mut peer = FaceShardStore::open(&peer_dir).unwrap(); + peer.put_image(1, "w600k_mbf", 1024, &[face(1, 1)]).unwrap(); + let peer_shard = peer.shard_path(0); + + let mine_dir = tempdir(); + let mut mine = FaceShardStore::open(&mine_dir).unwrap(); + assert_eq!(mine.merge_shard(&peer_shard).unwrap(), 1); + assert_eq!(mine.merge_shard(&peer_shard).unwrap(), 0); + } + + #[test] + fn the_adopted_ledger_stops_a_shard_being_fetched_twice() { + let dir = tempdir(); + let s = FaceShardStore::open(&dir).unwrap(); + assert!(!s.has_adopted("shard-abc-0000.sqlite", 100)); + s.mark_adopted("shard-abc-0000.sqlite", 100).unwrap(); + assert!(s.has_adopted("shard-abc-0000.sqlite", 100)); + // A shard that grew is the *open* one and must be re-read. + assert!(!s.has_adopted("shard-abc-0000.sqlite", 200)); + } + + #[test] + fn two_stores_get_different_client_ids() { + let a = FaceShardStore::open(&tempdir()).unwrap(); + let b = FaceShardStore::open(&tempdir()).unwrap(); + assert_ne!(a.client_id(), b.client_id()); + assert!(a.shard_path(0).to_string_lossy().contains(a.client_id())); + } + + #[test] + fn a_client_id_survives_reopening() { + let dir = tempdir(); + let first = FaceShardStore::open(&dir).unwrap().client_id().to_string(); + let second = FaceShardStore::open(&dir).unwrap().client_id().to_string(); + assert_eq!(first, second); + } +} + +#[cfg(test)] +mod catalog_round_trip { + use super::*; + use crate::{faces, schema}; + + fn tempdir(tag: &str) -> PathBuf { + use std::sync::atomic::{AtomicU64, Ordering}; + static N: AtomicU64 = AtomicU64::new(0); + let d = std::env::temp_dir().join(format!( + "dr-face-rt-{tag}-{}-{}", + std::process::id(), + N.fetch_add(1, Ordering::Relaxed) + )); + let _ = std::fs::remove_dir_all(&d); + std::fs::create_dir_all(&d).unwrap(); + d + } + + /// A device holding photographs keyed by `oc:fileid`, as they are on the + /// server. Row ids differ between devices; the file id is what agrees. + fn device(images: &[(i64, i64)]) -> Connection { + let c = Connection::open_in_memory().unwrap(); + schema::migrate(&c).unwrap(); + c.execute( + "INSERT INTO roots(id, kind, label) VALUES (1, 'local', 'lib')", + [], + ) + .unwrap(); + for (image_id, file_id) in images { + c.execute( + "INSERT INTO images(id, root_id, source_ref, added_at) VALUES (?1, 1, ?2, 0)", + rusqlite::params![image_id, format!("IMG_{file_id}.CR3")], + ) + .unwrap(); + c.execute( + "INSERT INTO remote(image_id, file_id) VALUES (?1, ?2)", + rusqlite::params![image_id, file_id], + ) + .unwrap(); + } + c + } + + fn detected(seed: u8) -> faces::DetectedFace { + faces::DetectedFace { + x: 0.1, + y: 0.2, + w: 0.15, + h: 0.2, + landmarks: [ + (0.11, 0.21), + (0.2, 0.21), + (0.15, 0.26), + (0.12, 0.3), + (0.19, 0.3), + ], + confidence: 0.9, + embedding: vec![seed; 1024], + crop_px: 180.0, + model_id: "w600k_mbf".into(), + } + } + + /// The whole point of the module: device B adopts A's work rather than + /// spending two hours re-detecting the same library. + #[test] + fn a_second_device_adopts_the_first_devices_indexing() { + let a = device(&[(1, 5001), (2, 5002), (3, 5003)]); + let b = device(&[(90, 5001), (91, 5002), (92, 5003)]); + + // A indexes everything: two images with faces, one without. + faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)]) + .unwrap(); + faces::record_detections( + &a, + dr_types::ImageId(2), + "w600k_mbf", + 1024, + &[detected(2), detected(3)], + ) + .unwrap(); + faces::record_detections(&a, dr_types::ImageId(3), "w600k_mbf", 1024, &[]).unwrap(); + + let mut store_a = FaceShardStore::open(&tempdir("a")).unwrap(); + assert_eq!(export_to_shards(&a, &mut store_a, "w600k_mbf").unwrap(), 3); + + assert_eq!(faces::coverage(&b, "w600k_mbf").unwrap().indexed, 0); + + let mut store_b = FaceShardStore::open(&tempdir("b")).unwrap(); + store_b.merge_shard(&store_a.shard_path(0)).unwrap(); + assert_eq!(import_from_shards(&b, &store_b, "w600k_mbf").unwrap(), 3); + + let cov = faces::coverage(&b, "w600k_mbf").unwrap(); + assert_eq!(cov.indexed, 3, "B did not adopt A's coverage"); + assert_eq!(cov.faces, 3); + // The faceless image must arrive as *indexed with none*, or B + // re-detects it for ever. + assert_eq!(cov.without_faces, 1); + assert_eq!(cov.outstanding(), 0); + + // The embeddings survived intact, or the two devices would cluster the + // same library differently. + let got = faces::for_image(&b, dr_types::ImageId(90)).unwrap(); + assert_eq!(got.len(), 1); + assert!((got[0].crop_px - 180.0).abs() < 1e-3); + assert!((got[0].landmarks[2].0 - 0.15).abs() < 1e-5); + let emb = faces::embeddings(&b, "w600k_mbf").unwrap(); + assert!(emb.iter().any(|(_, _, blob, _)| blob[0] == 1)); + } + + #[test] + fn importing_does_not_overwrite_work_this_device_already_did() { + let a = device(&[(1, 5001)]); + let b = device(&[(50, 5001)]); + + faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(7)]) + .unwrap(); + // B indexed the same photograph itself and got a different embedding. + faces::record_detections(&b, dr_types::ImageId(50), "w600k_mbf", 1024, &[detected(9)]) + .unwrap(); + + let mut store = FaceShardStore::open(&tempdir("keep")).unwrap(); + export_to_shards(&a, &mut store, "w600k_mbf").unwrap(); + + assert_eq!( + import_from_shards(&b, &store, "w600k_mbf").unwrap(), + 0, + "a peer's copy replaced work this device had already done" + ); + let emb = faces::embeddings(&b, "w600k_mbf").unwrap(); + assert_eq!(emb[0].2[0], 9, "B's own embedding was overwritten"); + } + + /// A device holding a subset of the library takes only its own part. + #[test] + fn a_device_ignores_faces_for_photographs_it_does_not_have() { + let a = device(&[(1, 5001), (2, 5002)]); + let b = device(&[(80, 5002)]); + + faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)]) + .unwrap(); + faces::record_detections(&a, dr_types::ImageId(2), "w600k_mbf", 1024, &[detected(2)]) + .unwrap(); + + let mut store = FaceShardStore::open(&tempdir("subset")).unwrap(); + export_to_shards(&a, &mut store, "w600k_mbf").unwrap(); + + assert_eq!(import_from_shards(&b, &store, "w600k_mbf").unwrap(), 1); + assert_eq!(faces::coverage(&b, "w600k_mbf").unwrap().indexed, 1); + } + + #[test] + fn exporting_twice_sends_only_what_is_new() { + let a = device(&[(1, 5001), (2, 5002)]); + faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)]) + .unwrap(); + + let mut store = FaceShardStore::open(&tempdir("incr")).unwrap(); + assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 1); + assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 0); + + faces::record_detections(&a, dr_types::ImageId(2), "w600k_mbf", 1024, &[detected(2)]) + .unwrap(); + assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 1); + } + + #[test] + fn a_local_only_image_with_no_file_id_is_skipped_rather_than_failing() { + let a = device(&[(1, 5001)]); + a.execute( + "INSERT INTO images(id, root_id, source_ref, added_at) VALUES (2, 1, 'local.CR3', 0)", + [], + ) + .unwrap(); + faces::record_detections(&a, dr_types::ImageId(1), "w600k_mbf", 1024, &[detected(1)]) + .unwrap(); + faces::record_detections(&a, dr_types::ImageId(2), "w600k_mbf", 1024, &[detected(2)]) + .unwrap(); + + let mut store = FaceShardStore::open(&tempdir("localonly")).unwrap(); + assert_eq!(export_to_shards(&a, &mut store, "w600k_mbf").unwrap(), 1); + } +} diff --git a/core/dr-catalog/src/lib.rs b/core/dr-catalog/src/lib.rs index 45e69e0..c99c94c 100644 --- a/core/dr-catalog/src/lib.rs +++ b/core/dr-catalog/src/lib.rs @@ -37,6 +37,7 @@ pub mod cache; pub mod collections; pub mod dedup; pub mod error; +pub mod face_shard; pub mod faces; pub mod jobs; pub mod keywords; @@ -53,6 +54,7 @@ pub use cache::{Budget, Cache, DEFAULT_BUDGET_BYTES}; pub use collections::{Collection, CollectionKind, TreeRow}; pub use dedup::{seen_by_content, seen_by_metadata, set_content_hash}; pub use error::CatalogError; +pub use face_shard::{FaceShardStore, SharedFace}; pub use faces::{Calibration, DetectedFace, Face, FaceId, Person, PersonId}; pub use jobs::{Job, JobKind, Priority}; pub use keywords::{Coverage, Keyword, KeywordId, SelectionKeyword}; diff --git a/docs/traceability.md b/docs/traceability.md index a7c89ae..2bd0936 100644 --- a/docs/traceability.md +++ b/docs/traceability.md @@ -9,8 +9,8 @@ Denominators are parsed from [`requirements.md`](requirements.md) at run time, n | Metric | Value | |---|---| -| Source files scanned | 242 | -| TRACES tags found | 628 | +| Source files scanned | 243 | +| TRACES tags found | 629 | | Requirements defined | 177 | | Requirements covered | 97 | | **Coverage** | **54.8%** (97/177) | @@ -54,7 +54,7 @@ _None._ | FR-CULL-12 | [`core/dr-catalog/src/faces.rs:1`](../core/dr-catalog/src/faces.rs#L1), [`core/dr-catalog/src/schema.rs:384`](../core/dr-catalog/src/schema.rs#L384), [`ui/dr-ui/src/identity.rs:1`](../ui/dr-ui/src/identity.rs#L1), [`ui/dr-ui/ui/identity.slint:1`](../ui/dr-ui/ui/identity.slint#L1) | | FR-CULL-2 | [`core/dr-decode/src/locate.rs:1`](../core/dr-decode/src/locate.rs#L1), [`core/dr-decode/src/preview.rs:161`](../core/dr-decode/src/preview.rs#L161), [`ui/dr-ui/src/import.rs:464`](../ui/dr-ui/src/import.rs#L464) | | FR-CULL-4 | [`core/dr-catalog/src/rating.rs:1`](../core/dr-catalog/src/rating.rs#L1), [`core/dr-pipeline/src/sidecar.rs:137`](../core/dr-pipeline/src/sidecar.rs#L137), [`ui/dr-ui/src/library.rs:203`](../ui/dr-ui/src/library.rs#L203), [`ui/dr-ui/src/library.rs:364`](../ui/dr-ui/src/library.rs#L364) | -| FR-CULL-8 | [`core/dr-catalog/src/faces.rs:1`](../core/dr-catalog/src/faces.rs#L1), [`core/dr-catalog/src/schema.rs:343`](../core/dr-catalog/src/schema.rs#L343), [`core/dr-catalog/src/schema.rs:384`](../core/dr-catalog/src/schema.rs#L384), [`ui/dr-ui/src/faces.rs:1`](../ui/dr-ui/src/faces.rs#L1) | +| FR-CULL-8 | [`core/dr-catalog/src/face_shard.rs:1`](../core/dr-catalog/src/face_shard.rs#L1), [`core/dr-catalog/src/faces.rs:1`](../core/dr-catalog/src/faces.rs#L1), [`core/dr-catalog/src/schema.rs:343`](../core/dr-catalog/src/schema.rs#L343), [`core/dr-catalog/src/schema.rs:384`](../core/dr-catalog/src/schema.rs#L384), [`ui/dr-ui/src/faces.rs:1`](../ui/dr-ui/src/faces.rs#L1) | | FR-CULL-9 | [`core/dr-catalog/src/faces.rs:1`](../core/dr-catalog/src/faces.rs#L1), [`core/dr-catalog/src/schema.rs:384`](../core/dr-catalog/src/schema.rs#L384), [`ui/dr-ui/src/faces.rs:1`](../ui/dr-ui/src/faces.rs#L1), [`ui/dr-ui/src/identity_ui.rs:1`](../ui/dr-ui/src/identity_ui.rs#L1) | | FR-DEV-2 | [`core/dr-pipeline/src/operation.rs:352`](../core/dr-pipeline/src/operation.rs#L352) | | FR-DEV-3 | [`core/dr-gpu/src/adjust.rs:2167`](../core/dr-gpu/src/adjust.rs#L2167), [`core/dr-gpu/src/adjust.rs:651`](../core/dr-gpu/src/adjust.rs#L651), [`core/dr-gpu/src/adjust.rs:770`](../core/dr-gpu/src/adjust.rs#L770), [`core/dr-gpu/src/adjust.rs:84`](../core/dr-gpu/src/adjust.rs#L84), [`core/dr-gpu/tests/tone_curve.rs:1`](../core/dr-gpu/tests/tone_curve.rs#L1), [`core/dr-pipeline/src/detail.rs:364`](../core/dr-pipeline/src/detail.rs#L364), [`core/dr-pipeline/src/detail.rs:439`](../core/dr-pipeline/src/detail.rs#L439), [`core/dr-pipeline/src/framing.rs:188`](../core/dr-pipeline/src/framing.rs#L188), [`core/dr-pipeline/src/framing.rs:602`](../core/dr-pipeline/src/framing.rs#L602), [`core/dr-pipeline/src/graph.rs:154`](../core/dr-pipeline/src/graph.rs#L154), [`core/dr-pipeline/src/graph.rs:457`](../core/dr-pipeline/src/graph.rs#L457), [`core/dr-pipeline/src/mask.rs:120`](../core/dr-pipeline/src/mask.rs#L120), [`core/dr-pipeline/src/operation.rs:299`](../core/dr-pipeline/src/operation.rs#L299), [`core/dr-pipeline/src/operation.rs:473`](../core/dr-pipeline/src/operation.rs#L473), [`core/dr-pipeline/src/ops/capture_sharpen.rs:1`](../core/dr-pipeline/src/ops/capture_sharpen.rs#L1), [`core/dr-pipeline/src/ops/capture_sharpen.rs:207`](../core/dr-pipeline/src/ops/capture_sharpen.rs#L207), [`core/dr-pipeline/src/ops/curve.rs:1`](../core/dr-pipeline/src/ops/curve.rs#L1), [`core/dr-pipeline/src/ops/curve.rs:218`](../core/dr-pipeline/src/ops/curve.rs#L218), [`core/dr-pipeline/src/ops/curve.rs:631`](../core/dr-pipeline/src/ops/curve.rs#L631), [`core/dr-pipeline/src/ops/curve.rs:99`](../core/dr-pipeline/src/ops/curve.rs#L99), [`core/dr-pipeline/src/ops/local_contrast.rs:1`](../core/dr-pipeline/src/ops/local_contrast.rs#L1), [`core/dr-pipeline/src/ops/noise_reduction.rs:1`](../core/dr-pipeline/src/ops/noise_reduction.rs#L1), [`core/dr-pipeline/src/ops/noise_reduction.rs:270`](../core/dr-pipeline/src/ops/noise_reduction.rs#L270), [`core/dr-pipeline/src/sidecar.rs:1343`](../core/dr-pipeline/src/sidecar.rs#L1343), [`core/dr-pipeline/src/sidecar.rs:1399`](../core/dr-pipeline/src/sidecar.rs#L1399), [`core/dr-pipeline/src/sidecar.rs:158`](../core/dr-pipeline/src/sidecar.rs#L158), [`core/dr-pipeline/tests/tone_curve.rs:1`](../core/dr-pipeline/tests/tone_curve.rs#L1), [`ui/dr-ui/src/develop.rs:101`](../ui/dr-ui/src/develop.rs#L101), [`ui/dr-ui/src/develop.rs:1111`](../ui/dr-ui/src/develop.rs#L1111), [`ui/dr-ui/src/develop.rs:132`](../ui/dr-ui/src/develop.rs#L132), [`ui/dr-ui/src/develop.rs:1584`](../ui/dr-ui/src/develop.rs#L1584), [`ui/dr-ui/src/develop.rs:1599`](../ui/dr-ui/src/develop.rs#L1599), [`ui/dr-ui/src/develop.rs:1621`](../ui/dr-ui/src/develop.rs#L1621), [`ui/dr-ui/src/develop.rs:1753`](../ui/dr-ui/src/develop.rs#L1753), [`ui/dr-ui/src/develop.rs:1831`](../ui/dr-ui/src/develop.rs#L1831), [`ui/dr-ui/src/develop.rs:238`](../ui/dr-ui/src/develop.rs#L238), [`ui/dr-ui/src/develop.rs:270`](../ui/dr-ui/src/develop.rs#L270), [`ui/dr-ui/src/develop.rs:2762`](../ui/dr-ui/src/develop.rs#L2762), [`ui/dr-ui/src/develop.rs:3194`](../ui/dr-ui/src/develop.rs#L3194), [`ui/dr-ui/src/develop.rs:3248`](../ui/dr-ui/src/develop.rs#L3248), [`ui/dr-ui/src/develop.rs:3292`](../ui/dr-ui/src/develop.rs#L3292), [`ui/dr-ui/src/develop.rs:3342`](../ui/dr-ui/src/develop.rs#L3342), [`ui/dr-ui/src/develop.rs:506`](../ui/dr-ui/src/develop.rs#L506), [`ui/dr-ui/src/develop.rs:544`](../ui/dr-ui/src/develop.rs#L544), [`ui/dr-ui/src/lib.rs:1323`](../ui/dr-ui/src/lib.rs#L1323), [`ui/dr-ui/src/lib.rs:1982`](../ui/dr-ui/src/lib.rs#L1982), [`ui/dr-ui/src/lib.rs:293`](../ui/dr-ui/src/lib.rs#L293), [`ui/dr-ui/src/library.rs:411`](../ui/dr-ui/src/library.rs#L411), [`ui/dr-ui/src/masks_ui.rs:218`](../ui/dr-ui/src/masks_ui.rs#L218), [`ui/dr-ui/src/masks_ui.rs:41`](../ui/dr-ui/src/masks_ui.rs#L41), [`ui/dr-ui/src/masks_ui.rs:816`](../ui/dr-ui/src/masks_ui.rs#L816), [`ui/dr-ui/src/masks_ui.rs:930`](../ui/dr-ui/src/masks_ui.rs#L930), [`ui/dr-ui/src/segmentation.rs:218`](../ui/dr-ui/src/segmentation.rs#L218), [`ui/dr-ui/ui/app.slint:1985`](../ui/dr-ui/ui/app.slint#L1985), [`ui/dr-ui/ui/app.slint:949`](../ui/dr-ui/ui/app.slint#L949) | @@ -92,7 +92,7 @@ _None._ | FR-NC-6a | [`core/dr-catalog/src/cache.rs:1`](../core/dr-catalog/src/cache.rs#L1), [`core/dr-catalog/src/schema.rs:555`](../core/dr-catalog/src/schema.rs#L555), [`core/dr-catalog/src/schema.rs:956`](../core/dr-catalog/src/schema.rs#L956), [`core/dr-types/src/selector.rs:1`](../core/dr-types/src/selector.rs#L1), [`core/dr-types/src/settings.rs:1`](../core/dr-types/src/settings.rs#L1), [`ui/dr-ui/src/collections_ui.rs:3195`](../ui/dr-ui/src/collections_ui.rs#L3195), [`ui/dr-ui/src/collections_ui.rs:621`](../ui/dr-ui/src/collections_ui.rs#L621), [`ui/dr-ui/src/collections_ui.rs:683`](../ui/dr-ui/src/collections_ui.rs#L683), [`ui/dr-ui/src/lib.rs:1658`](../ui/dr-ui/src/lib.rs#L1658), [`ui/dr-ui/src/lib.rs:2455`](../ui/dr-ui/src/lib.rs#L2455), [`ui/dr-ui/src/library.rs:1324`](../ui/dr-ui/src/library.rs#L1324), [`ui/dr-ui/src/library.rs:1347`](../ui/dr-ui/src/library.rs#L1347), [`ui/dr-ui/src/library.rs:1505`](../ui/dr-ui/src/library.rs#L1505), [`ui/dr-ui/src/library_ui.rs:1065`](../ui/dr-ui/src/library_ui.rs#L1065), [`ui/dr-ui/src/library_ui.rs:1138`](../ui/dr-ui/src/library_ui.rs#L1138), [`ui/dr-ui/src/library_ui.rs:1239`](../ui/dr-ui/src/library_ui.rs#L1239), [`ui/dr-ui/src/library_ui.rs:1294`](../ui/dr-ui/src/library_ui.rs#L1294), [`ui/dr-ui/src/library_ui.rs:1412`](../ui/dr-ui/src/library_ui.rs#L1412), [`ui/dr-ui/src/library_ui.rs:1531`](../ui/dr-ui/src/library_ui.rs#L1531), [`ui/dr-ui/src/library_ui.rs:2058`](../ui/dr-ui/src/library_ui.rs#L2058), [`ui/dr-ui/src/library_ui.rs:267`](../ui/dr-ui/src/library_ui.rs#L267), [`ui/dr-ui/src/library_ui.rs:278`](../ui/dr-ui/src/library_ui.rs#L278), [`ui/dr-ui/src/library_ui.rs:286`](../ui/dr-ui/src/library_ui.rs#L286), [`ui/dr-ui/src/library_ui.rs:298`](../ui/dr-ui/src/library_ui.rs#L298), [`ui/dr-ui/src/library_ui.rs:307`](../ui/dr-ui/src/library_ui.rs#L307), [`ui/dr-ui/src/library_ui.rs:399`](../ui/dr-ui/src/library_ui.rs#L399), [`ui/dr-ui/src/library_ui.rs:409`](../ui/dr-ui/src/library_ui.rs#L409), [`ui/dr-ui/src/library_ui.rs:451`](../ui/dr-ui/src/library_ui.rs#L451), [`ui/dr-ui/src/library_ui.rs:514`](../ui/dr-ui/src/library_ui.rs#L514), [`ui/dr-ui/src/library_ui.rs:5239`](../ui/dr-ui/src/library_ui.rs#L5239), [`ui/dr-ui/src/library_ui.rs:5257`](../ui/dr-ui/src/library_ui.rs#L5257), [`ui/dr-ui/src/library_ui.rs:5269`](../ui/dr-ui/src/library_ui.rs#L5269), [`ui/dr-ui/src/library_ui.rs:545`](../ui/dr-ui/src/library_ui.rs#L545), [`ui/dr-ui/src/library_ui.rs:557`](../ui/dr-ui/src/library_ui.rs#L557), [`ui/dr-ui/src/settings_store.rs:1`](../ui/dr-ui/src/settings_store.rs#L1), [`ui/dr-ui/src/settings_ui.rs:1`](../ui/dr-ui/src/settings_ui.rs#L1), [`ui/dr-ui/ui/app.slint:2490`](../ui/dr-ui/ui/app.slint#L2490), [`ui/dr-ui/ui/app.slint:612`](../ui/dr-ui/ui/app.slint#L612), [`ui/dr-ui/ui/collections.slint:249`](../ui/dr-ui/ui/collections.slint#L249), [`ui/dr-ui/ui/collections.slint:369`](../ui/dr-ui/ui/collections.slint#L369), [`ui/dr-ui/ui/collections.slint:52`](../ui/dr-ui/ui/collections.slint#L52), [`ui/dr-ui/ui/collections.slint:682`](../ui/dr-ui/ui/collections.slint#L682), [`ui/dr-ui/ui/collections.slint:84`](../ui/dr-ui/ui/collections.slint#L84), [`ui/dr-ui/ui/icons.slint:260`](../ui/dr-ui/ui/icons.slint#L260), [`ui/dr-ui/ui/library.slint:974`](../ui/dr-ui/ui/library.slint#L974) | | FR-NC-6b | [`ui/dr-ui/src/library_ui.rs:1294`](../ui/dr-ui/src/library_ui.rs#L1294) | | FR-NC-6c | [`core/dr-sync-nextcloud/src/desktop_client.rs:30`](../core/dr-sync-nextcloud/src/desktop_client.rs#L30), [`core/dr-types/src/lib.rs:119`](../core/dr-types/src/lib.rs#L119), [`core/dr-types/src/lib.rs:201`](../core/dr-types/src/lib.rs#L201), [`ui/dr-ui/src/activity.rs:1`](../ui/dr-ui/src/activity.rs#L1), [`ui/dr-ui/src/collections_ui.rs:3195`](../ui/dr-ui/src/collections_ui.rs#L3195), [`ui/dr-ui/src/collections_ui.rs:621`](../ui/dr-ui/src/collections_ui.rs#L621), [`ui/dr-ui/src/collections_ui.rs:683`](../ui/dr-ui/src/collections_ui.rs#L683), [`ui/dr-ui/src/library_ui.rs:1065`](../ui/dr-ui/src/library_ui.rs#L1065), [`ui/dr-ui/src/library_ui.rs:1138`](../ui/dr-ui/src/library_ui.rs#L1138), [`ui/dr-ui/ui/collections.slint:249`](../ui/dr-ui/ui/collections.slint#L249), [`ui/dr-ui/ui/collections.slint:682`](../ui/dr-ui/ui/collections.slint#L682), [`ui/dr-ui/ui/icons.slint:260`](../ui/dr-ui/ui/icons.slint#L260) | -| FR-NC-7 | [`core/dr-sync-nextcloud/src/lib.rs:95`](../core/dr-sync-nextcloud/src/lib.rs#L95), [`ui/dr-ui/src/derived_sync.rs:1`](../ui/dr-ui/src/derived_sync.rs#L1), [`ui/dr-ui/src/library.rs:2624`](../ui/dr-ui/src/library.rs#L2624), [`ui/dr-ui/src/library_ui.rs:3608`](../ui/dr-ui/src/library_ui.rs#L3608), [`ui/dr-ui/ui/settings.slint:340`](../ui/dr-ui/ui/settings.slint#L340) | +| FR-NC-7 | [`core/dr-catalog/src/face_shard.rs:1`](../core/dr-catalog/src/face_shard.rs#L1), [`core/dr-sync-nextcloud/src/lib.rs:95`](../core/dr-sync-nextcloud/src/lib.rs#L95), [`ui/dr-ui/src/derived_sync.rs:1`](../ui/dr-ui/src/derived_sync.rs#L1), [`ui/dr-ui/src/library.rs:2624`](../ui/dr-ui/src/library.rs#L2624), [`ui/dr-ui/src/library_ui.rs:3608`](../ui/dr-ui/src/library_ui.rs#L3608), [`ui/dr-ui/ui/settings.slint:340`](../ui/dr-ui/ui/settings.slint#L340) | | FR-NC-7a | [`core/dr-ingest/src/layout.rs:1`](../core/dr-ingest/src/layout.rs#L1), [`core/dr-sync/src/upload.rs:1`](../core/dr-sync/src/upload.rs#L1), [`core/dr-sync/src/upload.rs:40`](../core/dr-sync/src/upload.rs#L40), [`core/dr-types/src/settings.rs:116`](../core/dr-types/src/settings.rs#L116), [`ui/dr-ui/src/import.rs:1`](../ui/dr-ui/src/import.rs#L1), [`ui/dr-ui/src/import.rs:97`](../ui/dr-ui/src/import.rs#L97), [`ui/dr-ui/src/import_ui.rs:1`](../ui/dr-ui/src/import_ui.rs#L1), [`ui/dr-ui/src/lib.rs:1050`](../ui/dr-ui/src/lib.rs#L1050), [`ui/dr-ui/ui/import.slint:5`](../ui/dr-ui/ui/import.slint#L5) | | FR-NC-7b | [`core/dr-ingest/src/lib.rs:733`](../core/dr-ingest/src/lib.rs#L733), [`core/dr-sync/src/upload.rs:1`](../core/dr-sync/src/upload.rs#L1), [`ui/dr-ui/src/import.rs:123`](../ui/dr-ui/src/import.rs#L123), [`ui/dr-ui/src/import.rs:337`](../ui/dr-ui/src/import.rs#L337), [`ui/dr-ui/src/import.rs:585`](../ui/dr-ui/src/import.rs#L585), [`ui/dr-ui/src/import.rs:97`](../ui/dr-ui/src/import.rs#L97), [`ui/dr-ui/src/import_ui.rs:1`](../ui/dr-ui/src/import_ui.rs#L1), [`ui/dr-ui/src/lib.rs:1050`](../ui/dr-ui/src/lib.rs#L1050) | | FR-NC-8 | [`core/dr-pipeline/src/sidecar.rs:120`](../core/dr-pipeline/src/sidecar.rs#L120), [`core/dr-pipeline/src/sidecar.rs:90`](../core/dr-pipeline/src/sidecar.rs#L90), [`ui/dr-ui/src/lib.rs:1627`](../ui/dr-ui/src/lib.rs#L1627), [`ui/dr-ui/src/library.rs:364`](../ui/dr-ui/src/library.rs#L364), [`ui/dr-ui/src/library_ui.rs:466`](../ui/dr-ui/src/library_ui.rs#L466) | @@ -125,7 +125,7 @@ _None._ | NFR-R7 | [`core/dr-gpu/src/error.rs:1`](../core/dr-gpu/src/error.rs#L1) | | NFR-R8 | [`core/dr-gpu/src/error.rs:1`](../core/dr-gpu/src/error.rs#L1) | | NFR-RES-1 | [`core/dr-pipeline/src/history.rs:55`](../core/dr-pipeline/src/history.rs#L55), [`ui/dr-ui/src/lib.rs:62`](../ui/dr-ui/src/lib.rs#L62) | -| NFR-RES-4 | [`core/dr-catalog/src/cache.rs:1`](../core/dr-catalog/src/cache.rs#L1), [`core/dr-catalog/src/schema.rs:555`](../core/dr-catalog/src/schema.rs#L555), [`core/dr-thumbs/src/codec.rs:1`](../core/dr-thumbs/src/codec.rs#L1), [`core/dr-thumbs/src/lib.rs:1`](../core/dr-thumbs/src/lib.rs#L1), [`core/dr-thumbs/src/lib.rs:376`](../core/dr-thumbs/src/lib.rs#L376), [`ui/dr-ui/src/library.rs:2598`](../ui/dr-ui/src/library.rs#L2598) | +| NFR-RES-4 | [`core/dr-catalog/src/cache.rs:1`](../core/dr-catalog/src/cache.rs#L1), [`core/dr-catalog/src/face_shard.rs:1`](../core/dr-catalog/src/face_shard.rs#L1), [`core/dr-catalog/src/schema.rs:555`](../core/dr-catalog/src/schema.rs#L555), [`core/dr-thumbs/src/codec.rs:1`](../core/dr-thumbs/src/codec.rs#L1), [`core/dr-thumbs/src/lib.rs:1`](../core/dr-thumbs/src/lib.rs#L1), [`core/dr-thumbs/src/lib.rs:376`](../core/dr-thumbs/src/lib.rs#L376), [`ui/dr-ui/src/library.rs:2598`](../ui/dr-ui/src/library.rs#L2598) | | NFR-SEC-1 | [`core/dr-decode/src/error.rs:1`](../core/dr-decode/src/error.rs#L1) | | NFR-SEC-5 | [`core/dr-catalog/src/faces.rs:1`](../core/dr-catalog/src/faces.rs#L1), [`core/dr-catalog/src/schema.rs:384`](../core/dr-catalog/src/schema.rs#L384), [`ui/dr-ui/src/faces.rs:1`](../ui/dr-ui/src/faces.rs#L1), [`ui/dr-ui/src/identity.rs:1`](../ui/dr-ui/src/identity.rs#L1), [`ui/dr-ui/src/identity_ui.rs:1`](../ui/dr-ui/src/identity_ui.rs#L1), [`ui/dr-ui/ui/identity.slint:1`](../ui/dr-ui/ui/identity.slint#L1) | | R1 | [`tools/traceability/src/lib.rs:495`](../tools/traceability/src/lib.rs#L495), [`tools/traceability/src/lib.rs:499`](../tools/traceability/src/lib.rs#L499) | diff --git a/ui/dr-ui/examples/face_index.rs b/ui/dr-ui/examples/face_index.rs index a70de26..68a33b6 100644 --- a/ui/dr-ui/examples/face_index.rs +++ b/ui/dr-ui/examples/face_index.rs @@ -4,6 +4,7 @@ //! indexes whatever is missing one. //! //! cargo run -p dr-ui --example face_index -- CATALOG.db THUMBS_DIR [--run DET.onnx EMB.onnx] +//! cargo run -p dr-ui --example face_index -- CATALOG.db THUMBS_DIR --cluster //! //! # Why this exists beside the button in the Identity screen //! @@ -81,6 +82,23 @@ fn main() { println!(" ready to index {}", audit.ready); println!(" awaiting proxy {}", audit.awaiting_proxy); + // Grouping is a separate step from indexing on purpose: it is a + // whole-library operation over the embeddings detection produced, and it is + // worth running *after* a sweep rather than during one (catalog.md §10.2). + if args.iter().any(|a| a == "--cluster") { + match dr_ui::faces::recluster(&catalog, MODEL_ID, dr_face::DEFAULT_MERGE_PROBABILITY) { + Ok((suggested, created)) => { + println!("\nclustering: {suggested} suggestion(s), {created} new group(s)"); + report_people(&catalog); + } + Err(e) => { + eprintln!("clustering failed: {e}"); + std::process::exit(1); + } + } + return; + } + let run = args.iter().position(|a| a == "--run"); let Some(i) = run else { if audit.coverage.is_complete() { @@ -129,7 +147,7 @@ fn main() { seen += 1; // One line per image would be thousands of lines; one per // twenty-five is enough to see it moving and to estimate. - if seen % 25 == 0 || faces > 0 { + if seen.is_multiple_of(25) || faces > 0 { let rate = seen as f64 / start.elapsed().as_secs_f64().max(1e-6); println!( " {seen}/{total} {:.2} img/s ~{:.0} min left", @@ -154,4 +172,31 @@ fn main() { if let Ok(a) = faces::audit(&catalog, &store, MODEL_ID) { println!("{}", a.summary()); } + println!("\nrun again with --cluster to group these faces into people."); +} + +/// What the clustering proposed, largest group first. +fn report_people(catalog: &Catalog) { + let Ok(people) = dr_catalog::faces::people(catalog.connection()) else { + return; + }; + if people.is_empty() { + println!("no groups — too few faces, or none similar enough to group."); + return; + } + println!("\n{} group(s):", people.len()); + for p in people.iter().take(30) { + let name = if p.name.is_empty() { + "(unnamed)".to_string() + } else { + p.name.clone() + }; + println!( + " {name:<24} {} confirmed, {} suggested", + p.confirmed_faces, p.suggested_faces + ); + } + if people.len() > 30 { + println!(" … and {} more", people.len() - 30); + } } diff --git a/ui/dr-ui/src/derived_sync.rs b/ui/dr-ui/src/derived_sync.rs index 12de793..9333dbd 100644 --- a/ui/dr-ui/src/derived_sync.rs +++ b/ui/dr-ui/src/derived_sync.rs @@ -60,6 +60,16 @@ pub struct SyncReport { /// devices already had gains *no* collection, and reporting only the /// former left the sidebar showing no count beside a full collection. pub members_gained: usize, + + // Face data is counted apart from thumbnails for the same reason keywords + // are counted apart from collections: "adopted 4,812 faces" is a sentence + // the user can act on, and folding it into the thumbnail count would hide + // the one number that says whether this device still has hours of indexing + // ahead of it. + pub face_shards_uploaded: usize, + pub face_shards_downloaded: usize, + /// Images whose faces this device took from a peer instead of detecting. + pub faces_adopted: usize, } impl SyncReport { @@ -68,6 +78,8 @@ impl SyncReport { || self.shards_downloaded > 0 || self.catalog_uploaded || self.catalog_merged + || self.face_shards_uploaded > 0 + || self.face_shards_downloaded > 0 } } @@ -145,6 +157,9 @@ async fn run( let _ = tx.send(SyncMessage::Status("checking thumbnails…".into())); sync_shards(backend, &base, thumbs_dir, scratch, &mut report).await?; + let _ = tx.send(SyncMessage::Status("checking faces…".into())); + sync_face_shards(backend, &base, catalog_path, scratch, &mut report).await?; + let _ = tx.send(SyncMessage::Status("checking collections…".into())); sync_catalog(backend, &base, catalog_path, scratch, &mut report).await?; @@ -296,6 +311,154 @@ async fn sync_shards( Ok(()) } +/// Push and pull face shards, so a second device does not re-index the library. +/// +/// Mirrors [`sync_shards`] deliberately, down to the sealed-shard skip and the +/// adopted ledger: face data has exactly the properties that made that design +/// right for thumbnails. It is bulk, it is immutable once written, and it is +/// byte-identical on every device, because the same model over the same proxy +/// is deterministic. +/// +/// The catalog is the source and the destination; the shards are only the +/// carrier. So this exports the catalog's new faces into the local shard store +/// first, syncs the shards, and imports whatever arrived back into the catalog. +async fn sync_face_shards( + backend: &NextcloudBackend, + base: &RemotePath, + catalog_path: &Path, + scratch: &Path, + report: &mut SyncReport, +) -> Result<(), String> { + use dr_catalog::face_shard::{self, FaceShardStore}; + + // Beside the catalog, next to the thumbnails, and under the same folder the + // scanner already excludes. + let dir = match catalog_path.parent() { + Some(p) => p.join("faces"), + None => return Ok(()), + }; + let mut store = match FaceShardStore::open(&dir) { + Ok(s) => s, + Err(e) => { + // Not a sync failure: the next pass retries once it opens. + log::debug!("face shard store unavailable: {e}"); + return Ok(()); + } + }; + let client = store.client_id().to_string(); + + let model = crate::identity_ui::MODEL_ID; + + // ---- everything this device has detected, into the shards ------------ + if let Ok(catalog) = dr_catalog::Catalog::open(catalog_path) { + match face_shard::export_to_shards(catalog.connection(), &mut store, model) { + Ok(0) => {} + Ok(n) => log::info!("face sync: {n} newly indexed image(s) ready to upload"), + Err(e) => log::warn!("face sync: exporting to shards: {e}"), + } + } + + // Its own folder under the derived directory, so a client that does not + // care about faces lists thumbnails without paging past them. + let face_base = RemotePath::new(format!("{}/faces", base.as_str())); + let _ = backend.create_dir(&face_base).await; + + let remote: std::collections::HashMap = backend + .list(&face_base, None) + .await + .map(|entries| { + entries + .into_iter() + .filter(|e| e.kind == dr_sync::EntryKind::File) + .map(|e| (e.path.name().to_string(), e.size)) + .collect() + }) + .unwrap_or_default(); + + // ---- upload ---------------------------------------------------------- + let local = store.shards().map_err(|e| e.to_string())?; + for shard in &local { + let path = store.shard_path(shard.id); + let Ok(bytes) = std::fs::read(&path) else { + continue; + }; + let name = shard_name(&client, shard.id); + + // Sealed and present means byte-identical, and the client is in the + // name so nobody else could have written it. The open shard goes up + // again whenever its size differs, which is the only way it changes. + let skip = match remote.get(&name) { + Some(_) if shard.sealed => true, + Some(size) => *size == bytes.len() as u64, + None => false, + }; + if skip { + continue; + } + + let target = RemotePath::new(format!("{}/{name}", face_base.as_str())); + match backend.put(&target, bytes, None).await { + Ok(_) => report.face_shards_uploaded += 1, + Err(e) => log::warn!("uploading face shard {name}: {e}"), + } + } + + // ---- download -------------------------------------------------------- + for (name, size) in &remote { + let Some((owner, _)) = parse_shard(name) else { + continue; + }; + if owner == client || store.has_adopted(name, *size) { + continue; + } + + let source = RemotePath::new(format!("{}/{name}", face_base.as_str())); + let bytes = match backend.get(&RemoteId::Path(source), None).await { + Ok(b) => b, + Err(e) => { + log::warn!("downloading face shard {name}: {e}"); + continue; + } + }; + + // Into scratch and merged, never dropped into the store directory: a + // downloaded shard's id is the *other* device's numbering, and two + // devices independently fill shard 0. + let tmp = scratch.join(name); + if std::fs::write(&tmp, &bytes).is_err() { + continue; + } + match store.merge_shard(&tmp) { + Ok(_) => { + report.face_shards_downloaded += 1; + // Recorded only on success, so a failed merge is retried next + // pass rather than written off. + let _ = store.mark_adopted(name, *size); + } + Err(e) => log::warn!("merging face shard {name}: {e}"), + } + let _ = std::fs::remove_file(&tmp); + } + + // ---- and back into the catalog --------------------------------------- + // + // Last, and unconditionally rather than only when something downloaded: a + // previous pass may have merged shards into the store and then failed + // before importing, and this is what recovers from that. + if let Ok(catalog) = dr_catalog::Catalog::open(catalog_path) { + match face_shard::import_from_shards(catalog.connection(), &store, model) { + Ok(0) => {} + Ok(n) => { + report.faces_adopted = n; + log::info!("face sync: adopted {n} image(s) already indexed elsewhere"); + } + Err(e) => log::warn!("face sync: importing from shards: {e}"), + } + } + + Ok(()) +} + /// Whether a flat-named remote shard is this client's own earlier upload. /// /// Before the name carried a client every client wrote `shard-NNNN.sqlite`, so