wip: ingest

This commit is contained in:
2026-08-22 14:12:40 +02:00
parent 98e0ad0537
commit 735683b849
10 changed files with 1267 additions and 11 deletions
+1
View File
@@ -27,6 +27,7 @@ dr-plat.workspace = true
dr-sync.workspace = true
dr-sync-nextcloud.workspace = true
dr-export.workspace = true
dr-ingest.workspace = true
dr-pipeline.workspace = true
dr-catalog.workspace = true
dr-thumbs.workspace = true
+600
View File
@@ -0,0 +1,600 @@
//! 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_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<AtomicBool>;
/// 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<PathBuf>,
/// 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,
/// The library's root id *in the catalog*, which is what
/// [`dr_catalog::set_content_hash`] matches on. Unrelated to [`LIBRARY`],
/// which is this module's handle on the same folder.
pub catalog_root: u64,
/// Which file types to take off the card.
pub filter: FormatFilter,
pub options: Options,
}
/// 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,
},
/// 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<String>,
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<String>,
/// Card files that may now be deleted, for a move-import (FR-NC-7b).
///
/// Carried out rather than acted on here: this crate knows the local write
/// succeeded, and not whether the upload did.
pub retirable: 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<Message> {
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<Message>) -> 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(),
});
},
);
// The digests, so the *next* import can answer the content tier without
// reading anything. Best-effort and after the fact: the rows do not exist
// until a scan has caught up, and a hash that misses its row costs one
// wasted transfer next time rather than a failure now.
record_digests(catalog.connection(), request.catalog_root, &report.imported);
let _ = tx.send(Message::Finished(summarise(&report)));
Ok(())
}
/// 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<Vec<Candidate>, 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<Shot> {
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.
fn record_digests(conn: &rusqlite::Connection, root: u64, imported: &[Imported]) {
for image in imported {
if let Err(e) = dr_catalog::set_content_hash(conn, root, image.written.key(), &image.digest)
{
log::warn!("recording the digest of {}: {e}", image.name);
}
}
}
/// Reduce a report to what the interface says about it.
pub fn summarise(report: &Report) -> Outcome {
let mut folders: Vec<String> = 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(),
}
}
/// 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")),
}
// 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<Candidate> {
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::<Vec<_>>(),
["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"]);
}
}
+1
View File
@@ -26,6 +26,7 @@ mod develop;
mod export;
mod gradient;
mod histogram;
mod import;
mod labels;
mod library;
mod library_ui;
+20 -10
View File
@@ -1981,21 +1981,31 @@ fn collect_metadata(header: &[u8], req: &ThumbnailRequest, out: &mut Vec<Metadat
image_id: req.image_id,
captured_at: md.captured_at,
captured_offset: md.captured_offset,
camera: match (&md.make, &md.model) {
// Bodies repeat the make inside the model ("Canon EOS 6D"), so
// joining unconditionally yields "Canon Canon EOS 6D".
(Some(make), Some(model)) if model.starts_with(make.as_str()) => {
Some(model.trim().to_string())
}
(Some(make), Some(model)) => Some(format!("{} {}", make.trim(), model.trim())),
(None, Some(model)) => Some(model.trim().to_string()),
_ => None,
},
camera: camera_label(md.make.as_deref(), md.model.as_deref()),
lens: md.lens.map(|l| l.trim().to_string()),
iso: md.iso,
});
}
/// TRACES: FR-CAT-11
/// The camera string the catalog stores, from an EXIF make and model.
///
/// One definition rather than one per caller, because an import's duplicate
/// check compares against what a scan wrote (`dr_catalog::dedup`). Two
/// spellings of the same body would not fail loudly — they would silently
/// disable the cheap tier, and every re-inserted card would transfer in full
/// before the digest caught it.
pub fn camera_label(make: Option<&str>, model: Option<&str>) -> Option<String> {
match (make, model) {
// Bodies repeat the make inside the model ("Canon EOS 6D"), so
// joining unconditionally yields "Canon Canon EOS 6D".
(Some(make), Some(model)) if model.starts_with(make) => Some(model.trim().to_string()),
(Some(make), Some(model)) => Some(format!("{} {}", make.trim(), model.trim())),
(None, Some(model)) => Some(model.trim().to_string()),
_ => None,
}
}
/// Write a batch of dates and tell the UI, draining `found`.
///
/// Separate from the loop so the same path serves both the periodic flush and