//! TRACES: FR-CAT-10 | FR-CAT-11 | FR-NC-7a //! Bringing a card into the library: the work, off the interface thread. //! //! [`dr_ingest`] does the copying and knows nothing about cards, catalogs or //! windows — it is handed a list of files, a probe for their metadata and a //! question about duplicates. This supplies all three, and runs the result on //! a worker. //! //! The shape is the one every background operation in this crate has: a //! thread with an `mpsc` channel, drained by a `slint::Timer` on the UI //! thread, reporting into [`crate::activity`]. See `import_ui` for the //! draining half. //! //! # Two phases, one worker //! //! Surveying a card is itself slow — a full card is two thousand files across //! a directory tree on a bus that manages 40 MB/s on a good day — so it runs //! on the worker too and reports what it found before any copying starts. The //! user sees "1,847 photographs, 61.2 GB" and then a progress bar, rather than //! a frozen window and then a progress bar. //! //! # The import does not catalogue what it wrote //! //! It writes files into the library and then asks for a scan, rather than //! inserting rows itself. FR-CAT-1 already turns files on disk into catalog //! rows, and a second path into the `images` table would be a second set of //! bugs about folders, formats and metadata state. What the import *does* //! record afterwards is each file's digest (`dr_catalog::dedup`), which the //! scan cannot know because it never reads a whole file. use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::mpsc::{Receiver, Sender}; use std::sync::Arc; use dr_ingest::{Candidate, DupKey, Imported, Ingest, Options, Report, Shot}; use dr_plat::{DirRef, LocalStorage, Storage, WritableStorage}; use dr_sync_nextcloud::{AppCredentials, NextcloudBackend}; use dr_types::{FormatFilter, RootId}; /// Which root the card is granted as, and which the library is. /// /// Two distinct ids in one `LocalStorage` would also work; separate storages /// keep the read-only source and the writable destination separate all the /// way down, which is the distinction `WritableStorage` exists to make. const CARD: RootId = RootId(9001); const LIBRARY: RootId = RootId(9002); const BACKUP: RootId = RootId(9003); /// A cancel flag shared with the worker. pub type Cancel = Arc; /// Everything an import needs to run. #[derive(Debug, Clone)] pub struct Request { /// Where the card is mounted. pub card: PathBuf, /// The library root the template's folders are created under. pub library: PathBuf, /// The optional second destination (FR-CAT-10). pub backup: Option, /// The catalog, opened by the worker for its duplicate queries. /// /// A path rather than a connection: `rusqlite::Connection` is not `Sync`, /// and every other worker in this crate opens its own for the same reason. pub catalog: PathBuf, /// Which file types to take off the card. pub filter: FormatFilter, pub options: Options, /// Where to send the originals afterwards, if anywhere (FR-NC-7b). pub upload: Option, } /// TRACES: FR-NC-7a | FR-NC-7b /// Sending the imported originals on to the server. /// /// Deliberately a second phase rather than a destination the copy writes /// straight to. FR-NC-7b: files are copied locally and verified *first*, and /// only then queued for upload — so a network that fails costs an upload, not /// an import, and the photographs exist on disk either way. #[derive(Clone)] pub struct Upload { pub credentials: AppCredentials, pub user_id: String, /// The library folder on the server. The dated folders from the template /// are created beneath it, the same ones the local copy went into. pub library: String, } impl std::fmt::Debug for Upload { /// Hand-written so a credential cannot reach a log through a `{:?}`. fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("Upload") .field("user_id", &self.user_id) .field("library", &self.library) .finish_non_exhaustive() } } /// What the worker sends back. #[derive(Debug, Clone)] pub enum Message { /// The card has been walked. Sent once, before any copying. Surveyed { files: usize, bytes: u64 }, /// One file finished — imported, skipped or failed. Progress { done: usize, total: usize, bytes: u64, name: String, }, /// Originals going up, after every one of them is safely on disk. Uploading { done: usize, total: usize }, /// The run ended. Always the last message. Finished(Outcome), /// The run could not start at all. Failed(String), } /// What an import did, in the form the interface reports it. #[derive(Debug, Clone, Default, PartialEq, Eq)] pub struct Outcome { pub imported: usize, pub duplicates: usize, pub failed: usize, /// Files dated by modification time rather than by the shutter /// (FR-NC-7a). Named rather than counted: the user needs to know *which* /// photographs are filed under a date that is not when they were taken. pub undated: Vec, pub bytes: u64, pub cancelled: bool, /// The folders that were written into, for the message at the end and for /// the upload that follows (FR-NC-7a). pub folders: Vec, /// Card files that may now be deleted, for a move-import (FR-NC-7b). /// /// Only files that are *both* verified on disk and, where an upload was /// asked for, confirmed on the server. A card erased against a failed /// upload is unrecoverable, so this is the one count computed /// pessimistically. pub retirable: usize, /// Originals confirmed on the server. pub uploaded: usize, /// Originals that stayed local because the upload did not go through. /// /// Not a failure of the import: the photographs are on disk and verified, /// and the card has not been touched. Reported so the user knows the /// server does not have them yet. pub upload_failed: usize, } /// Start an import. /// /// Returns immediately; the work happens on a thread and reports through the /// channel. A dropped receiver does not stop the worker — the cancel flag /// does, and the caller holds it. pub fn spawn(request: Request, cancel: Cancel) -> Receiver { let (tx, rx) = std::sync::mpsc::channel(); std::thread::Builder::new() .name("import".into()) .spawn(move || { if let Err(e) = run(request, &cancel, &tx) { let _ = tx.send(Message::Failed(e)); } }) .expect("spawning the import worker"); rx } fn run(request: Request, cancel: &Cancel, tx: &Sender) -> Result<(), String> { let card = LocalStorage::with_root(CARD, request.card.clone()); let library = LocalStorage::with_root(LIBRARY, request.library.clone()); let backup = request .backup .as_ref() .map(|p| LocalStorage::with_root(BACKUP, p.clone())); let candidates = survey(&card, CARD, &request.filter, &|| { cancel.load(Ordering::Relaxed) }) .map_err(|e| format!("could not read the card: {e}"))?; let bytes: u64 = candidates.iter().map(|c| c.size).sum(); let _ = tx.send(Message::Surveyed { files: candidates.len(), bytes, }); // Opened after the survey rather than before: a card that turns out to be // empty should not also report a catalog it never needed. let catalog = dr_catalog::Catalog::open(&request.catalog) .map_err(|e| format!("could not open the catalog: {e}"))?; let library_root = DirRef::root(LIBRARY); let backup_root = DirRef::root(BACKUP); let ingest = Ingest { source: &card, dest: &library, dest_root: &library_root, backup: backup .as_ref() .map(|b| (b as &dyn WritableStorage, &backup_root)), options: request.options.clone(), }; let report = ingest.run( &candidates, &|c| probe(&card, c), &|key| is_duplicate(catalog.connection(), key), &|| cancel.load(Ordering::Relaxed), &mut |p| { let _ = tx.send(Message::Progress { done: p.done, total: p.total, bytes: p.bytes, name: p.name.clone(), }); }, ); let mut outcome = summarise(&report); // FR-NC-7b. Everything above has already finished: every file is on disk // and its digest checked, so nothing from here on can cost the user a // photograph — only a transfer. if let Some(upload) = &request.upload { let sent = upload_all(upload, &library, &report.imported, cancel, tx); outcome.uploaded = sent.len(); outcome.upload_failed = report.imported.len() - sent.len(); // The card keeps its copies of anything the server did not take. outcome.retirable = outcome.retirable.min(outcome.uploaded); // The digests, so the *next* import can answer the content tier // without reading anything. // // Recorded against the **remote** path rather than the local one, // because the catalog holds the server's library: the local copy has // no row and never will. Best-effort and after the fact — the row does // not exist until a sync has seen the upload, and a hash that misses // its row costs one wasted transfer next time rather than a failure // now. The metadata tier covers the interval, which is why this is not // worth waiting for a sync to do properly. record_digests(catalog.connection(), &upload.library, &sent); } let _ = tx.send(Message::Finished(outcome)); Ok(()) } /// Send the imported originals to the server, in their dated folders. /// /// Returns the remote path and digest of each original the server confirmed. /// A failure here is reported /// and not propagated: the import succeeded, and turning "the network was /// slow" into a failed import would misdescribe what happened to the /// photographs and, on a move-import, would be the difference between a card /// kept and a card emptied. fn upload_all( upload: &Upload, library: &LocalStorage, imported: &[Imported], cancel: &Cancel, tx: &Sender, ) -> Vec<(String, String)> { let rt = match crate::net_runtime::build() { Ok(rt) => rt, Err(e) => { log::warn!("no runtime for the upload: {e}"); return Vec::new(); } }; rt.block_on(async { let backend = match NextcloudBackend::new(&upload.credentials, &upload.user_id) { Ok(b) => b, Err(e) => { log::warn!("connecting to upload: {e}"); return Vec::new(); } }; let root = dr_sync::RemotePath::new(&upload.library); let mut sent = Vec::new(); for (i, image) in imported.iter().enumerate() { if cancel.load(Ordering::Relaxed) { // Everything not yet sent stays local, and the card keeps its // copies of all of it. break; } let _ = tx.send(Message::Uploading { done: i, total: imported.len(), }); // Read from the library rather than the card: this is the copy // that was verified, and the card may already be unplugged. let bytes = match read_all(library, image) { Ok(b) => b, Err(e) => { log::warn!("reading {} to upload it: {e}", image.name); continue; } }; // The same folder segments the local copy went into, so the two // libraries have the same shape (FR-NC-7a). match dr_sync::upload_original(&backend, &root, &image.folders, &image.name, bytes) .await { Ok((path, _)) => { log::info!("uploaded {path}"); sent.push((path.as_str().to_string(), image.digest.clone())); } Err(e) => log::warn!("uploading {}: {e}", image.name), } } sent }) } /// Read an imported file back out of the library. fn read_all(library: &LocalStorage, image: &Imported) -> Result, String> { use std::io::Read; let mut stream = library.open(&image.written).map_err(|e| e.to_string())?; let mut bytes = Vec::with_capacity(image.size as usize); stream.read_to_end(&mut bytes).map_err(|e| e.to_string())?; Ok(bytes) } /// Walk a card for files worth importing. /// /// Recursive rather than `DCIM`-only: a card carries `DCIM/100CANON`, a phone /// carries `DCIM/Camera`, and a card that has been through a computer carries /// whatever somebody put on it. Walking the whole volume finds all three, and /// the format filter is what keeps the result to photographs. pub fn survey( storage: &dyn Storage, root: RootId, filter: &FormatFilter, cancel: &dyn Fn() -> bool, ) -> Result, dr_plat::StorageError> { let mut out = Vec::new(); let mut queue = vec![storage.root_dir(root)?]; while let Some(dir) = queue.pop() { if cancel() { break; } // A directory that cannot be read is skipped rather than fatal: cards // carry vendor folders with odd permissions, and one of them must not // cost the user the other nine hundred photographs (FR-CAT-9). let entries = match storage.list(&dir) { Ok(e) => e, Err(e) => { log::warn!("skipping {dir}: {e}"); continue; } }; for entry in entries { if entry.meta.is_dir { if !dr_catalog::walk::is_excluded(&entry.meta.name) { if let dr_plat::Node::Dir(d) = entry.node { queue.push(d); } } continue; } if !filter.allows_name(&entry.meta.name) { continue; } if let dr_plat::Node::File(source) = entry.node { out.push(Candidate { source, name: entry.meta.name, size: entry.meta.size, // The listing's reading, which is the fallback the folder // template uses when a file carries no capture time. modified_at: entry.meta.mtime / 1000, }); } } } // A card walks in whatever order the filesystem hands back, which for FAT // is creation order per directory and arbitrary across them. Sorted so an // import reports its progress through a card in the order a photographer // would expect, and so two runs of the same card behave the same way. out.sort_by(|a, b| a.name.cmp(&b.name)); Ok(out) } /// Read one candidate's capture metadata. /// /// A header read rather than the whole file: `dr_decode::HEADER_BYTES` is /// enough for EXIF on every format in the tree, and reading 80 MB per file to /// learn a date would take longer than the import itself. fn probe(card: &dyn Storage, candidate: &Candidate) -> Option { let header = card .read_range(&candidate.source, 0..dr_decode::HEADER_BYTES) .ok()?; let md = dr_decode::metadata(&header).ok()?; Some(Shot { captured_at: md.captured_at, captured_offset: md.captured_offset, make: md.make, model: md.model, modified_at: candidate.modified_at, }) } /// FR-CAT-11's two tiers, against the catalog. /// /// The camera string is composed the way the scan composes it /// ([`crate::library::camera_label`]) — the comparison is against what a scan /// wrote, so a second spelling here would silently disable the cheap tier. fn is_duplicate(conn: &rusqlite::Connection, key: &DupKey) -> bool { if let Some(digest) = &key.digest { return dr_catalog::seen_by_content(conn, digest).unwrap_or(false); } let camera = crate::library::camera_label(key.make.as_deref(), key.model.as_deref()); dr_catalog::seen_by_metadata( conn, key.captured_at, camera.as_deref(), key.size, &key.original_name, ) .unwrap_or(false) } /// Store what the import learned that a scan cannot. /// /// A scan never reads a whole file, so `content_hash` is only ever filled in /// by something that had a reason to read every byte — which an import did. fn record_digests(conn: &rusqlite::Connection, library: &str, sent: &[(String, String)]) { // The library's row, resolved the same way the remote scan resolves it: // one root per library folder, keyed by label. let root: Option = conn .query_row( "SELECT id FROM roots WHERE label = ?1 AND kind = 'remote'", [library], |r| r.get(0), ) .ok(); let Some(root) = root else { // No root yet means no sync has run against this library, so there is // nothing to attach a digest to. Not an error. return; }; for (path, digest) in sent { if let Err(e) = dr_catalog::set_content_hash(conn, root as u64, path, digest) { log::warn!("recording the digest of {path}: {e}"); } } } /// Reduce a report to what the interface says about it. pub fn summarise(report: &Report) -> Outcome { let mut folders: Vec = report .imported .iter() .map(|i| i.folders.join("/")) .collect(); folders.sort(); folders.dedup(); Outcome { imported: report.imported.len(), duplicates: report.duplicates.len(), failed: report.failed.len(), undated: report.undated.clone(), bytes: report.bytes, cancelled: report.cancelled, folders, retirable: report.retirable.len(), // Filled in by the upload phase, which runs after this. uploaded: 0, upload_failed: 0, } } /// The sentence shown when an import ends. /// /// Says what happened to every file, because the counts that are zero are the /// ones worth not mentioning and the ones that are not are the whole point. /// A run that imported nothing says so rather than reporting success in the /// abstract. pub fn describe(outcome: &Outcome) -> String { let mut parts = Vec::new(); if outcome.imported > 0 { parts.push(format!( "{} imported ({})", outcome.imported, crate::activity::describe_bytes(outcome.bytes) )); } if outcome.duplicates > 0 { parts.push(format!("{} already in the library", outcome.duplicates)); } if outcome.failed > 0 { parts.push(format!("{} failed", outcome.failed)); } let head = if parts.is_empty() { "Nothing to import".to_string() } else { parts.join(", ") }; let mut out = if outcome.cancelled { format!("Stopped — {head}") } else { head }; // Where they went, which is the question the folder template raises. match outcome.folders.len() { 0 => {} 1 => out.push_str(&format!(" into {}", outcome.folders[0])), n => out.push_str(&format!(" into {n} folders")), } // What the server got. Said plainly rather than folded into the import // count: "imported" and "uploaded" are different promises, and a user on a // failing connection needs to know which one held. if outcome.uploaded > 0 || outcome.upload_failed > 0 { out.push_str(&format!(". {} uploaded", outcome.uploaded)); if outcome.upload_failed > 0 { out.push_str(&format!( ", {} still only on this computer", outcome.upload_failed )); } } // FR-NC-7a: the mtime fallback is reported rather than silent. if !outcome.undated.is_empty() { out.push_str(&format!( ". {} had no capture time and {} filed by modification date", outcome.undated.len(), if outcome.undated.len() == 1 { "was" } else { "were" } )); } out } /// Whether a folder looks like somewhere a card would be. /// /// Used to warn before an import writes into the wrong place: a library root /// and a card mount point are both just directories, and the two are very /// easy to swap in a dialog. pub fn looks_like_a_card(path: &Path) -> bool { path.join("DCIM").is_dir() } #[cfg(test)] mod tests { use super::*; use dr_ingest::DateSource; struct Tree(PathBuf); impl Tree { fn new(name: &str) -> Self { let dir = std::env::temp_dir().join(format!( "dr-ui-import-{name}-{}-{:?}", std::process::id(), std::thread::current().id() )); let _ = std::fs::remove_dir_all(&dir); std::fs::create_dir_all(&dir).unwrap(); Tree(dir) } fn file(&self, rel: &str, bytes: &[u8]) -> &Self { let p = self.0.join(rel); std::fs::create_dir_all(p.parent().unwrap()).unwrap(); std::fs::write(p, bytes).unwrap(); self } } impl Drop for Tree { fn drop(&mut self) { let _ = std::fs::remove_dir_all(&self.0); } } fn survey_of(t: &Tree) -> Vec { let storage = LocalStorage::with_root(CARD, t.0.clone()); survey(&storage, CARD, &FormatFilter::all(), &|| false).unwrap() } #[test] fn a_survey_finds_photographs_wherever_the_camera_put_them() { let t = Tree::new("survey"); // Three real layouts at once: a Canon card, a phone, and files // somebody dropped at the top level. t.file("DCIM/100CANON/IMG_0001.CR3", b"a") .file("DCIM/Camera/PXL_0002.dng", b"bb") .file("loose.NEF", b"ccc"); let found = survey_of(&t); assert_eq!(found.len(), 3); assert_eq!( found.iter().map(|c| c.name.as_str()).collect::>(), ["IMG_0001.CR3", "PXL_0002.dng", "loose.NEF"] ); } #[test] fn a_survey_leaves_what_is_not_a_photograph() { let t = Tree::new("survey-filter"); t.file("DCIM/100CANON/IMG_0001.CR3", b"a") // Every card carries these, and none of them is an import. .file("DCIM/100CANON/IMG_0001.THM", b"x") .file("MISC/settings.dat", b"x") .file("readme.txt", b"x"); let found = survey_of(&t); assert_eq!(found.len(), 1); assert_eq!(found[0].name, "IMG_0001.CR3"); } #[test] fn a_survey_is_ordered_the_same_way_twice() { // FAT hands directories back in creation order and volumes in none at // all, so an unsorted survey reports its progress through a card in an // order that changes between runs. let t = Tree::new("survey-order"); t.file("DCIM/100CANON/IMG_0003.CR3", b"c") .file("DCIM/100CANON/IMG_0001.CR3", b"a") .file("DCIM/101CANON/IMG_0002.CR3", b"b"); assert_eq!(survey_of(&t), survey_of(&t)); } #[test] fn a_survey_carries_what_the_folder_template_needs() { let t = Tree::new("survey-meta"); t.file("DCIM/100CANON/IMG_0001.CR3", b"12345"); let found = survey_of(&t); assert_eq!(found[0].size, 5); // Seconds, not milliseconds — a template handed a millisecond reading // files everything in the year 56000. let year = dr_types::civil_from_unix(found[0].modified_at).year; assert!((2000..2200).contains(&year), "mtime looked like {year}"); } #[test] fn a_survey_can_be_stopped() { let t = Tree::new("survey-cancel"); t.file("DCIM/100CANON/IMG_0001.CR3", b"a"); let storage = LocalStorage::with_root(CARD, t.0.clone()); let found = survey(&storage, CARD, &FormatFilter::all(), &|| true).unwrap(); assert!(found.is_empty()); } #[test] fn a_card_is_recognised_by_its_dcim_folder() { let t = Tree::new("looks-like"); t.file("DCIM/100CANON/IMG_0001.CR3", b"a"); assert!(looks_like_a_card(&t.0)); assert!(!looks_like_a_card(&t.0.join("DCIM/100CANON"))); } // ---- what the interface says ----------------------------------------- fn outcome() -> Outcome { Outcome { imported: 3, bytes: 3 * 1024 * 1024, folders: vec!["2026/2026-08-22".into()], ..Default::default() } } #[test] fn a_finished_import_says_what_it_did_and_where() { let text = describe(&outcome()); assert!(text.starts_with("3 imported ("), "{text}"); assert!(text.ends_with(" into 2026/2026-08-22"), "{text}"); } #[test] fn a_card_that_was_already_imported_says_so_rather_than_nothing() { let o = Outcome { imported: 0, duplicates: 12, bytes: 0, folders: vec![], ..Default::default() }; // The failure mode this guards against is a dialog that closes with a // cheerful "done" after transferring nothing. assert_eq!(describe(&o), "12 already in the library"); } #[test] fn an_empty_card_is_not_reported_as_a_success() { assert_eq!(describe(&Outcome::default()), "Nothing to import"); } #[test] fn a_cancelled_run_leads_with_the_fact_that_it_stopped() { let o = Outcome { cancelled: true, ..outcome() }; assert!( describe(&o).starts_with("Stopped — 3 imported"), "{}", describe(&o) ); } #[test] fn undated_files_are_named_in_the_result() { // FR-NC-7a: a file dated by mtime is filed under when it was last // written, which for a card through a reader is often today. let o = Outcome { undated: vec!["SCAN_01.TIF".into()], ..outcome() }; let text = describe(&o); assert!(text.contains("1 had no capture time"), "{text}"); assert!(text.contains("was filed by modification date"), "{text}"); let many = Outcome { undated: vec!["a".into(), "b".into()], ..outcome() }; assert!(describe(&many).contains("2 had no capture time")); assert!(describe(&many).contains("were filed")); } #[test] fn several_days_on_one_card_are_counted_rather_than_listed() { let o = Outcome { folders: vec!["2026/2026-08-21".into(), "2026/2026-08-22".into()], ..outcome() }; assert!(describe(&o).ends_with("into 2 folders"), "{}", describe(&o)); } #[test] fn a_summary_collapses_a_days_worth_of_files_into_one_folder() { let report = Report { imported: (0..3) .map(|i| Imported { source: dr_types::SourceRef::Local { root: CARD, relative: format!("IMG_{i}.CR3"), }, written: dr_types::SourceRef::Local { root: LIBRARY, relative: format!("2026/2026-08-22/IMG_{i}.CR3"), }, folders: vec!["2026".into(), "2026-08-22".into()], name: format!("IMG_{i}.CR3"), digest: "x".into(), size: 10, dated_from: DateSource::Capture, backup: None, }) .collect(), bytes: 30, ..Default::default() }; let o = summarise(&report); assert_eq!(o.imported, 3); assert_eq!(o.folders, ["2026/2026-08-22"]); } }