wip: ingest

This commit is contained in:
2026-08-22 15:34:43 +02:00
parent 735683b849
commit 1526c957cf
9 changed files with 1570 additions and 26 deletions
+193 -19
View File
@@ -35,6 +35,7 @@ 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.
@@ -63,13 +64,37 @@ pub struct Request {
/// 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,
/// Where to send the originals afterwards, if anywhere (FR-NC-7b).
pub upload: Option<Upload>,
}
/// 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.
@@ -84,6 +109,8 @@ pub enum Message {
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.
@@ -107,9 +134,19 @@ pub struct Outcome {
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.
/// 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.
@@ -160,7 +197,9 @@ fn run(request: Request, cancel: &Cancel, tx: &Sender<Message>) -> Result<(), St
source: &card,
dest: &library,
dest_root: &library_root,
backup: backup.as_ref().map(|b| (b as &dyn WritableStorage, &backup_root)),
backup: backup
.as_ref()
.map(|b| (b as &dyn WritableStorage, &backup_root)),
options: request.options.clone(),
};
@@ -179,16 +218,115 @@ fn run(request: Request, cancel: &Cancel, tx: &Sender<Message>) -> Result<(), St
},
);
// 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 mut outcome = summarise(&report);
let _ = tx.send(Message::Finished(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<Message>,
) -> 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<Vec<u8>, 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
@@ -289,11 +427,27 @@ fn is_duplicate(conn: &rusqlite::Connection, key: &DupKey) -> bool {
}
/// 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);
///
/// 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<i64> = 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}");
}
}
}
@@ -317,6 +471,9 @@ pub fn summarise(report: &Report) -> Outcome {
cancelled: report.cancelled,
folders,
retirable: report.retirable.len(),
// Filled in by the upload phase, which runs after this.
uploaded: 0,
upload_failed: 0,
}
}
@@ -361,6 +518,19 @@ pub fn describe(outcome: &Outcome) -> String {
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!(
@@ -537,7 +707,11 @@ mod tests {
cancelled: true,
..outcome()
};
assert!(describe(&o).starts_with("Stopped — 3 imported"), "{}", describe(&o));
assert!(
describe(&o).starts_with("Stopped — 3 imported"),
"{}",
describe(&o)
);
}
#[test]