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