//! TRACES: FR-CULL-5 //! Burst grouping, as the library screen uses it. //! //! [`dr_catalog::bursts`] holds the grouping itself and knows nothing about //! pixels. This is the other half: where the similarity signal comes from, and //! how a group reaches a grid cell. //! //! # The signal comes out of the thumbnail store //! //! A perceptual signature needs pixels, and the cheapest pixels in the app are //! the ones already sitting in `dr-thumbs`: a 256px JPEG per photograph, built //! for the grid, shared between devices, and vastly more resolution than a 9×8 //! reduction can use. So this pass decodes thumbnails, never originals. A //! library that has been browsed — or that has synced somebody else's shards — //! has already paid for every signature it is about to get. //! //! The consequence, stated rather than hidden: **an image with no thumbnail //! gets no signature, and a frame with no signature never joins a burst.** That //! is self-correcting rather than permanent — the next pass finds the thumbnail //! the sweep has since built — and it is the reason this runs when the //! thumbnail sweep finishes rather than on a timer. //! //! # Why a pass and not a job //! //! Hashing is per-image and would make a perfectly good job kind. Grouping is //! not: a burst is a property of a *run* of frames, so a per-image job would //! regroup the library once per photograph. Since the two have to happen in that //! order and the second cannot be split, both live in one pass — the same //! argument docs/dev/catalog.md §10.2 makes for face clustering. //! //! # It is never on the UI thread //! //! Decoding tens of thousands of thumbnails is bounded only by library size, and //! the one thing that must not grow with library size is how long the window //! stops answering (NFR-P9). Cancellation is dropping the receiver; a pass //! abandoned half way leaves the signatures it did compute — they are permanent //! and correct — and the previous grouping intact. use std::cell::{Cell, RefCell}; use std::collections::HashMap; use std::path::PathBuf; use std::sync::mpsc::Receiver; use dr_catalog::bursts::{self, Rules, Signature}; use dr_catalog::Catalog; use dr_thumbs::{ThumbSize, ThumbStore}; use dr_types::ImageId; use slint::{ComponentHandle as _, Model as _}; use crate::{AppWindow, Library}; /// How many signatures are written per transaction. /// /// One transaction per image costs a WAL commit per thumbnail and turns a pass /// over a real library into minutes of fsync; one transaction for the whole pass /// holds a write lock for the duration and loses everything if the app closes. /// A few hundred is the usual answer to that trade. const WRITE_BATCH: usize = 256; /// Progress from a grouping pass. #[derive(Debug, Clone, PartialEq)] pub enum BurstMessage { /// How many images still need a signature. Sent once, before any decoding. Started { to_hash: usize }, /// Cumulative signatures written. Progress { hashed: usize }, /// The pass finished, and this is what the library now looks like. Finished { hashed: usize, bursts: usize, frames: usize, }, /// It did not. Failed(String), } /// Hash whatever is missing a signature, then rebuild the grouping. /// /// Both halves run on a worker thread. The catalog is opened here rather than /// shared with the UI's connection: SQLite connections are not `Send`, and WAL /// is what makes a second one safe while the grid reads (NFR-R1). pub fn spawn_grouping(catalog_path: PathBuf, thumbs_dir: PathBuf) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); std::thread::spawn(move || { let catalog = match Catalog::open(&catalog_path) { Ok(c) => c, Err(e) => { let _ = tx.send(BurstMessage::Failed(format!("cannot open catalog: {e}"))); return; } }; let conn = catalog.connection(); let outstanding = match bursts::images_without_signature(conn) { Ok(v) => v, Err(e) => { let _ = tx.send(BurstMessage::Failed(e.to_string())); return; } }; if tx .send(BurstMessage::Started { to_hash: outstanding.len(), }) .is_err() { return; } // A broken or absent store costs signatures, never correctness: the // grouping still runs over whatever is already hashed, and the images // that missed out are picked up by the next pass. let store = match ThumbStore::open(&thumbs_dir) { Ok(s) => Some(s), Err(e) => { log::warn!("thumbnail store unavailable, not hashing: {e}"); None } }; let hashed = match store { Some(store) => hash_all(conn, &store, &outstanding, &tx), None => 0, }; match bursts::regroup(conn, Rules::default()) { Ok(report) => { log::info!( "bursts: {hashed} signature(s) added, {} group(s) over {} frame(s), \ largest {}", report.bursts, report.frames, report.largest ); let _ = tx.send(BurstMessage::Finished { hashed, bursts: report.bursts, frames: report.frames, }); } Err(e) => { let _ = tx.send(BurstMessage::Failed(e.to_string())); } } }); rx } thread_local! { /// The running pass's drain timer, and whether one is running. /// /// Module-local rather than a pair of fields on the library controller, so /// that everything this feature needs to run lives in this file and the /// screen that starts it is left holding nothing. Safe as a thread local /// because Slint's event loop is single-threaded (NFR-P9) and this is only /// ever touched from it. /// /// The flag is separate because stopping a timer does not drop it: a slot /// tested for emptiness would refuse every pass after the first. static DRAIN: RefCell> = const { RefCell::new(None) }; static RUNNING: Cell = const { Cell::new(false) }; } /// Hash what has become hashable, then rebuild the library's burst grouping. /// /// Fired when the thumbnail sweep finishes, because that is the moment the /// signatures can all be computed: the pass reads thumbnails, and until the /// sweep has run most images have none. Not on a timer, and not after every /// scan — a regroup is cheap but not free, and nothing is waiting on it. /// /// `grouped` is called once, with the number of bursts the library now has, if /// the pass finishes. It is where the caller reloads the grid: the cells hold /// the same photographs they held before — a new burst arrives open — but every /// run of frames now carries a mark it did not have a moment ago, and only a /// reload carries it. /// /// Deliberately silent otherwise. The sweeps around it report progress because /// they run for tens of minutes; this is seconds, and a status line for it would /// be a line the user must read in order to learn nothing. pub fn start_pass(catalog_path: PathBuf, thumbs_dir: PathBuf, grouped: impl Fn(usize) + 'static) { // A second pass would read the same rows and write the same answer over the // first one's transactions. if RUNNING.get() { return; } RUNNING.set(true); let rx = spawn_grouping(catalog_path, thumbs_dir); let timer = slint::Timer::default(); timer.start( slint::TimerMode::Repeated, std::time::Duration::from_millis(400), move || loop { let msg = match rx.try_recv() { Ok(m) => m, Err(std::sync::mpsc::TryRecvError::Empty) => return, Err(std::sync::mpsc::TryRecvError::Disconnected) => { finish(); return; } }; match msg { BurstMessage::Started { to_hash } => { log::info!("burst grouping: {to_hash} image(s) still to hash"); } // Nothing on screen is showing this. Drained rather than // ignored, because an unread channel is a worker that stalls. BurstMessage::Progress { hashed } => { log::debug!("burst grouping: {hashed} hashed so far"); } BurstMessage::Finished { hashed, bursts, frames, } => { log::info!( "burst grouping: {hashed} signature(s) added, \ {bursts} group(s) over {frames} frame(s)" ); finish(); grouped(bursts); return; } BurstMessage::Failed(e) => { log::warn!("burst grouping: {e}"); finish(); return; } } }, ); DRAIN.with(|slot| *slot.borrow_mut() = Some(timer)); } /// Stop draining, and let another pass start. /// /// Stops the timer without dropping it — dropping one from inside its own /// callback is not something to rely on — which is why the flag beside it is /// what actually says whether a pass is running. fn finish() { RUNNING.set(false); DRAIN.with(|slot| { if let Some(timer) = slot.borrow().as_ref() { timer.stop(); } }); } /// Decode each image's stored thumbnail and record its signature. /// /// Returns how many were written. Anything without a stored thumbnail, or whose /// blob will not decode, is simply skipped — it keeps its NULL and comes back /// next time, which is the same treatment `spawn_thumbnails` gives a corrupt /// blob. fn hash_all( conn: &rusqlite::Connection, store: &ThumbStore, outstanding: &[ImageId], tx: &std::sync::mpsc::Sender, ) -> usize { let keys = thumbnail_keys(conn); let mut pending: Vec<(ImageId, Signature)> = Vec::with_capacity(WRITE_BATCH); let mut written = 0usize; for image in outstanding { let Some(file_id) = keys.get(image).copied() else { continue; }; // The grid class, not the large one: 256px is already thirty times the // detail the reduction keeps, and asking for `Large` would miss most of // the store, which is filled at `Grid`. let Ok(Some(thumb)) = store.get(file_id, ThumbSize::Grid) else { continue; }; let Ok((w, h, rgba)) = dr_thumbs::decode_rgba(&thumb.bytes) else { continue; }; let Some(signature) = bursts::signature_of_rgba(&rgba, w, h) else { continue; }; pending.push((*image, signature)); if pending.len() >= WRITE_BATCH { written += flush(conn, &mut pending); if tx.send(BurstMessage::Progress { hashed: written }).is_err() { // The receiver is gone: the screen has moved on, and finishing // the pass would be work nobody is waiting for. return written; } } } written += flush(conn, &mut pending); written } /// Write one batch of signatures, emptying `pending`. fn flush(conn: &rusqlite::Connection, pending: &mut Vec<(ImageId, Signature)>) -> usize { if pending.is_empty() { return 0; } let n = pending.len(); let tx = match conn.unchecked_transaction() { Ok(t) => t, Err(e) => { log::warn!("recording signatures: {e}"); pending.clear(); return 0; } }; for (image, signature) in pending.drain(..) { if let Err(e) = bursts::set_signature(&tx, image, signature) { log::debug!("recording signature for {}: {e}", image.0); } } match tx.commit() { Ok(()) => n, Err(e) => { log::warn!("committing signatures: {e}"); 0 } } } /// Every image's thumbnail-store key, in one query. /// /// Read whole rather than asked per image: the table is one small row per /// photograph, and a query per image would be tens of thousands of statements /// to save a megabyte. fn thumbnail_keys(conn: &rusqlite::Connection) -> HashMap { let mut out = HashMap::new(); let Ok(mut stmt) = conn.prepare("SELECT image_id, file_id FROM remote") else { return out; }; let Ok(rows) = stmt.query_map([], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?))) else { return out; }; for (image, file) in rows.flatten() { out.insert(ImageId(image as u64), file as u64); } out } /// Fill in the burst badge on every cell of the loaded window. /// /// One query for the window, in the same style as the collection badges and the /// rating counts beside it — a grid that asks the catalog a question per cell is /// a grid that stutters under a finger. /// /// `ids` is the window in grid order, so the row index into the model is the /// index into it. pub fn sync_badges(window: &AppWindow, catalog: &Catalog, ids: &[ImageId]) { if ids.is_empty() { return; } let found = match bursts::memberships(catalog.connection(), ids) { Ok(m) => m, Err(e) => { log::debug!("reading burst membership: {e}"); return; } }; let model = window.global::().get_library_cells(); for (row, id) in ids.iter().enumerate() { let (count, expanded, representative) = match found.get(id) { Some(m) => (m.size as i32, m.expanded, m.representative), // Not in a burst at all, which is most of a library. None => (0, false, false), }; if let Some(mut cell) = model.row_data(row) { if cell.burst_count != count || cell.burst_expanded != expanded || cell.burst_representative != representative { cell.burst_count = count; cell.burst_expanded = expanded; cell.burst_representative = representative; model.set_row_data(row, cell); } } } } /// Open or close the burst one cell belongs to. /// /// Which direction is decided from what the catalog says the group is currently /// doing, rather than from the cell's own `burst-expanded`. The cell is a copy /// of that state and can be one reload behind; the table cannot. /// /// The caller reloads the grid rather than repainting it, because collapsing /// changes what the grid's *query* returns: the row count, the scrollbar and /// the ordinal a scrub resolves all move together, and they can only stay in /// step by being read again together. /// /// Returns whether anything changed, so a click on a cell that is in no burst — /// an entirely normal thing to happen — costs no round trip through the grid. pub fn toggle(catalog: &Catalog, image: ImageId) -> bool { let conn = catalog.connection(); let Ok(found) = bursts::memberships(conn, &[image]) else { return false; }; let Some(membership) = found.get(&image) else { return false; }; match bursts::set_expanded(conn, membership.burst_id, !membership.expanded) { Ok(()) => true, Err(e) => { log::warn!("toggling burst {}: {e}", membership.burst_id.0); false } } } /// Name the frame the burst a cell belongs to folds up to. /// /// The whole of the arbitration is [`bursts::choose_representative`]'s: a group /// stands for one frame, so the previous pick has to go in the same transaction /// the new one arrives in, and the flag the grouping already carries has to be /// corrected there too or the cell would not change until the next pass. /// /// Nothing here decides *which* frame deserves it, and nothing ever will — the /// default is the earliest one, which is a fact about the clock rather than a /// judgement about the photograph (FR-CULL-5). This is the photographer saying /// otherwise about a group they are looking at. /// /// Returns whether the catalog took it, so a caller only repaints when there is /// something to repaint. Choosing a frame that is in no burst is not an error /// and writes nothing; it simply cannot arrive from the grid, where the mark is /// drawn on the members of an open group and nowhere else. pub fn choose(catalog: &Catalog, image: ImageId) -> bool { match bursts::choose_representative(catalog.connection(), image) { Ok(()) => true, Err(e) => { log::warn!("choosing the representative of image {}: {e}", image.0); false } } } #[cfg(test)] mod tests { use super::*; /// A catalog with two frames of one burst, and a thumbnail store holding a /// picture for each. /// /// `name` keeps two tests from sharing a directory, since they run in /// parallel. Same shape as the store tests in `library`, which is also why /// this reaches for `temp_dir` rather than a crate: nothing else here needs /// one. fn library(name: &str) -> (Catalog, PathBuf) { let cat = Catalog::in_memory().unwrap(); let c = cat.connection(); c.execute( "INSERT INTO roots(id, kind, label) VALUES (1, 'local', 'lib')", [], ) .unwrap(); for (id, at) in [(1i64, 1000i64), (2, 1001)] { c.execute( "INSERT INTO images(id, root_id, source_ref, captured_at, camera, added_at) VALUES (?1, 1, ?2, ?3, 'Canon EOS R5', 0)", rusqlite::params![id, format!("IMG_{id}.CR3"), at], ) .unwrap(); c.execute( "INSERT INTO remote(image_id, file_id) VALUES (?1, ?2)", rusqlite::params![id, 100 + id], ) .unwrap(); } let dir = std::env::temp_dir().join(format!("dr-ui-bursts-{name}-{}", std::process::id())); let _ = std::fs::remove_dir_all(&dir); { let mut store = ThumbStore::open(&dir).unwrap(); for (id, shift) in [(1u64, 0usize), (2, 1)] { let rgba = picture(shift); let bytes = dr_thumbs::encode_rgba(64, 64, &rgba).unwrap(); store .put( 100 + id, ThumbSize::Grid, &dr_thumbs::Thumbnail { width: 64, height: 64, bytes, }, ) .unwrap(); } } (cat, dir) } /// A blocky scene, moved sideways by `shift` pixels — one frame of a burst /// and then the next. fn picture(shift: usize) -> Vec { let mut out = vec![0u8; 64 * 64 * 4]; for y in 0..64 { for x in 0..64 { let v = (((x + shift) / 8) * 37 + (y / 8) * 91) as u8; let p = (y * 64 + x) * 4; out[p] = v; out[p + 1] = v; out[p + 2] = v; out[p + 3] = 255; } } out } #[test] fn the_pass_hashes_from_thumbnails_and_groups_what_it_hashed() { let (cat, dir) = library("hashes"); let conn = cat.connection(); let store = ThumbStore::open(&dir).unwrap(); let outstanding = bursts::images_without_signature(conn).unwrap(); assert_eq!(outstanding.len(), 2); let (tx, _rx) = std::sync::mpsc::channel(); let hashed = hash_all(conn, &store, &outstanding, &tx); assert_eq!(hashed, 2, "both thumbnails should have yielded a signature"); let report = bursts::regroup(conn, Rules::default()).unwrap(); assert_eq!(report.bursts, 1, "the two frames were not grouped"); assert_eq!(report.frames, 2); } #[test] fn an_image_with_no_thumbnail_keeps_its_null() { let (cat, dir) = library("no-thumb"); let conn = cat.connection(); conn.execute( "INSERT INTO images(id, root_id, source_ref, captured_at, added_at) VALUES (9, 1, 'IMG_9.CR3', 1002, 0)", [], ) .unwrap(); let store = ThumbStore::open(&dir).unwrap(); let outstanding = bursts::images_without_signature(conn).unwrap(); let (tx, _rx) = std::sync::mpsc::channel(); assert_eq!(hash_all(conn, &store, &outstanding, &tx), 2); // And it is still offered next time, rather than being written off. assert_eq!( bursts::images_without_signature(conn).unwrap(), vec![ImageId(9)] ); } #[test] fn toggling_a_cell_that_is_in_no_burst_changes_nothing() { let (cat, _dir) = library("no-burst"); assert!(!toggle(&cat, ImageId(1))); } #[test] fn choosing_a_frame_moves_the_mark_the_grid_draws() { // The tick a cell wears is filled in from `memberships`, so that is // where the choice has to appear. A pick recorded only in `burst_pick` // would be honoured at the next pass and invisible until then, which is // the whole reason `choose_representative` corrects the grouping in the // same transaction. let (cat, dir) = library("choose"); let conn = cat.connection(); let store = ThumbStore::open(&dir).unwrap(); let outstanding = bursts::images_without_signature(conn).unwrap(); let (tx, _rx) = std::sync::mpsc::channel(); assert_eq!(hash_all(conn, &store, &outstanding, &tx), 2); bursts::regroup(conn, Rules::default()).unwrap(); let mark = |id: ImageId| bursts::memberships(conn, &[id]).unwrap()[&id].representative; assert!( mark(ImageId(1)), "the earliest frame stands for it by default" ); assert!(choose(&cat, ImageId(2))); assert!(mark(ImageId(2)), "the chosen frame did not take the mark"); assert!( !mark(ImageId(1)), "two frames cannot both stand for one moment" ); } #[test] fn choosing_a_frame_that_is_in_no_burst_is_not_an_error() { // Unreachable from the grid, where the mark is drawn on the members of // an open group and nowhere else — but a lone frame arriving here is a // race with a regroup, not a fault. let (cat, _dir) = library("choose-no-burst"); assert!(choose(&cat, ImageId(1))); } }