diff --git a/Cargo.lock b/Cargo.lock index 962825b..bc52970 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1714,6 +1714,7 @@ dependencies = [ "rusqlite", "serde_json", "serde_norway", + "sha2", "slint", "slint-build", "thiserror 2.0.20", diff --git a/ui/dr-ui/Cargo.toml b/ui/dr-ui/Cargo.toml index 0e39603..29a6b35 100644 --- a/ui/dr-ui/Cargo.toml +++ b/ui/dr-ui/Cargo.toml @@ -37,6 +37,9 @@ dr-export.workspace = true # runtime is the same tract the faces and masks already carry. dr-pano = { workspace = true, features = ["xfeat", "embedded-model"] } dr-ingest.workspace = true +# The sameness probe of a catalog duplicate (FR-CAT-11a): SHA-256 over the +# ends of each copy, the digest the import already uses for whole files. +sha2 = "0.10" dr-film.workspace = true # The lens profile database, here for the same reason dr-film is: dr-pipeline # knows the maths of lens correction and deliberately has no dependency with diff --git a/ui/dr-ui/src/duplicates.rs b/ui/dr-ui/src/duplicates.rs new file mode 100644 index 0000000..9506d8b --- /dev/null +++ b/ui/dr-ui/src/duplicates.rs @@ -0,0 +1,1215 @@ +//! 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 std::path::PathBuf; +use std::sync::mpsc::{Receiver, Sender}; + +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, + } + } + + /// Whether the group goes ahead when the button is pressed. + pub fn planned(&self) -> bool { + self.include && matches!(self.check, Check::Same(_)) + } +} + +/// 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)>, +) -> Receiver { + let (tx, rx) = std::sync::mpsc::channel(); + std::thread::spawn(move || { + let stopped = run_check(&conn, &catalog_path, cache_dir, groups, &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)>, + 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() { + // 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. +/// +/// Closing the review stops the job *between* groups, never inside one. +pub fn spawn_consolidate( + conn: Connection, + catalog_path: PathBuf, + cache_dir: PathBuf, + plans: Vec, +) -> Receiver { + let (tx, rx) = std::sync::mpsc::channel(); + std::thread::spawn(move || { + let stopped = run_consolidate(&conn, &catalog_path, cache_dir, plans, &tx).err(); + let _ = tx.send(DupMessage::Finished { stopped }); + }); + rx +} + +fn run_consolidate( + conn: &Connection, + catalog_path: &std::path::Path, + cache_dir: PathBuf, + plans: Vec, + 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() { + 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, + ); + 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, + ); + 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")); + } +} diff --git a/ui/dr-ui/src/lib.rs b/ui/dr-ui/src/lib.rs index e3940c8..e030212 100644 --- a/ui/dr-ui/src/lib.rs +++ b/ui/dr-ui/src/lib.rs @@ -26,6 +26,7 @@ mod bursts; mod collections_ui; #[cfg(test)] mod decoder_seam; +mod duplicates; mod derived_sync; mod develop; mod develop_ui;