Files
DarkRoom/ui/dr-ui/src/duplicates.rs
T
dtourolle b1d1c47261 Start every worker thread through the executors module
Thirty-nine spawn sites in dr-ui, and one in the Android entry point,
called std::thread::spawn or a Builder of their own, and most of the
threads they started were <unnamed> in a panic message or a profiler.
Each now calls executors::spawn with its executor and a role, so the thread is
named <executor>:<role> — net:sync, decode:thumbs, io:catalog-open —
and knows which executor it is on. The three that already set a name
(automation, import, prefetch) keep their name as the role.

Behaviour is unchanged: each job still gets a thread of its own when it
starts, and spawn panics where std::thread::spawn did.

The module's documentation now says how a job is assigned: by what it
spends its time on, so a sweep that fetches bytes and then decodes them
is Decode, and a sidecar write that touches the catalog is Network.

Left as they were: the segmentation and refine workers in masks_ui.rs,
which another change is reworking, and test-only threads.
2026-09-27 07:08:37 -04:00

1255 lines
42 KiB
Rust

//! 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<bool>,
}
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<usize> {
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<String>,
},
}
// ── 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<String> {
let mut sc = sidecar.clone();
sc.fuse_default_versions(Some("default"));
let mut blocks: Vec<String> = 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<Sidecar>, source: &Sidecar, survivor_uuid: &str) -> Sidecar {
let mut base = survivor.unwrap_or_default();
base.fuse_default_versions(Some(survivor_uuid));
let theirs_default: Vec<String> = 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<Option<Sidecar>, 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<u64>,
size: u64,
) -> Result<Vec<u8>, 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<String, Fail> {
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<String, Fail> {
match probe(backend, &copy.source_ref, copy.file_size).await {
Err(Fail::Remote(RemoteError::NotFound(e))) => {
let there = trash::trash_path(root, copy.image, &copy.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<Probe>), String> {
let mut probes = Vec::new();
// The bytes. A full digest decides where every copy has one.
let digests: Vec<String> = 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<AtomicBool>,
) -> Receiver<DupMessage> {
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<DupMessage>,
) -> 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<String>,
}
/// Build the plan for the groups going ahead.
pub fn plan(
catalog: &Catalog,
root: &str,
review: &[Reviewed],
) -> Result<Vec<Plan>, 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<Option<Vec<u8>>>,
) {
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<Outcome, String> {
// 1. The edit, where the survivor lacks the group's one edit.
let mut sidecar_before: Option<Option<Vec<u8>>> = 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<Plan>,
stop: Arc<AtomicBool>,
) -> Receiver<DupMessage> {
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<Plan>,
stop: &AtomicBool,
tx: &Sender<DupMessage>,
) -> 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<Reviewed> {
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<String>) {
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<u8> {
(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::<String>::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::<String>::new()));
assert_eq!(trashed(&lib.catalog()), 1);
assert!(lib.exists(".darkroom-trash/2-IMG_0001.CR2"));
}
}