From 8d4ecb75c12c43dcaa175da100d72c2ba3c2a035 Mon Sep 17 00:00:00 2001 From: Duncan Tourolle Date: Fri, 25 Sep 2026 23:07:45 -0400 Subject: [PATCH] Prove duplicate originals the same and consolidate them on workers dr_ui::duplicates is the half of #67 that touches files. The check reads each copy's first and last megabyte by range through the backend and hashes them (or compares stored content hashes where every copy has one), keeps the probes in the catalog, and reads each copy's sidecar: a group whose bytes differ, whose copies cannot be read, or whose develop edits differ is left out and the review says why. Consolidating a group carries the one edit onto the survivor's sidecar where it has none, moves the other copies into the trash, and then commits dr_catalog::duplicates::consolidate. A failure after the first move puts the files and the sidecar back; a run that died between the moves and the commit is finished by the next one, which finds each moved file at its trash path. Tested end to end on a folder library of real files: the copies land in .darkroom-trash, the skipped groups are untouched, the edit reaches the survivor, a catalog failure moves everything back, and restore returns the copies byte for byte. --- Cargo.lock | 1 + ui/dr-ui/Cargo.toml | 3 + ui/dr-ui/src/duplicates.rs | 1215 ++++++++++++++++++++++++++++++++++++ ui/dr-ui/src/lib.rs | 1 + 4 files changed, 1220 insertions(+) create mode 100644 ui/dr-ui/src/duplicates.rs 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;