//! TRACES: FR-CAT-11a | FR-CAT-15 | NFR-P9 //! Proving catalog duplicates the same, and consolidating them, on workers. //! //! `dr_catalog::duplicates` finds the groups and owns the merge; this is the //! half that touches files: reading enough of each copy to prove the bytes //! agree, reading each copy's sidecar to see whether its edit does, and the //! `MOVE` into the trash that a consolidation commits behind. //! //! # The proof, and what it costs //! //! A group is already equal in camera, capture instant and byte count. On top //! of that the check wants the bytes: a stored full digest where every copy //! has one, and otherwise a SHA-256 over the first and last //! [`PROBE_WINDOW`] of each file, read by range through the backend — two //! megabytes a copy rather than the whole RAW. What is read is kept in the //! catalog against the size and mtime it was read at, so a second review //! costs one query. A group where any copy differs, or cannot be read, is //! dropped from the plan and the review says why. //! //! The develop edit is part of the proof. Edits live in the `.drsc` beside //! each copy, so each is read too: every copy's edit empty or identical, and //! the group consolidates keeping one; two different edits, and the group is //! kept whole for virtual copies (FR-CAT-12). //! //! # A group is consolidated whole or not at all //! //! Per group, in this order: //! //! 1. where the survivor has no edit and a copy has the one edit the group //! carries, the survivor's sidecar is written with it; //! 2. every other copy is `MOVE`d into the trash; //! 3. `dr_catalog::duplicates::consolidate` merges and records the trash in //! one transaction. //! //! A failure at 2 or 3 moves back what was moved and puts the survivor's //! sidecar back as it was. A process that dies between 2 and 3 leaves files //! in the trash that the catalog still lists where they were; the trash path //! is a function of the image id, so the next review finds each such file //! there, and consolidating again treats it as already moved. use crate::executors::{self, Executor}; use std::path::PathBuf; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::mpsc::{Receiver, Sender}; use std::sync::Arc; use dr_catalog::duplicates::{Copy, Group, Outcome, Probe, PROBE_WINDOW}; use dr_catalog::{trash, Catalog}; use dr_pipeline::{Sidecar, Version}; use dr_sync::{Connection, RemoteBackend, RemoteError, RemoteId, RemotePath}; use dr_types::ImageId; use sha2::{Digest, Sha256}; use crate::library::sidecar_path; use crate::sidecar_cache::SidecarCache; /// What the check found for one group. #[derive(Debug, Clone, PartialEq, Eq)] pub enum Check { /// Not looked at yet. Unchecked, /// The bytes agree, and so do the edits. Same(Edits), /// Left out of the plan, and why — in words for the review. Skip(String), } /// The develop edits in a group whose edits agree. #[derive(Debug, Clone, PartialEq, Eq, Default)] pub struct Edits { /// Per copy, in the group's order: whether its sidecar carries an edit. pub edited: Vec, } impl Edits { /// A copy holding the group's one edit, to take it from where the /// survivor has none. pub fn source_for(&self, survivor: usize) -> Option { if self.edited.get(survivor).copied().unwrap_or(false) { return None; } self.edited.iter().position(|e| *e) } } /// Why a group was left out, as the review says it. pub const EDITS_DIFFER: &str = "Different edits, kept for virtual copies"; /// A group, the copy chosen to stay, and whether it is in the plan. #[derive(Debug, Clone)] pub struct Reviewed { pub group: Group, pub survivor: usize, pub include: bool, pub check: Check, } impl Reviewed { pub fn new(group: Group) -> Self { let survivor = dr_catalog::duplicates::survivor(&group.copies); Self { group, survivor, include: true, check: Check::Unchecked, } } } /// The dry-run figures the review shows above the list. #[derive(Debug, Clone, Copy, PartialEq, Eq, Default)] pub struct Summary { pub groups: usize, /// Copies that would move to the trash. pub files: usize, /// Groups the check dropped. pub skipped: usize, /// Groups left out by hand. pub excluded: usize, pub unchecked: usize, } pub fn summarise(review: &[Reviewed]) -> Summary { let mut s = Summary::default(); for r in review { match &r.check { Check::Unchecked => s.unchecked += 1, Check::Skip(_) => s.skipped += 1, Check::Same(_) if !r.include => s.excluded += 1, Check::Same(_) => { s.groups += 1; s.files += r.group.copies.len() - 1; } } } s } /// What a worker reports. #[derive(Debug)] pub enum DupMessage { /// One group checked, by its index in the review. Checked { index: usize, check: Check, }, /// One group consolidated. Consolidated { index: usize, survivor: ImageId, outcome: Outcome, }, /// One group could not be consolidated and was left as it was. Failed { index: usize, reason: String, }, Progress { done: usize, total: usize, }, /// The job ended. `stopped` when the server went away or the review /// was closed. Finished { stopped: Option, }, } // ── the edit ───────────────────────────────────────────────────────────── /// A version with everything that is not the edit taken out: identity, /// bookkeeping, and the judgement, which is merged by the catalog instead. fn strip(v: &Version) -> Version { let mut c = v.clone(); c.uuid = String::new(); c.name = String::new(); c.revision = 0; c.device = String::new(); c.modified = 0; c.rating = 0; c.flag = 0; c.label = 0; c.snapshot_of = None; c } /// The edit a sidecar carries, as text two copies can be compared by, or /// `None` where it carries none. /// /// Judgements, revisions, devices and times are not the edit and are left /// out: two copies rated differently but developed alike have one edit. pub fn edit_key(sidecar: &Sidecar) -> Option { let mut sc = sidecar.clone(); sc.fuse_default_versions(Some("default")); let mut blocks: Vec = Vec::new(); for v in sc.versions.values() { let mut s = strip(v); s.is_default = v.is_default; // A snapshot or a virtual copy is named by the photographer; the // default's name is bookkeeping. if !v.is_default { s.name = v.name.clone(); } if s == strip(&Version { is_default: s.is_default, name: s.name.clone(), ..Default::default() }) && v.is_default { continue; } let mut one = Sidecar::new(); s.uuid = "v".into(); one.put(s); blocks.push(one.to_text()); } blocks.sort(); if blocks.is_empty() && sc.derived_from.is_empty() && sc.merge.is_none() { return None; } let mut key = blocks.join("\n"); for d in &sc.derived_from { key.push_str("\nderived_from "); key.push_str(d); } if let Some(m) = &sc.merge { key.push_str("\nmerge "); key.push_str(m); } Some(key) } /// The survivor's sidecar with a copy's edit carried onto it. /// /// The survivor's own judgement stays — the catalog merges those and the /// write after the commit brings them here — and the edit, its snapshots and /// its virtual copies come across under the survivor's identity. pub fn carry_edit(survivor: Option, source: &Sidecar, survivor_uuid: &str) -> Sidecar { let mut base = survivor.unwrap_or_default(); base.fuse_default_versions(Some(survivor_uuid)); let theirs_default: Vec = source .versions .values() .filter(|v| v.is_default) .map(|v| v.uuid.clone()) .collect(); let mut src = source.clone(); src.fuse_default_versions(Some(survivor_uuid)); if let Some(mut edit) = src.versions.remove(survivor_uuid) { if let Some(ours) = base.versions.get(survivor_uuid) { edit.rating = ours.rating; edit.flag = ours.flag; edit.label = ours.label; edit.revision = edit.revision.max(ours.revision); } edit.revision = edit.revision.saturating_add(1); edit.modified = crate::library::now_secs(); edit.is_default = true; base.put(edit); } let mut snapshots = Vec::new(); for (uuid, v) in src.versions { if base.versions.contains_key(&uuid) { continue; } match &v.snapshot_of { Some(of) if theirs_default.contains(of) => snapshots.push(v), _ => base.put(v), } } base.replace_snapshots(survivor_uuid, snapshots, &[]); if base.derived_from.is_empty() { base.derived_from = source.derived_from.clone(); } if base.merge.is_none() { base.merge = source.merge.clone(); } base } // ── reading ────────────────────────────────────────────────────────────── /// Why a read failed: the server, which may mean it has gone away, or the /// file, which is a verdict about this group alone. #[derive(Debug)] pub enum Fail { Remote(RemoteError), Other(String), } impl Fail { /// Whether every later read would fail the same way. fn offline(&self) -> bool { matches!(self, Fail::Remote(e) if e.indicates_offline()) } } impl std::fmt::Display for Fail { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { Fail::Remote(RemoteError::NotMaterialised(_)) => { f.write_str("not downloaded on this device") } Fail::Remote(e) => write!(f, "{e}"), Fail::Other(s) => f.write_str(s), } } } /// A copy's sidecar: the outbox's copy while it has unsent work, the /// server's otherwise. `Ok(None)` where it has none. async fn read_sidecar( backend: &dyn RemoteBackend, cache: &SidecarCache, image_path: &str, ) -> Result, Fail> { let path = sidecar_path(image_path); if cache.is_pending(&path) { if let Some(s) = cache.load(&path) { return Ok(Some(s)); } } match backend .get(&RemoteId::Path(RemotePath::new(path.clone())), None) .await { Ok(bytes) if bytes.is_empty() => Ok(None), Ok(bytes) => Sidecar::parse(&String::from_utf8_lossy(&bytes)) .map(Some) .map_err(|e| Fail::Other(format!("{path} is unreadable: {e}"))), Err(RemoteError::NotFound(_)) => Ok(None), Err(e) => Err(Fail::Remote(e)), } } /// Read a range, slicing it out where the backend sent the whole file. async fn read_range( backend: &dyn RemoteBackend, path: &str, range: std::ops::Range, size: u64, ) -> Result, Fail> { let want = (range.end - range.start) as usize; let got = backend .get(&RemoteId::Path(RemotePath::new(path)), Some(range.clone())) .await .map_err(Fail::Remote)?; if got.len() == want { return Ok(got); } // `get`'s range is a hint: a backend without range reads answers with // the whole object. if got.len() as u64 == size { return Ok(got[range.start as usize..range.end as usize].to_vec()); } Err(Fail::Other(format!( "{path} is not the size the catalog records" ))) } /// The sameness probe of one file: SHA-256 over its first and last /// [`PROBE_WINDOW`] bytes — the whole file, where it is smaller than two. pub async fn probe(backend: &dyn RemoteBackend, path: &str, size: u64) -> Result { let mut h = Sha256::new(); let head_end = size.min(PROBE_WINDOW); if head_end > 0 { h.update(read_range(backend, path, 0..head_end, size).await?); } let tail_start = size.saturating_sub(PROBE_WINDOW).max(head_end); if tail_start < size { h.update(read_range(backend, path, tail_start..size, size).await?); } Ok(h.finalize().iter().map(|b| format!("{b:02x}")).collect()) } /// Probe a copy where it is, or where an interrupted run left it: the /// trash path is a function of the image id, so a file a stopped run moved /// before its catalog write is found there. async fn probe_copy(backend: &dyn RemoteBackend, root: &str, copy: &Copy) -> Result { match probe(backend, ©.source_ref, copy.file_size).await { Err(Fail::Remote(RemoteError::NotFound(e))) => { let there = trash::trash_path(root, copy.image, ©.source_ref); probe(backend, &there, copy.file_size) .await .map_err(|_| Fail::Other(format!("{} is missing ({e})", copy.source_ref))) } other => other, } } /// Check one group: its bytes, then its edits. /// /// Returns the verdict and the probes it had to take, for the catalog, or /// `Err` where the server went away and no verdict can be given. pub async fn check_group( backend: &dyn RemoteBackend, cache: &SidecarCache, root: &str, group: &Group, ) -> Result<(Check, Vec), String> { let mut probes = Vec::new(); // The bytes. A full digest decides where every copy has one. let digests: Vec = if group.copies.iter().all(|c| c.content_hash.is_some()) { group .copies .iter() .filter_map(|c| c.content_hash.clone()) .collect() } else { let reads = crate::library::futures_join_all(group.copies.iter().map(|c| async move { match &c.probe { Some(p) => Ok((p.clone(), false)), None => probe_copy(backend, root, c).await.map(|d| (d, true)), } })) .await; let mut out = Vec::new(); for (copy, read) in group.copies.iter().zip(reads) { match read { Ok((digest, fresh)) => { if fresh { probes.push(Probe { image: copy.image, file_size: copy.file_size, file_mtime: copy.file_mtime, digest: digest.clone(), }); } out.push(digest); } Err(e) if e.offline() => return Err(e.to_string()), Err(e) => { let name = copy.file_name(); return Ok((Check::Skip(format!("Could not read {name}: {e}")), probes)); } } } out }; if digests.windows(2).any(|w| w[0] != w[1]) { return Ok(( Check::Skip("The files differ in their first or last megabyte".into()), probes, )); } // The edits. let sidecars = crate::library::futures_join_all( group .copies .iter() .map(|c| read_sidecar(backend, cache, &c.source_ref)), ) .await; let mut keys = Vec::new(); for s in sidecars { match s { Ok(s) => keys.push(s.as_ref().and_then(edit_key)), Err(e) if e.offline() => return Err(e.to_string()), Err(e) => { return Ok(( Check::Skip(format!("A sidecar could not be read: {e}")), probes, )) } } } let mut distinct: Vec<&String> = keys.iter().flatten().collect(); distinct.sort(); distinct.dedup(); if distinct.len() > 1 { return Ok((Check::Skip(EDITS_DIFFER.into()), probes)); } Ok(( Check::Same(Edits { edited: keys.iter().map(Option::is_some).collect(), }), probes, )) } /// Check groups on a worker, one at a time, reporting each. /// /// Probes are written back per group, so a check stopped half way keeps /// what it read. pub fn spawn_check( conn: Connection, catalog_path: PathBuf, cache_dir: PathBuf, groups: Vec<(usize, Group)>, stop: Arc, ) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::Network, "dupes", move || { let stopped = run_check(&conn, &catalog_path, cache_dir, groups, &stop, &tx).err(); let _ = tx.send(DupMessage::Finished { stopped }); }); rx } fn run_check( conn: &Connection, catalog_path: &std::path::Path, cache_dir: PathBuf, groups: Vec<(usize, Group)>, stop: &AtomicBool, tx: &Sender, ) -> Result<(), String> { let catalog = Catalog::open(catalog_path).map_err(|e| e.to_string())?; let cache = SidecarCache::open(cache_dir); let rt = crate::net_runtime::build().map_err(|e| e.to_string())?; rt.block_on(async { let backend = crate::remote::connect(conn).map_err(|e| e.to_string())?; let total = groups.len(); for (done, (index, group)) in groups.into_iter().enumerate() { if stop.load(Ordering::Relaxed) { return Err("Stopped".into()); } // The server going away ends the job: every group after this // one would fail the same way, and none of them is a verdict. let (check, probes) = check_group(&*backend, &cache, &conn.account.root, &group).await?; if let Err(e) = dr_catalog::duplicates::record_probes( catalog.connection(), &probes, crate::library::now_secs(), ) { log::warn!("duplicates: keeping probes: {e}"); } if tx.send(DupMessage::Checked { index, check }).is_err() || tx .send(DupMessage::Progress { done: done + 1, total, }) .is_err() { return Err("closed".into()); } } Ok(()) }) } // ── consolidating ──────────────────────────────────────────────────────── /// One group's consolidation, resolved on the UI thread before the worker /// starts. #[derive(Debug, Clone)] pub struct Plan { /// Index in the review, for reporting. pub index: usize, pub survivor: ImageId, pub survivor_path: String, /// The survivor's default version uuid, which a carried edit takes. pub survivor_uuid: String, /// `(image, from, to)`: each other copy and where in the trash it goes. pub copies: Vec<(ImageId, String, String)>, /// Where the survivor has no edit and a copy has the group's one edit: /// that copy's path. pub edit_from: Option, } /// Build the plan for the groups going ahead. pub fn plan( catalog: &Catalog, root: &str, review: &[Reviewed], ) -> Result, dr_catalog::CatalogError> { let mut out = Vec::new(); for (index, r) in review.iter().enumerate() { let Check::Same(edits) = &r.check else { continue; }; if !r.include { continue; } let survivor = &r.group.copies[r.survivor]; let version = dr_catalog::rating::default_version_id(catalog.connection(), survivor.image)?; let survivor_uuid: String = catalog.connection().query_row( "SELECT uuid FROM versions WHERE id = ?1", [version], |row| row.get(0), )?; out.push(Plan { index, survivor: survivor.image, survivor_path: survivor.source_ref.clone(), survivor_uuid, copies: r .group .copies .iter() .enumerate() .filter(|(i, _)| *i != r.survivor) .map(|(_, c)| { ( c.image, c.source_ref.clone(), trash::trash_path(root, c.image, &c.source_ref), ) }) .collect(), edit_from: edits .source_for(r.survivor) .map(|i| r.group.copies[i].source_ref.clone()), }); } Ok(out) } /// Move one file, counting "already there" as done: the run that moved /// it stopped before its catalog write. async fn move_file(backend: &dyn RemoteBackend, from: &str, to: &str) -> Result<(), RemoteError> { match backend .move_to(&RemoteId::Path(RemotePath::new(from)), &RemotePath::new(to)) .await { Ok(()) => Ok(()), Err(RemoteError::NotFound(e)) => { match backend .get(&RemoteId::Path(RemotePath::new(to)), Some(0..1)) .await { Ok(_) => Ok(()), Err(_) => Err(RemoteError::NotFound(e)), } } Err(e) => Err(e), } } /// Put back what a failed group moved, newest first. async fn undo_moves(backend: &dyn RemoteBackend, moved: &[(String, String)]) { for (from, to) in moved.iter().rev() { if let Err(e) = move_file(backend, to, from).await { // Logged loudly: the file is safe in the trash, but the catalog // does not say so, and this line is what says where it is. log::error!("duplicates: {to} could not be moved back to {from}: {e}"); } } } /// Put the survivor's sidecar back as it was before a carried edit, when /// the group it was carried for did not go ahead. async fn restore_sidecar( backend: &dyn RemoteBackend, cache: &SidecarCache, path: &str, before: Option>>, ) { let Some(before) = before else { return }; let remote = RemotePath::new(path); let undone = match before { Some(bytes) => { if let Ok(s) = Sidecar::parse(&String::from_utf8_lossy(&bytes)) { let _ = cache.store(path, &s, false); } backend.put(&remote, bytes, None).await.map(|_| ()) } None => backend.delete(&RemoteId::Path(remote), None).await, }; if let Err(e) = undone { log::warn!("duplicates: putting {path} back: {e}"); } } /// Consolidate one group, whole or not at all. pub async fn consolidate_group( backend: &dyn RemoteBackend, catalog: &Catalog, cache: &SidecarCache, plan: &Plan, ) -> Result { // 1. The edit, where the survivor lacks the group's one edit. let mut sidecar_before: Option>> = None; let survivor_sidecar = sidecar_path(&plan.survivor_path); if let Some(from) = &plan.edit_from { let source = read_sidecar(backend, cache, from) .await .map_err(|e| e.to_string())? .ok_or_else(|| format!("{from} no longer carries its edit"))?; let before = match backend .get( &RemoteId::Path(RemotePath::new(survivor_sidecar.clone())), None, ) .await { Ok(bytes) => Some(bytes), Err(RemoteError::NotFound(_)) => None, Err(e) => return Err(format!("{survivor_sidecar}: {e}")), }; let ours = match &before { Some(bytes) => Some( Sidecar::parse(&String::from_utf8_lossy(bytes)) .map_err(|e| format!("{survivor_sidecar} is unreadable: {e}"))?, ), None => None, }; let carried = carry_edit(ours, &source, &plan.survivor_uuid); backend .put( &RemotePath::new(survivor_sidecar.clone()), carried.to_text().into_bytes(), None, ) .await .map_err(|e| format!("{survivor_sidecar}: {e}"))?; cache.store(&survivor_sidecar, &carried, false)?; sidecar_before = Some(before); } // 2. The files. let mut moved: Vec<(String, String)> = Vec::new(); for (_, from, to) in &plan.copies { match move_file(backend, from, to).await { Ok(()) => moved.push((from.clone(), to.clone())), Err(e) => { undo_moves(backend, &moved).await; restore_sidecar(backend, cache, &survivor_sidecar, sidecar_before.take()).await; let name = from.rsplit('/').next().unwrap_or(from); return Err(format!("moving {name}: {e}")); } } } // 3. The catalog, in one transaction. let recorded: Vec<(ImageId, String)> = plan .copies .iter() .map(|(id, _, to)| (*id, to.clone())) .collect(); match dr_catalog::duplicates::consolidate( catalog.connection(), plan.survivor, &recorded, crate::library::now_secs(), ) { Ok(outcome) => Ok(outcome), Err(e) => { undo_moves(backend, &moved).await; restore_sidecar(backend, cache, &survivor_sidecar, sidecar_before).await; Err(format!("catalog: {e}")) } } } /// Consolidate groups on a worker, one at a time. /// /// `stop` ends the job *between* groups, never inside one. pub fn spawn_consolidate( conn: Connection, catalog_path: PathBuf, cache_dir: PathBuf, plans: Vec, stop: Arc, ) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); executors::spawn(Executor::Network, "consolidate", move || { let stopped = run_consolidate(&conn, &catalog_path, cache_dir, plans, &stop, &tx).err(); let _ = tx.send(DupMessage::Finished { stopped }); }); rx } fn run_consolidate( conn: &Connection, catalog_path: &std::path::Path, cache_dir: PathBuf, plans: Vec, stop: &AtomicBool, tx: &Sender, ) -> Result<(), String> { let catalog = Catalog::open(catalog_path).map_err(|e| e.to_string())?; let cache = SidecarCache::open(cache_dir); let rt = crate::net_runtime::build().map_err(|e| e.to_string())?; rt.block_on(async { let backend = crate::remote::connect(conn).map_err(|e| e.to_string())?; let total = plans.len(); for (done, plan) in plans.iter().enumerate() { // Between groups, never inside one. if stop.load(Ordering::Relaxed) { return Err("Stopped".into()); } let msg = match consolidate_group(&*backend, &catalog, &cache, plan).await { Ok(outcome) => DupMessage::Consolidated { index: plan.index, survivor: plan.survivor, outcome, }, Err(reason) => { log::warn!("duplicates: {}: {reason}", plan.survivor_path); DupMessage::Failed { index: plan.index, reason, } } }; if tx.send(msg).is_err() || tx .send(DupMessage::Progress { done: done + 1, total, }) .is_err() { return Err("closed".into()); } } Ok(()) }) } #[cfg(test)] mod tests { use super::*; fn sidecar(params: &[(&str, &str, f32)], rating: u8) -> Sidecar { let mut s = Sidecar::new(); let mut v = Version { uuid: "a".into(), name: "Default".into(), is_default: true, rating, revision: 3, device: "desk".into(), modified: 99, ..Default::default() }; for (op, p, x) in params { v.params.insert((op.to_string(), p.to_string()), *x); } s.put(v); s } #[test] fn a_judgement_alone_is_not_an_edit() { assert_eq!(edit_key(&sidecar(&[], 4)), None); assert_eq!(edit_key(&Sidecar::new()), None); } #[test] fn the_same_edit_under_other_identities_is_one_edit() { let a = sidecar(&[("exposure", "ev", 0.5)], 1); let mut b = sidecar(&[("exposure", "ev", 0.5)], 5); // Another device's uuid and bookkeeping. let mut v = b.versions.remove("a").unwrap(); v.uuid = "zzz".into(); v.revision = 40; v.device = "tablet".into(); b.put(v); assert!(edit_key(&a).is_some()); assert_eq!(edit_key(&a), edit_key(&b)); } #[test] fn different_parameters_are_different_edits() { let a = sidecar(&[("exposure", "ev", 0.5)], 0); let b = sidecar(&[("exposure", "ev", 0.7)], 0); assert_ne!(edit_key(&a), edit_key(&b)); } #[test] fn a_carried_edit_takes_the_survivors_identity_and_keeps_its_judgement() { let ours = sidecar(&[], 2); let theirs = sidecar(&[("exposure", "ev", 0.5)], 5); let out = carry_edit(Some(ours), &theirs, "survivor-uuid"); let v = out.default_version().unwrap(); assert_eq!(v.uuid, "survivor-uuid"); assert_eq!(v.rating, 2, "the catalog merges judgements, not this"); assert_eq!(edit_key(&out), edit_key(&theirs)); assert_eq!(out.versions.len(), 1); } #[test] fn the_summary_counts_what_the_button_will_do() { let copy = |id: u64| Copy { image: ImageId(id), source_ref: format!("a/{id}.CR2"), file_size: 1, file_mtime: None, added_at: 0, content_hash: None, probe: None, file_id: None, }; let group = |ids: &[u64]| Group { root_id: 1, camera: "c".into(), captured_at: 1, file_size: 1, copies: ids.iter().map(|i| copy(*i)).collect(), }; let mut review = vec![ Reviewed::new(group(&[1, 2, 3])), Reviewed::new(group(&[4, 5])), Reviewed::new(group(&[6, 7])), Reviewed::new(group(&[8, 9])), ]; review[0].check = Check::Same(Edits::default()); review[1].check = Check::Skip(EDITS_DIFFER.into()); review[2].check = Check::Same(Edits::default()); review[2].include = false; let s = summarise(&review); assert_eq!( s, Summary { groups: 1, files: 2, skipped: 1, excluded: 1, unchecked: 1 } ); } #[test] fn an_edit_is_taken_only_where_the_survivor_lacks_it() { let e = Edits { edited: vec![false, true, true], }; assert_eq!(e.source_for(0), Some(1)); assert_eq!(e.source_for(2), None); let none = Edits { edited: vec![false, false], }; assert_eq!(none.source_for(0), None); } // ── end to end, on a folder library ──────────────────────────────── /// A folder library holding real files, and a catalog describing them. struct Library { dir: PathBuf, lib: PathBuf, catalog_path: PathBuf, cache: PathBuf, } impl Library { fn new(name: &str) -> Self { let dir = std::env::temp_dir().join(format!("dr-ui-dups-{name}-{}", std::process::id())); let _ = std::fs::remove_dir_all(&dir); let lib = dir.join("lib"); std::fs::create_dir_all(&lib).unwrap(); let catalog_path = dir.join("catalog.sqlite"); let cat = Catalog::open(&catalog_path).unwrap(); cat.connection() .execute( "INSERT INTO roots(id, kind, label) VALUES (1, 'local', ?1)", [lib.to_string_lossy().as_ref()], ) .unwrap(); Self { cache: dir.join("sidecars"), dir, lib, catalog_path, } } fn conn(&self) -> Connection { Connection::new( dr_sync::Account::new("folder", self.lib.to_string_lossy()), None, ) } fn catalog(&self) -> Catalog { Catalog::open(&self.catalog_path).unwrap() } /// Write a file and catalogue it. fn add(&self, id: i64, path: &str, at: i64, bytes: &[u8]) { let full = self.lib.join(path); std::fs::create_dir_all(full.parent().unwrap()).unwrap(); std::fs::write(&full, bytes).unwrap(); let cat = self.catalog(); cat.connection() .execute( "INSERT INTO images(id, root_id, source_ref, captured_at, camera, file_size, file_mtime, added_at) VALUES (?1, 1, ?2, ?3, 'Canon EOS 6D', ?4, 1, ?1)", rusqlite::params![id, path, at, bytes.len() as i64], ) .unwrap(); dr_catalog::rating::ensure_default_versions(cat.connection()).unwrap(); } fn edit(&self, path: &str, ev: f32) { let text = sidecar(&[("exposure", "ev", ev)], 0).to_text(); std::fs::write(self.lib.join(sidecar_path(path)), text).unwrap(); } fn exists(&self, path: &str) -> bool { self.lib.join(path).exists() } fn review(&self) -> Vec { let cat = self.catalog(); dr_catalog::duplicates::candidates(cat.connection()) .unwrap() .into_iter() .map(Reviewed::new) .collect() } fn check(&self, review: &mut [Reviewed]) { let groups = review .iter() .enumerate() .map(|(i, r)| (i, r.group.clone())) .collect(); let rx = spawn_check( self.conn(), self.catalog_path.clone(), self.cache.clone(), groups, Arc::default(), ); for msg in rx { match msg { DupMessage::Checked { index, check } => review[index].check = check, DupMessage::Finished { stopped } => assert_eq!(stopped, None), _ => {} } } } fn consolidate(&self, review: &[Reviewed]) -> (usize, Vec) { let plans = plan(&self.catalog(), "", review).unwrap(); let rx = spawn_consolidate( self.conn(), self.catalog_path.clone(), self.cache.clone(), plans, Arc::default(), ); let (mut done, mut failed) = (0, Vec::new()); for msg in rx { match msg { DupMessage::Consolidated { .. } => done += 1, DupMessage::Failed { reason, .. } => failed.push(reason), _ => {} } } (done, failed) } } impl Drop for Library { fn drop(&mut self) { let _ = std::fs::remove_dir_all(&self.dir); } } /// Three megabytes that differ from any other seed's, so the probe reads /// a real head and a real tail. fn raw(seed: u8) -> Vec { (0..3 * 1024 * 1024u32) .map(|i| (i.wrapping_mul(2_654_435_761) >> 13) as u8 ^ seed) .collect() } fn trashed(cat: &Catalog) -> i64 { cat.connection() .query_row( "SELECT COUNT(*) FROM images WHERE trashed_at IS NOT NULL", [], |r| r.get(0), ) .unwrap() } #[test] fn a_folder_library_is_checked_consolidated_and_restored() { let lib = Library::new("e2e"); // The same frame three times, with the one edit on a backup copy. let frame = raw(1); lib.add(1, "2023/bck/_MG_4623.CR2", 100, &frame); lib.add(2, "2023/2023-06-24/_MG_4623.CR2", 100, &frame); lib.add(3, "alps trip/Raw/20230628_0642.CR2", 100, &frame); lib.edit("2023/bck/_MG_4623.CR2", 0.5); // Same size, camera and instant; a different last byte. let mut other = raw(2); lib.add(4, "2023/2023-06-24/_MG_4700.CR2", 200, &other); *other.last_mut().unwrap() ^= 1; lib.add(5, "2023/bck/_MG_4700.CR2", 200, &other); // Identical bytes, two different edits. let edited = raw(3); lib.add(6, "2023/2023-06-24/_MG_4800.CR2", 300, &edited); lib.add(7, "2023/bck/_MG_4800.CR2", 300, &edited); lib.edit("2023/2023-06-24/_MG_4800.CR2", 0.3); lib.edit("2023/bck/_MG_4800.CR2", -0.3); let mut review = lib.review(); assert_eq!(review.len(), 3); lib.check(&mut review); assert!( matches!(review[0].check, Check::Same(_)), "{:?}", review[0].check ); assert_eq!( review[1].check, Check::Skip("The files differ in their first or last megabyte".into()) ); assert_eq!(review[2].check, Check::Skip(EDITS_DIFFER.into())); assert_eq!( review[0].group.copies[review[0].survivor].source_ref, "2023/2023-06-24/_MG_4623.CR2" ); assert_eq!( summarise(&review), Summary { groups: 1, files: 2, skipped: 2, excluded: 0, unchecked: 0 } ); // What was read is kept: a second review needs no bytes. assert!(lib.review()[0] .group .copies .iter() .all(|c| c.probe.is_some())); let (done, failed) = lib.consolidate(&review); assert_eq!((done, failed), (1, Vec::::new())); // The survivor stays; the copies are in the trash, not deleted. assert!(lib.exists("2023/2023-06-24/_MG_4623.CR2")); assert!(!lib.exists("2023/bck/_MG_4623.CR2")); assert!(!lib.exists("alps trip/Raw/20230628_0642.CR2")); assert!(lib.exists(".darkroom-trash/1-_MG_4623.CR2")); assert!(lib.exists(".darkroom-trash/3-20230628_0642.CR2")); // The skipped groups are untouched. for path in [ "2023/bck/_MG_4700.CR2", "2023/bck/_MG_4800.CR2", "2023/2023-06-24/_MG_4800.CR2", ] { assert!(lib.exists(path), "{path}"); } let cat = lib.catalog(); assert_eq!(trashed(&cat), 2); // The edit the backup copy carried now lives beside the survivor. let carried = std::fs::read_to_string(lib.lib.join(sidecar_path("2023/2023-06-24/_MG_4623.CR2"))) .unwrap(); let carried = Sidecar::parse(&carried).unwrap(); assert_eq!( edit_key(&carried), edit_key(&sidecar(&[("exposure", "ev", 0.5)], 0)) ); // And the trash gives them back. let moves = crate::trash::plan_restore(&cat, &[ImageId(1), ImageId(3)]).unwrap(); let rx = crate::trash::spawn_move( lib.conn(), moves, crate::trash::Direction::Restore, lib.catalog_path.clone(), ); let mut restored = 0; for msg in rx { if let crate::trash::TrashMessage::Done { moved, failed } = msg { assert!(failed.is_empty(), "{failed:?}"); restored = moved; } } assert_eq!(restored, 2); assert!(lib.exists("2023/bck/_MG_4623.CR2")); assert!(lib.exists("alps trip/Raw/20230628_0642.CR2")); assert_eq!( std::fs::read(lib.lib.join("2023/bck/_MG_4623.CR2")).unwrap(), frame ); assert_eq!(trashed(&lib.catalog()), 0); } #[test] fn a_catalog_failure_moves_the_files_back() { let lib = Library::new("rollback"); let frame = raw(4); lib.add(1, "2023/bck/IMG_0001.CR2", 100, &frame); lib.add(2, "2023/a/IMG_0001.CR2", 100, &frame); lib.add(3, "2023/b/IMG_0001-2.CR2", 100, &frame); let mut review = lib.review(); lib.check(&mut review); assert!(matches!(review[0].check, Check::Same(_))); // The trash write for the second copy fails, after the first copy's // judgements have been merged in the same transaction. lib.catalog() .connection() .execute_batch( "CREATE TRIGGER fail_on_3 BEFORE UPDATE OF source_ref ON images WHEN NEW.id = 3 BEGIN SELECT RAISE(ABORT, 'disk full'); END;", ) .unwrap(); let (done, failed) = lib.consolidate(&review); assert_eq!(done, 0); assert_eq!(failed.len(), 1, "{failed:?}"); for path in [ "2023/bck/IMG_0001.CR2", "2023/a/IMG_0001.CR2", "2023/b/IMG_0001-2.CR2", ] { assert!(lib.exists(path), "{path} is back where it was"); } assert!(!lib.exists(".darkroom-trash/1-IMG_0001.CR2")); assert_eq!(trashed(&lib.catalog()), 0); } #[test] fn a_run_that_stopped_after_its_moves_is_finished_by_the_next() { let lib = Library::new("resume"); let frame = raw(5); lib.add(1, "2023/a/IMG_0001.CR2", 100, &frame); lib.add(2, "2023/bck/IMG_0001.CR2", 100, &frame); // A previous run moved the backup copy, then the process died // before the catalog heard about it. std::fs::create_dir_all(lib.lib.join(".darkroom-trash")).unwrap(); std::fs::rename( lib.lib.join("2023/bck/IMG_0001.CR2"), lib.lib.join(".darkroom-trash/2-IMG_0001.CR2"), ) .unwrap(); let mut review = lib.review(); lib.check(&mut review); assert!( matches!(review[0].check, Check::Same(_)), "{:?}", review[0].check ); let (done, failed) = lib.consolidate(&review); assert_eq!((done, failed), (1, Vec::::new())); assert_eq!(trashed(&lib.catalog()), 1); assert!(lib.exists(".darkroom-trash/2-IMG_0001.CR2")); } }