diff --git a/core/dr-thumbs/src/lib.rs b/core/dr-thumbs/src/lib.rs index 0a36e0e..351af9e 100644 --- a/core/dr-thumbs/src/lib.rs +++ b/core/dr-thumbs/src/lib.rs @@ -96,6 +96,16 @@ impl ThumbSize { } } + /// The class a stored discriminant names, or `None` for one this build + /// does not know. + pub fn from_stored(v: i64) -> Option { + match v { + 0 => Some(ThumbSize::Grid), + 1 => Some(ThumbSize::Large), + _ => None, + } + } + fn from_i64(v: i64) -> Self { match v { 1 => ThumbSize::Large, diff --git a/ui/dr-ui/src/develop_ui.rs b/ui/dr-ui/src/develop_ui.rs index 79aee35..458b674 100644 --- a/ui/dr-ui/src/develop_ui.rs +++ b/ui/dr-ui/src/develop_ui.rs @@ -315,7 +315,7 @@ fn wire_export(window: &AppWindow, w: &DevelopWiring) { total, to, move |written| { - drain_outbox(&library_for_drain); + drain_outbox(&library_for_drain, weak.clone()); // What the album now holds, and which photograph // each file came from. if let (Some(albums), Some(w)) = (albums.as_ref(), weak.upgrade()) { diff --git a/ui/dr-ui/src/export.rs b/ui/dr-ui/src/export.rs index ed8a1d4..e8c22e9 100644 --- a/ui/dr-ui/src/export.rs +++ b/ui/dr-ui/src/export.rs @@ -146,6 +146,20 @@ impl Pending { } } +/// TRACES: FR-MRG-6 +/// Where a file staged for `remote_dir` of the library at `root`, called +/// `name`, will be once it is uploaded — spelled as the scan will list it, +/// because the catalog row written ahead of the scan is keyed on it. +pub(crate) fn staged_remote_path(root: &str, remote_dir: &str, name: &str) -> RemotePath { + Pending { + local: PathBuf::new(), + remote_dir: remote_dir.to_string(), + account: false, + name: name.to_string(), + } + .remote_path(root) +} + /// Where an export was put, for the interface to report. #[derive(Debug, Clone, PartialEq, Eq)] pub enum Placed { @@ -332,6 +346,75 @@ pub(crate) fn destination_record(payload: &Path) -> PathBuf { )) } +/// TRACES: FR-MRG-6 +/// Where the thumbnails made with a staged payload wait for its upload: +/// `name.ext.thumbs` beside it. +/// +/// Not a record: [`pending`] reads only `.dest` files, so this never goes to +/// the server. The upload reads it once the server has named the file +/// ([`take_thumbnails`]) and puts its contents in the thumbnail store under +/// that name; until then there is no key to store them under. +pub(crate) fn thumbnails_record(payload: &Path) -> PathBuf { + let mut name = payload.as_os_str().to_owned(); + name.push(".thumbs"); + PathBuf::from(name) +} + +/// The magic line a thumbnails record starts with. +const THUMBS_MAGIC: &[u8] = b"darkroom-thumbs 1\n"; + +/// Write the thumbnails for a staged payload. A size-class byte and three +/// little-endian `u32`s — width, height, length — before each JPEG. +pub(crate) fn write_thumbnails( + payload: &Path, + thumbnails: &[(dr_thumbs::ThumbSize, dr_thumbs::Thumbnail)], +) -> std::io::Result<()> { + let mut out = THUMBS_MAGIC.to_vec(); + for (size, t) in thumbnails { + out.push(*size as i64 as u8); + for v in [t.width, t.height, t.bytes.len() as u32] { + out.extend_from_slice(&v.to_le_bytes()); + } + out.extend_from_slice(&t.bytes); + } + std::fs::write(thumbnails_record(payload), out) +} + +/// Read and remove the thumbnails waiting beside a payload. Nothing, for a +/// payload that has none or a record that does not parse — a composite whose +/// thumbnails are lost is thumbnailed the ordinary way. +pub(crate) fn take_thumbnails(payload: &Path) -> Vec<(dr_thumbs::ThumbSize, dr_thumbs::Thumbnail)> { + let path = thumbnails_record(payload); + let Ok(bytes) = std::fs::read(&path) else { + return Vec::new(); + }; + let _ = std::fs::remove_file(&path); + parse_thumbnails(&bytes).unwrap_or_default() +} + +fn parse_thumbnails(bytes: &[u8]) -> Option> { + let mut rest = bytes.strip_prefix(THUMBS_MAGIC)?; + let mut out = Vec::new(); + while let Some((&class, tail)) = rest.split_first() { + let word = |at: usize| -> Option { + Some(u32::from_le_bytes(tail.get(at..at + 4)?.try_into().ok()?)) + }; + let (width, height, len) = (word(0)?, word(4)?, word(8)? as usize); + let data = tail.get(12..12 + len)?; + let size = dr_thumbs::ThumbSize::from_stored(i64::from(class))?; + out.push(( + size, + dr_thumbs::Thumbnail { + width, + height, + bytes: data.to_vec(), + }, + )); + rest = &tail[12 + len..]; + } + Some(out) +} + /// Everything currently waiting in the outbox. /// /// A payload with no record is skipped rather than guessed at — see the write @@ -381,11 +464,12 @@ pub fn pending_count(outbox: &Path) -> usize { pending(outbox).len() } -/// Remove an entry and its record, once it is safely on the server. +/// Remove an entry and its records, once it is safely on the server. fn clear(entry: &Pending) { let record = PathBuf::from(format!("{}.dest", entry.local.display())); let _ = std::fs::remove_file(&entry.local); let _ = std::fs::remove_file(&record); + let _ = std::fs::remove_file(thumbnails_record(&entry.local)); } /// Progress from the upload worker. @@ -393,13 +477,170 @@ fn clear(entry: &Pending) { pub enum UploadMessage { Status(String), /// Uploaded, still pending, and the first error if there was one. + /// `landed` counts the uploads that went into the library itself rather + /// than an album outside it: those are what a scan has something to + /// find. Finished { uploaded: usize, remaining: usize, error: Option, + landed: usize, }, } +/// Log a drain as it goes and, once files have landed in the library, ask +/// the grid to rescan — the scan that finds them, run after the upload +/// rather than beside it (FR-MRG-6). +/// +/// On a thread of its own, since the drain reports only when it finishes and +/// that can be minutes for a composite; the window is reached through the +/// event loop. +pub(crate) fn watch_upload(rx: Receiver, window: slint::Weak) { + executors::spawn(Executor::Io, "upload-log", move || { + while let Ok(msg) = rx.recv() { + match msg { + UploadMessage::Status(s) => log::info!("export: {s}"), + UploadMessage::Finished { + uploaded, + remaining, + error, + landed, + } => { + log::info!("export: {uploaded} uploaded, {remaining} still queued"); + if let Some(e) = error { + log::warn!("export upload stopped: {e}"); + } + if landed > 0 { + let _ = window.upgrade_in_event_loop(|w| { + w.global::().invoke_library_rescan(); + }); + } + } + } + } + }); +} + +/// Where a library's catalog and thumbnail store are, for the upload to +/// record what the server made of a file. +pub(crate) struct LibraryFiles { + pub catalog: PathBuf, + pub thumbs: PathBuf, +} + +impl LibraryFiles { + fn of(account: &Account) -> Self { + LibraryFiles { + catalog: crate::library::catalog_path(account), + thumbs: crate::library::thumbs_dir(account), + } + } +} + +/// What became of one outbox entry. +#[derive(Debug, PartialEq, Eq)] +enum Sent { + Uploaded, + /// Its payload had gone; the record went with it. + Gone, +} + +/// Send one outbox entry, and clear it once it is on the server. +async fn send( + backend: &dyn dr_sync::RemoteBackend, + library: &LibraryFiles, + root: &str, + entry: &Pending, +) -> Result { + let Ok(bytes) = std::fs::read(&entry.local) else { + // The payload vanished under us. Drop the record too; retrying + // forever against a file that is gone helps nobody. + clear(entry); + return Ok(Sent::Gone); + }; + + // The folder may not exist — this is the first export into it — and + // `create_dir` treats "already there" as success, so it is unconditional + // rather than guarded by a check that would cost a request every time. + backend + .create_dir(&entry.remote_folder(root)) + .await + .map_err(|e| e.to_string())?; + backend + .put(&entry.remote_path(root), bytes, None) + .await + .map_err(|e| e.to_string())?; + let thumbnails = take_thumbnails(&entry.local); + register_upload(backend, library, root, entry, thumbnails).await; + clear(entry); + Ok(Sent::Uploaded) +} + +/// TRACES: FR-MRG-6 +/// Tell the catalog what the server made of a file it has just been given. +/// +/// Only for a file the catalog already has a row for — a composite the merge +/// catalogued before its upload. Anything else is the scan's to find, and +/// this spends no request on it. The listing is one request for the folder, +/// and it is the only way to learn the id a server assigns on upload: the +/// thumbnail store keys on it, and the thumbnails made during the merge wait +/// beside the payload until it is known. +async fn register_upload( + backend: &dyn dr_sync::RemoteBackend, + library: &LibraryFiles, + root: &str, + entry: &Pending, + thumbnails: Vec<(dr_thumbs::ThumbSize, dr_thumbs::Thumbnail)>, +) { + if entry.account { + return; + } + let path = entry.remote_path(root); + let catalog = match dr_catalog::Catalog::open(&library.catalog) { + Ok(c) => c, + Err(e) => { + log::warn!("upload: opening the catalog to record {}: {e}", entry.name); + return; + } + }; + let Some(image) = crate::library::image_at(&catalog, root, path.as_str()) else { + return; + }; + let listed = match backend.list(&entry.remote_folder(root), None).await { + Ok(entries) => entries.into_iter().find(|e| e.path == path), + Err(e) => { + log::warn!("upload: listing {} after sending it: {e}", entry.name); + None + } + }; + let Some(listed) = listed else { + log::warn!("upload: {} is not listed where it was sent", entry.name); + return; + }; + if let Err(e) = crate::library::record_uploaded(&catalog, image, &listed) { + log::warn!("upload: recording {}: {e}", entry.name); + } + let dr_sync::RemoteId::Stable(file_id) = listed.id else { + return; + }; + if thumbnails.is_empty() { + return; + } + match dr_thumbs::ThumbStore::open(&library.thumbs) { + Ok(mut store) => { + for (size, thumb) in &thumbnails { + crate::library::store_thumbnail(&mut store, file_id, *size, thumb); + } + log::info!( + "upload: {} thumbnail(s) of {} stored under file {file_id}", + thumbnails.len(), + entry.name + ); + } + Err(e) => log::warn!("upload: opening the thumbnail store: {e}"), + } +} + /// Held by whichever drain is sending the outbox. See [`spawn_upload`]. static OUTBOX_DRAIN: std::sync::Mutex<()> = std::sync::Mutex::new(()); @@ -436,6 +677,7 @@ pub fn spawn_upload( uploaded: 0, remaining: pending_count(&outbox), error: Some(e.to_string()), + landed: 0, }); return; } @@ -449,14 +691,17 @@ pub fn spawn_upload( uploaded: 0, remaining: pending_count(&outbox), error: Some(e.to_string()), + landed: 0, }); return; } }; + let library = LibraryFiles::of(&conn.account); let queue = pending(&outbox); let total = queue.len(); let mut uploaded = 0; + let mut landed = 0; let mut error = None; for (i, entry) in queue.iter().enumerate() { @@ -466,30 +711,16 @@ pub fn spawn_upload( i + 1 ))); - let Ok(bytes) = std::fs::read(&entry.local) else { - // The payload vanished under us. Drop the record too; - // retrying forever against a file that is gone helps - // nobody. - clear(entry); - continue; - }; - - // The folder may not exist — this is the first export into it - // — and `create_dir` treats "already there" as success, so it - // is unconditional rather than guarded by a check that would - // cost a request every time. - if let Err(e) = backend.create_dir(&entry.remote_folder(&root)).await { - error = Some(e.to_string()); - break; - } - - match backend.put(&entry.remote_path(&root), bytes, None).await { - Ok(_) => { - clear(entry); + match send(&*backend, &library, &root, entry).await { + Ok(Sent::Gone) => {} + Ok(Sent::Uploaded) => { uploaded += 1; + if !entry.account { + landed += 1; + } } Err(e) => { - error = Some(e.to_string()); + error = Some(e); break; } } @@ -499,6 +730,7 @@ pub fn spawn_upload( uploaded, remaining: pending_count(&outbox), error, + landed, }); }); }); @@ -1295,6 +1527,328 @@ fn stop_timer(slot: &Rc>>) { mod tests { use super::*; + /// A server that names a file when it is given one, as Nextcloud does: + /// the id exists only once the upload has landed, and a listing is the + /// way to learn it. + #[derive(Default)] + struct Assigning { + files: std::sync::Mutex)>>, + lists: std::sync::atomic::AtomicUsize, + caps: std::sync::OnceLock, + } + + impl Assigning { + fn listing(&self, dir: &str) -> Vec { + self.files + .lock() + .unwrap() + .iter() + .filter(|(p, _)| p.rsplit_once('/').map_or("", |(d, _)| d) == dir) + .map(|(p, (id, body))| dr_sync::RemoteEntry { + id: dr_sync::RemoteId::Stable(*id), + path: RemotePath::new(p.clone()), + kind: dr_sync::EntryKind::File, + validator: dr_sync::Validator::new(format!("etag-{id}")), + size: body.len() as u64, + modified: None, + has_preview: false, + materialised: true, + }) + .collect() + } + } + + #[async_trait::async_trait] + impl dr_sync::RemoteBackend for Assigning { + fn capabilities(&self) -> &dr_sync::Capabilities { + self.caps.get_or_init(dr_sync::Capabilities::minimal) + } + fn name(&self) -> &str { + "assigning" + } + async fn list( + &self, + dir: &RemotePath, + _since: Option<&dr_sync::Validator>, + ) -> Result, dr_sync::RemoteError> { + self.lists.fetch_add(1, Ordering::SeqCst); + Ok(self.listing(dir.as_str())) + } + async fn dir_validator( + &self, + _dir: &RemotePath, + ) -> Result { + Err(dr_sync::RemoteError::Unsupported("test")) + } + async fn delta( + &self, + _c: &dr_sync::Cursor, + ) -> Result<(Vec, dr_sync::Cursor), dr_sync::RemoteError> { + Err(dr_sync::RemoteError::Unsupported("test")) + } + async fn get( + &self, + _id: &dr_sync::RemoteId, + _r: Option>, + ) -> Result, dr_sync::RemoteError> { + Err(dr_sync::RemoteError::Unsupported("test")) + } + async fn put( + &self, + path: &RemotePath, + body: Vec, + _pc: Option, + ) -> Result { + let mut files = self.files.lock().unwrap(); + let id = files + .get(path.as_str()) + .map(|(id, _)| *id) + .unwrap_or(40_000 + files.len() as u64); + files.insert(path.as_str().to_string(), (id, body)); + Ok(dr_sync::Validator::new(format!("etag-{id}"))) + } + async fn delete( + &self, + _id: &dr_sync::RemoteId, + _pc: Option, + ) -> Result<(), dr_sync::RemoteError> { + Ok(()) + } + async fn move_to( + &self, + _f: &dr_sync::RemoteId, + _t: &RemotePath, + ) -> Result<(), dr_sync::RemoteError> { + Ok(()) + } + async fn create_dir(&self, _p: &RemotePath) -> Result<(), dr_sync::RemoteError> { + Ok(()) + } + } + + /// A library of one sweep, its catalog on disk, a store beside it, and + /// an outbox holding the composite the merge staged: the row written + /// ahead of the scan, the payload, its record, and its thumbnails. + struct Staged { + dir: PathBuf, + library: LibraryFiles, + entry: Pending, + image: i64, + } + + fn staged(name: &str, root: &str) -> Staged { + let dir = + std::env::temp_dir().join(format!("dr-composite-upload-{name}-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(dir.join("outbox")).unwrap(); + let library = LibraryFiles { + catalog: dir.join("catalog.sqlite"), + thumbs: dir.join("thumbs"), + }; + let at = |n: &str| { + if root.is_empty() { + format!("Alps/{n}") + } else { + format!("{root}/Alps/{n}") + } + }; + let catalog = dr_catalog::Catalog::open(&library.catalog).unwrap(); + crate::library::test_support::scanned_listing( + &catalog, + root, + vec![crate::library::test_support::entry( + &at("_MG_8320.CR2"), + 1, + 30, + )], + ); + let image = crate::library::catalogue_composite( + &catalog, + root, + &crate::library::CompositeRow { + source_ref: staged_remote_path(root, "Alps", "_MG_8320-pano.dng") + .as_str() + .to_string(), + width: 4000, + height: 1000, + captured_at: Some(1_000), + captured_offset: None, + camera: None, + lens: None, + iso: None, + file_size: 5, + sources: vec![at("_MG_8320.CR2")], + }, + ) + .unwrap(); + // The scan that ran while the upload was still going: the folder is + // listed without the composite, and the row is still there after it. + crate::library::test_support::scanned_listing( + &catalog, + root, + vec![crate::library::test_support::entry( + &at("_MG_8320.CR2"), + 1, + 30, + )], + ); + let local = dir.join("outbox").join("_MG_8320-pano.dng"); + std::fs::write(&local, b"pano!").unwrap(); + std::fs::write(destination_record(&local), "Alps\n_MG_8320-pano.dng\n").unwrap(); + let jpeg = dr_thumbs::encode_rgba(4, 1, &[200u8; 16]).unwrap(); + write_thumbnails( + &local, + &[ + ( + dr_thumbs::ThumbSize::Grid, + dr_thumbs::Thumbnail { + width: 4, + height: 1, + bytes: jpeg.clone(), + }, + ), + ( + dr_thumbs::ThumbSize::Large, + dr_thumbs::Thumbnail { + width: 4, + height: 1, + bytes: jpeg, + }, + ), + ], + ) + .unwrap(); + let entry = pending(&dir.join("outbox")).pop().unwrap(); + Staged { + dir, + library, + entry, + image, + } + } + + fn file_id_of(library: &LibraryFiles, image: i64) -> Option { + dr_catalog::Catalog::open(&library.catalog) + .unwrap() + .connection() + .query_row( + "SELECT file_id FROM remote WHERE image_id = ?1", + [image], + |r| r.get(0), + ) + .ok() + } + + #[tokio::test] + async fn an_upload_that_lands_after_the_scan_gives_the_row_its_server_id() { + let s = staged("nextcloud", "PhotosRaw"); + let server = Assigning::default(); + assert_eq!(file_id_of(&s.library, s.image), None); + + let sent = send(&server, &s.library, "PhotosRaw", &s.entry) + .await + .unwrap(); + assert_eq!(sent, Sent::Uploaded); + + // The id the server assigned, on the row the merge wrote, and the + // merge's thumbnails in the store under it. + let listed = server.listing("PhotosRaw/Alps"); + let dr_sync::RemoteId::Stable(id) = listed[0].id else { + panic!("the server names its files"); + }; + assert_eq!(file_id_of(&s.library, s.image), Some(id as i64)); + let store = dr_thumbs::ThumbStore::open(&s.library.thumbs).unwrap(); + assert!(store.contains(id, dr_thumbs::ThumbSize::Grid)); + assert!(store.contains(id, dr_thumbs::ThumbSize::Large)); + // Nothing left in the outbox: not the payload, not its records. + assert_eq!( + std::fs::read_dir(s.dir.join("outbox")).unwrap().count(), + 0, + "the outbox is empty" + ); + + // The scan that follows the upload finds the row already there. + let catalog = dr_catalog::Catalog::open(&s.library.catalog).unwrap(); + crate::library::test_support::scanned_listing(&catalog, "PhotosRaw", listed); + let n: i64 = catalog + .connection() + .query_row( + "SELECT count(*) FROM images WHERE source_ref LIKE '%pano%'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(n, 1); + let _ = std::fs::remove_dir_all(&s.dir); + } + + #[tokio::test] + async fn a_folder_library_learns_the_id_the_folder_gives_its_path() { + let s = staged("folder", ""); + let library_dir = s.dir.join("library"); + std::fs::create_dir_all(library_dir.join("Alps")).unwrap(); + let folder = dr_sync_folder::FolderBackend::new(&library_dir).unwrap(); + + send(&folder, &s.library, "", &s.entry).await.unwrap(); + + assert_eq!( + std::fs::read(library_dir.join("Alps/_MG_8320-pano.dng")).unwrap(), + b"pano!" + ); + use dr_sync::RemoteBackend as _; + let listed = folder + .list(&RemotePath::new("Alps"), None) + .await + .unwrap() + .into_iter() + .find(|e| e.path.name() == "_MG_8320-pano.dng") + .unwrap(); + let dr_sync::RemoteId::Stable(id) = listed.id else { + panic!("a folder names a file by its path"); + }; + assert_eq!(file_id_of(&s.library, s.image), Some(id as i64)); + let store = dr_thumbs::ThumbStore::open(&s.library.thumbs).unwrap(); + assert!(store.contains(id, dr_thumbs::ThumbSize::Grid)); + let _ = std::fs::remove_dir_all(&s.dir); + } + + #[tokio::test] + async fn an_ordinary_export_costs_no_listing() { + let s = staged("export", "PhotosRaw"); + let other = s.dir.join("outbox").join("print.jpg"); + std::fs::write(&other, b"jpeg").unwrap(); + std::fs::write(destination_record(&other), "Alps\nprint.jpg\n").unwrap(); + let entry = pending(&s.dir.join("outbox")) + .into_iter() + .find(|p| p.name == "print.jpg") + .unwrap(); + let server = Assigning::default(); + send(&server, &s.library, "PhotosRaw", &entry) + .await + .unwrap(); + assert_eq!(server.lists.load(Ordering::SeqCst), 0); + let _ = std::fs::remove_dir_all(&s.dir); + } + + #[test] + fn thumbnails_survive_the_outbox_and_a_damaged_record_is_nothing() { + let dir = tmp(); + let payload = dir.join("x.dng"); + let t = dr_thumbs::Thumbnail { + width: 3, + height: 2, + bytes: vec![1, 2, 3, 4, 5], + }; + write_thumbnails(&payload, &[(dr_thumbs::ThumbSize::Large, t.clone())]).unwrap(); + assert_eq!( + take_thumbnails(&payload), + vec![(dr_thumbs::ThumbSize::Large, t)] + ); + assert!(!thumbnails_record(&payload).exists(), "taken, not copied"); + std::fs::write(thumbnails_record(&payload), b"darkroom-thumbs 1\n\x00\x01").unwrap(); + assert!(take_thumbnails(&payload).is_empty()); + } + fn encoded(name: &str, bytes: &[u8]) -> Encoded { Encoded { name: name.to_string(), diff --git a/ui/dr-ui/src/lib.rs b/ui/dr-ui/src/lib.rs index 4a2ac15..cab43d1 100644 --- a/ui/dr-ui/src/lib.rs +++ b/ui/dr-ui/src/lib.rs @@ -731,7 +731,7 @@ fn batch_request( /// the first is what makes an upload feel immediate, and the second is what /// eventually delivers the exports made while the train was in a tunnel. /// Running it twice over an empty outbox costs a directory listing. -fn drain_outbox(library: &Rc) { +fn drain_outbox(library: &Rc, window: slint::Weak) { // Offline is not a failure worth reporting here — the entries stay // staged and the next pass takes them. if library.is_offline() { @@ -747,23 +747,7 @@ fn drain_outbox(library: &Rc) { let root = conn.account.root.clone(); let rx = export::spawn_upload(conn, root, outbox); - executors::spawn(executors::Executor::Io, "upload-log", move || { - while let Ok(msg) = rx.recv() { - match msg { - export::UploadMessage::Status(s) => log::info!("export: {s}"), - export::UploadMessage::Finished { - uploaded, - remaining, - error, - } => { - log::info!("export: {uploaded} uploaded, {remaining} still queued"); - if let Some(e) = error { - log::warn!("export upload stopped: {e}"); - } - } - } - } - }); + export::watch_upload(rx, window); } /// What the export button should say, given where an export would go. @@ -1901,17 +1885,31 @@ fn wire_import_and_merge( move || library_for_sources.export_sources(&collections_for_sources.selected()), move || { let conn = library_for_context.session()?; + let catalog = library_for_context.catalog(); + let root = conn.account.root.clone(); Some(merge_ui::Context { outbox: export::outbox_dir(&conn.account), conn, + names_in: Box::new(move |folder| { + catalog + .borrow() + .as_ref() + .map(|c| library::names_in_folder(c, &root, folder)) + .unwrap_or_default() + }), }) }, - move |w| { - // Staged beside its sources: upload it now rather than - // on the next sync pass, then look for it, exactly as an - // import does. - drain_outbox(&library_for_done); - w.global::().invoke_library_rescan(); + move |w, placed| { + // TRACES: FR-MRG-6 + // In the grid now, from what the merge knows: a scan + // started here raced the upload and did not find it. + if let Some(placed) = placed { + library_ui::catalogue_composite(w, &library_for_done, placed); + } + // Upload it now rather than on the next sync pass; the + // scan that follows the upload gives the row the file id + // the server assigned. + drain_outbox(&library_for_done, w.as_weak()); }, ); } diff --git a/ui/dr-ui/src/library/composite.rs b/ui/dr-ui/src/library/composite.rs new file mode 100644 index 0000000..2a21471 --- /dev/null +++ b/ui/dr-ui/src/library/composite.rs @@ -0,0 +1,461 @@ +//! TRACES: FR-MRG-6 +//! A merged panorama in the catalog the moment it is written, rather than +//! whenever a scan next happens to find it. +//! +//! # Why the merge catalogues its own file +//! +//! The composite reaches the library through the outbox, like an export, +//! and the grid used to learn of it the way it learns of anything: a rescan, +//! fired as the merge finished. That scan raced the upload it followed. On a +//! folder library the copy of an 800 MB file was still running when the +//! folder was listed; on Nextcloud the upload takes minutes. Either way the +//! listing did not have the file, the folder's validator was recorded as it +//! stood, and nothing looked again until the next sync pass — the panorama +//! the photographer had just watched being made was not in the grid. +//! +//! So the job hands the library what it knows — the name the file will have, +//! the size of the picture, when it was taken — and the row is written at +//! once, in one transaction. The upload later supplies what only the server +//! knows, its file id, through [`record_uploaded`]; the scan that follows +//! finds a row already there and updates it in place, on the same +//! `(root_id, source_ref)` key it would have inserted under. + +use dr_catalog::Catalog; + +use super::scan::now_secs; + +/// What the catalog is told about a composite before any scan has seen it. +#[derive(Debug, Clone, PartialEq)] +pub struct CompositeRow { + /// Where it will be in the library, spelled as the scan spells it. + pub source_ref: String, + /// The picture it opens on — the crop, where the merge cropped — upright. + pub width: u32, + pub height: u32, + /// When it was taken: the middle of its sweep, as the DNG says. `None` + /// where no frame carried a time, and then the earliest of `sources` as + /// the catalog has them. + pub captured_at: Option, + pub captured_offset: Option, + pub camera: Option, + pub lens: Option, + pub iso: Option, + pub file_size: u64, + /// The frames it was merged from, as source refs. + pub sources: Vec, +} + +/// Write the composite's row, and return its image id. +/// +/// One transaction: the root, the capture time the sources lend it if it has +/// none of its own, and the row. At `metadata_state = 2`, since everything a +/// header read would add is already here — the sweep has no reason to fetch +/// the tail of an 800 MB file for a date the merge wrote into it. +/// +/// A row already under the name is brought up to date rather than +/// duplicated. The merge picks a name the catalog does not hold, so that is +/// a safety net: the file at that name is the one just written. +pub fn catalogue_composite( + catalog: &Catalog, + root: &str, + row: &CompositeRow, +) -> Result { + let conn = catalog.connection(); + let tx = conn.unchecked_transaction()?; + + tx.execute( + "INSERT INTO roots(kind, label, last_seen) VALUES ('remote', ?1, ?2) + ON CONFLICT DO NOTHING", + rusqlite::params![root, now_secs()], + )?; + let root_id: i64 = tx.query_row( + "SELECT id FROM roots WHERE label = ?1 AND kind = 'remote'", + [root], + |r| r.get(0), + )?; + + // The folder, if a scan has recorded it; the sources are in it, so it + // almost always has. `None` is what the scan itself writes for a file at + // the root, and the next scan fills it either way. + let folder_id: Option = match row.source_ref.rsplit_once('/') { + Some((parent, _)) => tx + .query_row( + "SELECT id FROM folders WHERE root_id = ?1 AND path = ?2", + rusqlite::params![root_id, parent], + |r| r.get(0), + ) + .ok(), + None => None, + }; + + // The sources' earliest, where the composite carries no time. One + // aggregate over the handful of frames, by their key. + let captured_at = match row.captured_at { + Some(t) => Some(t), + None if row.sources.is_empty() => None, + None => { + let placeholders = std::iter::repeat_n("?", row.sources.len()) + .collect::>() + .join(","); + let mut params: Vec = + vec![rusqlite::types::Value::Integer(root_id)]; + params.extend( + row.sources + .iter() + .map(|s| rusqlite::types::Value::Text(s.clone())), + ); + tx.query_row( + &format!( + "SELECT min(captured_at) FROM images + WHERE root_id = ? AND source_ref IN ({placeholders})" + ), + rusqlite::params_from_iter(params.iter()), + |r| r.get::<_, Option>(0), + )? + } + }; + + let format = row + .source_ref + .rsplit_once('.') + .map(|(_, e)| e.to_ascii_lowercase()); + tx.execute( + "INSERT INTO images(root_id, folder_id, source_ref, format, w, h, + captured_at, captured_offset, camera, lens, iso, + file_size, availability, metadata_state, added_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, 0, 2, ?13) + ON CONFLICT(root_id, source_ref) DO UPDATE SET + folder_id = coalesce(excluded.folder_id, images.folder_id), + w = excluded.w, + h = excluded.h, + captured_at = coalesce(excluded.captured_at, images.captured_at), + captured_offset = coalesce(excluded.captured_offset, images.captured_offset), + camera = coalesce(excluded.camera, images.camera), + lens = coalesce(excluded.lens, images.lens), + iso = coalesce(excluded.iso, images.iso), + file_size = excluded.file_size, + metadata_state = 2", + rusqlite::params![ + root_id, + folder_id, + row.source_ref, + format, + row.width as i64, + row.height as i64, + captured_at, + row.captured_offset, + row.camera, + row.lens, + row.iso.map(i64::from), + row.file_size as i64, + now_secs(), + ], + )?; + let image_id: i64 = tx.query_row( + "SELECT id FROM images WHERE root_id = ?1 AND source_ref = ?2", + rusqlite::params![root_id, row.source_ref], + |r| r.get(0), + )?; + tx.commit()?; + Ok(image_id) +} + +/// The image a library path names, if the catalog has it. What the upload +/// asks before it spends a listing on a file: a file the catalog has no row +/// for is the scan's to find, and only a row written ahead of the scan is +/// waiting for its server identity. +pub fn image_at(catalog: &Catalog, root: &str, source_ref: &str) -> Option { + catalog + .connection() + .query_row( + "SELECT i.id FROM images i JOIN roots r ON r.id = i.root_id + WHERE r.label = ?1 AND r.kind = 'remote' AND i.source_ref = ?2", + rusqlite::params![root, source_ref], + |r| r.get(0), + ) + .ok() +} + +/// The names the catalog holds directly in `folder` of the library at +/// `root`, whatever their state — what a new file there must not be called. +/// +/// A range over the `(root_id, source_ref)` key, `/` to `0` being the byte +/// after it, so a folder of a thousand costs a seek and a thousand index +/// entries, not the library. +pub fn names_in_folder( + catalog: &Catalog, + root: &str, + folder: &str, +) -> std::collections::HashSet { + let folder = folder.trim_matches('/'); + let (from, to) = if folder.is_empty() { + (String::new(), "\u{10ffff}".to_string()) + } else { + (format!("{folder}/"), format!("{folder}0")) + }; + let Ok(mut stmt) = catalog.connection().prepare( + "SELECT i.source_ref FROM images i JOIN roots r ON r.id = i.root_id + WHERE r.label = ?1 AND r.kind = 'remote' + AND i.source_ref >= ?2 AND i.source_ref < ?3", + ) else { + return Default::default(); + }; + let Ok(rows) = stmt.query_map(rusqlite::params![root, from, to], |r| r.get::<_, String>(0)) + else { + return Default::default(); + }; + rows.flatten() + .filter_map(|p| { + let name = p.get(from.len()..)?; + (!name.contains('/')).then(|| name.to_string()) + }) + .collect() +} + +/// Give a catalogued file the identity the server assigned it on upload: +/// its file id, which the thumbnail store keys on, and its validator. +/// +/// The same row the scan's `remote` upsert writes, so whichever of the two +/// runs second agrees with the first. +pub fn record_uploaded( + catalog: &Catalog, + image_id: i64, + entry: &dr_sync::RemoteEntry, +) -> Result<(), dr_catalog::CatalogError> { + let conn = catalog.connection(); + let tx = conn.unchecked_transaction()?; + if let dr_sync::RemoteId::Stable(file_id) = entry.id { + tx.execute( + "INSERT INTO remote(image_id, file_id, etag, remote_path) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT(image_id) DO UPDATE SET + file_id = excluded.file_id, etag = excluded.etag, + remote_path = excluded.remote_path", + rusqlite::params![ + image_id, + file_id as i64, + entry.validator.as_str(), + entry.path.as_str() + ], + )?; + } + tx.execute( + "UPDATE images SET file_size = ?2 WHERE id = ?1 AND file_size IS NOT ?2", + rusqlite::params![image_id, entry.size as i64], + )?; + tx.commit()?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::library::read_cells; + use crate::library::scan::persist; + use crate::library::test_support::entry; + use dr_sync::RemotePath; + + /// Twelve frames of a sweep a second apart, scanned and dated. + fn a_sweep() -> Catalog { + let catalog = Catalog::in_memory().unwrap(); + let images = (0..12) + .map(|i| { + entry( + &format!("Alps/_MG_{:04}.CR2", 8320 + i), + 100 + i, + 30_000_000, + ) + }) + .chain(std::iter::once(entry("Alps/later.CR2", 200, 30_000_000))) + .collect(); + persist( + &catalog, + "", + &dr_sync::ScanResult { + images, + directories: vec![(RemotePath::new("Alps"), dr_sync::Validator::new("e1"))], + progress: Default::default(), + sidecars: Vec::new(), + }, + ) + .unwrap(); + for i in 0..12 { + catalog + .connection() + .execute( + "UPDATE images SET captured_at = ?2, metadata_state = 2 + WHERE source_ref = ?1", + rusqlite::params![format!("Alps/_MG_{:04}.CR2", 8320 + i), 1_000 + i], + ) + .unwrap(); + } + catalog + .connection() + .execute( + "UPDATE images SET captured_at = 5000 WHERE source_ref = 'Alps/later.CR2'", + [], + ) + .unwrap(); + catalog + } + + fn composite(captured_at: Option) -> CompositeRow { + CompositeRow { + source_ref: "Alps/_MG_8320-pano.dng".into(), + width: 22_000, + height: 5_600, + captured_at, + captured_offset: None, + camera: Some("Canon EOS 6D".into()), + lens: None, + iso: Some(100), + file_size: 800_000_000, + sources: (0..12) + .map(|i| format!("Alps/_MG_{:04}.CR2", 8320 + i)) + .collect(), + } + } + + fn names(catalog: &Catalog) -> Vec { + read_cells(catalog, 0, 100) + .unwrap() + .into_iter() + .map(|c| c.name) + .collect() + } + + #[test] + fn a_composite_sits_among_its_sources_the_moment_it_is_written() { + let catalog = a_sweep(); + catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); + let names = names(&catalog); + let at = names.iter().position(|n| n == "_MG_8320-pano.dng").unwrap(); + // Dated at the middle of the sweep, so it is among the frames and not + // after the photograph taken an hour later. + assert!(at > 0 && at < 12, "{names:?}"); + assert_eq!(names.last().unwrap(), "later.CR2"); + } + + #[test] + fn a_composite_with_no_time_takes_its_earliest_source() { + let catalog = a_sweep(); + let id = catalogue_composite(&catalog, "", &composite(None)).unwrap(); + let (t, w, h, state): (Option, i64, i64, i64) = catalog + .connection() + .query_row( + "SELECT captured_at, w, h, metadata_state FROM images WHERE id = ?1", + [id], + |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)), + ) + .unwrap(); + assert_eq!(t, Some(1_000)); + assert_eq!((w, h, state), (22_000, 5_600, 2)); + } + + #[test] + fn a_scan_that_ran_before_the_upload_does_not_lose_the_row() { + // Nextcloud: the merge finishes, the row is written, and a sync pass + // lists the folder before the upload has finished — the file is not + // on the server yet. + let catalog = a_sweep(); + let id = catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); + persist( + &catalog, + "", + &dr_sync::ScanResult { + images: vec![entry("Alps/later.CR2", 200, 30_000_000)], + directories: vec![(RemotePath::new("Alps"), dr_sync::Validator::new("e2"))], + progress: Default::default(), + sidecars: Vec::new(), + }, + ) + .unwrap(); + assert!(names(&catalog).contains(&"_MG_8320-pano.dng".to_string())); + + // The upload lands and the server names it; then the scan finds it. + let listed = entry("Alps/_MG_8320-pano.dng", 9_001, 812_345_678); + record_uploaded(&catalog, id, &listed).unwrap(); + persist( + &catalog, + "", + &dr_sync::ScanResult { + images: vec![listed], + directories: vec![(RemotePath::new("Alps"), dr_sync::Validator::new("e3"))], + progress: Default::default(), + sidecars: Vec::new(), + }, + ) + .unwrap(); + + let rows: Vec<(i64, Option, i64, Option)> = catalog + .connection() + .prepare( + "SELECT i.id, r.file_id, i.file_size, i.captured_at FROM images i + LEFT JOIN remote r ON r.image_id = i.id + WHERE i.source_ref = 'Alps/_MG_8320-pano.dng'", + ) + .unwrap() + .query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?))) + .unwrap() + .map(Result::unwrap) + .collect(); + assert_eq!(rows, vec![(id, Some(9_001), 812_345_678, Some(1_005))]); + let cell = read_cells(&catalog, 0, 100) + .unwrap() + .into_iter() + .find(|c| c.image_id == id) + .unwrap(); + assert_eq!(cell.file_id, Some(9_001)); + } + + #[test] + fn merging_again_under_the_same_name_updates_the_one_row() { + let catalog = a_sweep(); + let first = catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); + let mut again = composite(Some(1_006)); + again.width = 21_000; + let second = catalogue_composite(&catalog, "", &again).unwrap(); + assert_eq!(first, second); + let w: i64 = catalog + .connection() + .query_row("SELECT w FROM images WHERE id = ?1", [first], |r| r.get(0)) + .unwrap(); + assert_eq!(w, 21_000); + } + + #[test] + fn the_names_in_a_folder_are_its_own_and_not_its_subfolders() { + let catalog = a_sweep(); + catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); + persist( + &catalog, + "", + &dr_sync::ScanResult { + images: vec![ + entry("Alps/2019/deep.CR2", 300, 1), + entry("Alpsine/near.CR2", 301, 1), + ], + directories: Vec::new(), + progress: Default::default(), + sidecars: Vec::new(), + }, + ) + .unwrap(); + let names = names_in_folder(&catalog, "", "Alps"); + assert!(names.contains("_MG_8320-pano.dng")); + assert!(names.contains("_MG_8320.CR2")); + assert!(!names + .iter() + .any(|n| n.contains("deep") || n.contains("near"))); + assert_eq!(names.len(), 14); + } + + #[test] + fn only_a_row_the_merge_wrote_is_found_for_the_upload() { + let catalog = a_sweep(); + assert!(image_at(&catalog, "", "Alps/_MG_8320-pano.dng").is_none()); + let id = catalogue_composite(&catalog, "", &composite(Some(1_005))).unwrap(); + assert_eq!(image_at(&catalog, "", "Alps/_MG_8320-pano.dng"), Some(id)); + assert!(image_at(&catalog, "elsewhere", "Alps/_MG_8320-pano.dng").is_none()); + } +} diff --git a/ui/dr-ui/src/library/mod.rs b/ui/dr-ui/src/library/mod.rs index 0d43b8e..c0e2bea 100644 --- a/ui/dr-ui/src/library/mod.rs +++ b/ui/dr-ui/src/library/mod.rs @@ -28,6 +28,7 @@ //! `library` needs to change. mod cells; +mod composite; mod filters; mod paths; mod scan; @@ -39,6 +40,7 @@ mod thumbnails_gen; mod xmp; pub use cells::*; +pub use composite::*; pub use filters::*; pub use paths::*; pub use scan::*; @@ -84,6 +86,21 @@ pub(crate) mod test_support { } } + /// Write a scan's listing of `images` under `root`, as a scan would. + pub fn scanned_listing(catalog: &Catalog, root: &str, images: Vec) { + persist( + catalog, + root, + &dr_sync::ScanResult { + images, + directories: Vec::new(), + progress: Default::default(), + sidecars: Vec::new(), + }, + ) + .unwrap(); + } + /// A catalog with `n` images, ready to file into collections. pub fn with_images(n: usize) -> Catalog { let catalog = Catalog::in_memory().unwrap(); diff --git a/ui/dr-ui/src/library_ui/mod.rs b/ui/dr-ui/src/library_ui/mod.rs index 40a2640..2651023 100644 --- a/ui/dr-ui/src/library_ui/mod.rs +++ b/ui/dr-ui/src/library_ui/mod.rs @@ -48,8 +48,8 @@ mod window; pub use controller::LibraryController; pub use grid::wire; +pub use open::{catalogue_composite, open, reload}; pub(crate) use open::{forget_catalog, show_catalog_now, start_rescan}; -pub use open::{open, reload}; pub use ratings_keywords::paste_settings_to_selection; pub(crate) use ratings_keywords::{save_judgements, start_sidecar_writes}; pub use timeline::format_date; diff --git a/ui/dr-ui/src/library_ui/open.rs b/ui/dr-ui/src/library_ui/open.rs index bb878a1..fd5e8c0 100644 --- a/ui/dr-ui/src/library_ui/open.rs +++ b/ui/dr-ui/src/library_ui/open.rs @@ -41,6 +41,52 @@ pub fn reload(window: &AppWindow, ctl: &Rc) { note_place(window, ctl); } +/// TRACES: FR-MRG-6 +/// Put a composite the merge has just staged into the catalog and the grid, +/// before its upload has finished and before any scan could find it. +/// +/// See `library::composite` for why this does not wait for the scan. The +/// row is keyed where the scan will list the file, so the scan after the +/// upload updates it rather than adding a second. +pub fn catalogue_composite( + window: &AppWindow, + ctl: &Rc, + placed: &crate::merge_ui::Placed, +) { + let Some((conn, _)) = ctl.session.borrow().clone() else { + return; + }; + let root = conn.account.root.clone(); + let c = &placed.composite; + let source_ref = crate::export::staged_remote_path(&root, &placed.remote_dir, &c.name); + let row = library::CompositeRow { + source_ref: source_ref.as_str().to_string(), + width: c.width, + height: c.height, + captured_at: c.captured_at, + captured_offset: c.captured_offset, + camera: c.camera.clone(), + lens: c.lens.clone(), + iso: c.iso, + file_size: c.file_size, + sources: placed.sources.clone(), + }; + let written = { + let borrow = ctl.catalog.borrow(); + let Some(catalog) = borrow.as_ref() else { + return; + }; + library::catalogue_composite(catalog, &root, &row) + }; + match written { + Ok(image) => { + log::info!("merge: {} catalogued as image {image}", row.source_ref); + reload(window, ctl); + } + Err(e) => log::warn!("merge: cataloguing {}: {e}", row.source_ref), + } +} + /// Open a library: show the grid, start a scan, then fill in thumbnails. /// /// Called from the launch screen's "Open library" button — the callback that diff --git a/ui/dr-ui/src/library_ui/sync.rs b/ui/dr-ui/src/library_ui/sync.rs index fe40b18..496a048 100644 --- a/ui/dr-ui/src/library_ui/sync.rs +++ b/ui/dr-ui/src/library_ui/sync.rs @@ -112,23 +112,7 @@ pub(super) fn start_derived_sync(window: &AppWindow, ctl: &Rc let outbox = crate::export::outbox_dir(&conn.account); if crate::export::pending_count(&outbox) > 0 { let rx = crate::export::spawn_upload(conn.clone(), conn.account.root.clone(), outbox); - executors::spawn(Executor::Io, "upload-log", move || { - while let Ok(msg) = rx.recv() { - match msg { - crate::export::UploadMessage::Status(s) => log::info!("export: {s}"), - crate::export::UploadMessage::Finished { - uploaded, - remaining, - error, - } => { - log::info!("export: {uploaded} uploaded, {remaining} still queued"); - if let Some(e) = error { - log::warn!("export upload stopped: {e}"); - } - } - } - } - }); + crate::export::watch_upload(rx, window.as_weak()); } } diff --git a/ui/dr-ui/src/merge.rs b/ui/dr-ui/src/merge.rs index cae4bf8..53433dc 100644 --- a/ui/dr-ui/src/merge.rs +++ b/ui/dr-ui/src/merge.rs @@ -87,6 +87,11 @@ pub enum MergeDestination { outbox: PathBuf, /// The sources' folder, relative to the library root. remote_dir: String, + /// The names the library already holds in that folder, as the + /// catalog has them. The upload replaces whatever is at its name, so + /// a second merge of the same frames must not be called what the + /// first one was. + taken: std::collections::HashSet, }, } @@ -228,18 +233,41 @@ pub enum MergeEvent { /// for the photographer to confirm, and what names a failure. Aligned(AlignmentReport), /// The composite is written: on the device at `path`, or staged in the - /// outbox for the drain to upload, in which case `staged` is true and - /// the library learns of it when the folder is next scanned. + /// outbox for the drain to upload, in which case `staged` is true. + /// `composite` is what the library needs to show it before the upload + /// has finished and a scan has found it (FR-MRG-6). Done { path: PathBuf, staged: bool, width: u32, height: u32, + composite: Box, }, Failed(String), Cancelled, } +/// TRACES: FR-MRG-6 +/// A finished composite, described for the catalog: everything a scan and a +/// header read would have found, known here without either. +#[derive(Debug, Clone, PartialEq)] +pub struct Composite { + /// Its name where it goes — in the library folder for a staged merge. + pub name: String, + /// The picture it opens on: the crop where the border was cropped, the + /// whole composite where it was filled. + pub width: u32, + pub height: u32, + /// The middle of the sweep, as written into the DNG; `None` where no + /// frame carried a time. + pub captured_at: Option, + pub captured_offset: Option, + pub camera: Option, + pub lens: Option, + pub iso: Option, + pub file_size: u64, +} + impl MergeEvent { /// Whether this is the job's last word: after one of these the worker /// has nothing more to say, so its channel closing is expected. @@ -598,10 +626,14 @@ fn run_inner( // 5. Merge, into a DNG beside the first frame. let name = format!("{}-pano.dng", stem(Path::new(&first_name))); let (out_path, staged) = match &request.destination { - MergeDestination::Local(dir) => (unused_name(dir, &name), false), - MergeDestination::Outbox { outbox, remote_dir } => { + MergeDestination::Local(dir) => (unused_name(dir, &name, &Default::default()), false), + MergeDestination::Outbox { + outbox, + remote_dir, + taken, + } => { std::fs::create_dir_all(outbox).map_err(|e| format!("{}: {e}", outbox.display()))?; - let path = unused_name(outbox, &name); + let path = unused_name(outbox, &name, taken); // The record first here, unlike an export: the payload is // written over minutes and a record naming a half-written file // is worse than a payload with no record, so the record is @@ -649,6 +681,16 @@ fn run_inner( carried.captured_at = Some((sum / stamps.len() as i128) as i64); } + // What the catalog is told about the file (FR-MRG-6), taken before the + // header moves into the writer. + let described = ( + carried.captured_at, + carried.captured_offset, + crate::library::camera_label(carried.make.as_deref(), carried.model.as_deref()), + carried.lens.as_ref().map(|l| l.trim().to_string()), + carried.iso, + ); + let output = MergeOutput { projection, scale: focal_full, @@ -923,6 +965,16 @@ fn run_inner( cleanup(&out_path); return Err(format!("{}: {e}", out_path.display())); } + // The rectangle the file opens on, as the writer recorded it. + let picture = if fill_cam.is_some() { + None + } else { + inscribed + .lock() + .ok() + .map(|i| clamp_crop(i.best(), out_w, out_h)) + .filter(|r| r.width > 0 && r.height > 0) + }; log::info!( "merge: {}×{} written to {} in {:?}", out_w, @@ -971,14 +1023,43 @@ fn run_inner( } } + let (captured_at, captured_offset, camera, lens, iso) = described; + let composite = Composite { + name: out_path + .file_name() + .map(|n| n.to_string_lossy().into_owned()) + .unwrap_or_default(), + width: picture.map_or(out_w, |r| r.width), + height: picture.map_or(out_h, |r| r.height), + captured_at, + captured_offset, + camera, + lens, + iso, + file_size: std::fs::metadata(&out_path).map(|m| m.len()).unwrap_or(0), + }; + Ok(Some(MergeEvent::Done { path: out_path, staged, width: out_w, height: out_h, + composite: Box::new(composite), })) } +/// The crop as `write_linear_dng` writes it: inside the frame. +fn clamp_crop(r: dr_export::Rect, width: u32, height: u32) -> dr_export::Rect { + let x = r.x.min(width.saturating_sub(1)); + let y = r.y.min(height.saturating_sub(1)); + dr_export::Rect { + x, + y, + width: r.width.min(width - x), + height: r.height.min(height - y), + } +} + /// The proxy the detector reads: the frame through the camera-space tap at /// proxy size, upright, as gamma-encoded grey. /// @@ -1493,16 +1574,20 @@ fn part_name(path: &Path) -> PathBuf { } /// `name` in `dir`, numbered if that name is taken: a merge never -/// overwrites (FR-MRG-3). -fn unused_name(dir: &Path, name: &str) -> PathBuf { - let mut candidate = dir.join(name); +/// overwrites (FR-MRG-3). `name` in `dir`, or `name-2`, `name-3`… — the first that neither `dir` +/// nor `taken` (the names already at the destination) holds. +fn unused_name(dir: &Path, name: &str, taken: &std::collections::HashSet) -> PathBuf { + let mut candidate = name.to_string(); let base = stem(Path::new(name)); let mut n = 2; - while candidate.exists() { - candidate = dir.join(format!("{base}-{n}.dng")); + while dir.join(&candidate).exists() + || part_name(&dir.join(&candidate)).exists() + || taken.contains(&candidate) + { + candidate = format!("{base}-{n}.dng"); n += 1; } - candidate + dir.join(candidate) } /// Drain everything a job has said so far. @@ -1589,6 +1674,34 @@ mod tests { ); } + #[test] + fn a_second_merge_is_not_named_over_the_first_in_the_library() { + // The outbox is empty once the first has uploaded, so only the + // catalog knows the name is taken — and the upload replaces + // whatever is at its name. + let dir = std::env::temp_dir().join(format!("dr-merge-names-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + std::fs::create_dir_all(&dir).unwrap(); + let none = std::collections::HashSet::new(); + assert_eq!( + unused_name(&dir, "a-pano.dng", &none), + dir.join("a-pano.dng") + ); + let taken: std::collections::HashSet = + ["a-pano.dng".to_string(), "a-pano-2.dng".to_string()].into(); + assert_eq!( + unused_name(&dir, "a-pano.dng", &taken), + dir.join("a-pano-3.dng") + ); + // And one still being written in the outbox holds its name too. + std::fs::write(part_name(&dir.join("a-pano.dng")), b"").unwrap(); + assert_eq!( + unused_name(&dir, "a-pano.dng", &none), + dir.join("a-pano-2.dng") + ); + let _ = std::fs::remove_dir_all(&dir); + } + #[test] fn a_worker_that_hangs_up_mid_job_is_reported_as_gone() { // The failure the page could not see: a panic drops the sender with diff --git a/ui/dr-ui/src/merge_ui.rs b/ui/dr-ui/src/merge_ui.rs index a5957ff..ceba979 100644 --- a/ui/dr-ui/src/merge_ui.rs +++ b/ui/dr-ui/src/merge_ui.rs @@ -8,9 +8,10 @@ //! opens the page; a timer drains the job's events into the page's //! properties. When the alignment arrives the page shows it and waits. //! A frame's box leaves it out or brings it back, and the job aligns again -//! over the rest. "Merge" sends the decision; "Stop" or "Back" cancels. When the file is -//! staged, the outbox drains and the library rescans, and the composite -//! appears in the grid beside its sources. +//! over the rest. "Merge" sends the decision; "Stop" or "Back" cancels. When +//! the file is staged the composite is catalogued at once and appears in the +//! grid beside its sources; the outbox drains, and the scan that follows the +//! upload gives it the identity the server assigned (FR-MRG-6). use crate::executors::{self, Executor}; use std::cell::{Cell, RefCell}; @@ -36,6 +37,29 @@ struct Job { cancel: Cancel, names: Vec, activity: Activity, + /// Where in the library the composite goes, for a merge of library + /// frames; `None` for one written to a folder on the device. + library: Option, +} + +/// A library merge's destination, kept for the moment it finishes. +#[derive(Debug, Clone)] +struct Destined { + remote_dir: String, + sources: Vec, +} + +/// TRACES: FR-MRG-6 +/// A composite that has been staged for the library: where it goes, what it +/// was made from, and what the job knows about it — enough for the library +/// to catalogue it before the upload has finished. +#[derive(Debug, Clone)] +pub struct Placed { + /// The sources' folder, relative to the library root. + pub remote_dir: String, + /// The frames, as library paths. + pub sources: Vec, + pub composite: merge::Composite, } pub struct MergeController { @@ -74,15 +98,18 @@ impl MergeController { } /// What the page needs from the library to start: the account for the -/// fetch, and where the outbox is. +/// fetch, where the outbox is, and the names a library folder already +/// holds, so the composite is not named over one of them. pub struct Context { pub conn: dr_sync::Connection, pub outbox: std::path::PathBuf, + pub names_in: Box std::collections::HashSet>, } /// Wire the page. `sources` yields the selection as fetchable library /// sources; `context` the account; `on_done` runs when a composite has -/// been staged, so the caller can drain the outbox and rescan. +/// been written, with where it went in the library when it was staged for +/// one, so the caller can catalogue it and drain the outbox. pub fn wire( window: &AppWindow, ctl: Rc, @@ -93,7 +120,7 @@ pub fn wire( ) where S: Fn() -> Vec + 'static, C: Fn() -> Option + 'static, - F: Fn(&AppWindow) + 'static, + F: Fn(&AppWindow, Option<&Placed>) + 'static, { wire_start(window, &ctl, gpu, sources, context, on_done); wire_decision(window, &ctl); @@ -111,7 +138,7 @@ fn wire_start( ) where S: Fn() -> Vec + 'static, C: Fn() -> Option + 'static, - F: Fn(&AppWindow) + 'static, + F: Fn(&AppWindow, Option<&Placed>) + 'static, { let on_done = Rc::new(on_done); let gpu_for_start = gpu.clone(); @@ -165,10 +192,15 @@ fn wire_start( let remote_dir = first_dir .strip_prefix(&root) .map(|s| s.trim_start_matches('/').to_string()) - .unwrap_or(first_dir); + .unwrap_or_else(|| first_dir.clone()); let destination = MergeDestination::Outbox { outbox: context.outbox, + remote_dir: remote_dir.clone(), + taken: (context.names_in)(&first_dir), + }; + let destined = Destined { remote_dir, + sources: sources.iter().map(|(p, _)| p.clone()).collect(), }; let names: Vec = sources @@ -226,7 +258,16 @@ fn wire_start( } Some(frames) }; - start(&w, &ctl, gpu, names, destination, fetch, &on_done); + start( + &w, + &ctl, + gpu, + names, + destination, + Some(destined), + fetch, + &on_done, + ); }); } @@ -273,6 +314,7 @@ fn wire_start( gpu, names, MergeDestination::Local(dir), + None, fetch, &on_done, ); @@ -459,14 +501,18 @@ fn wire_stop_and_leave(window: &AppWindow, ctl: &Rc) { /// Start a job: `fetch` runs first on the job's thread and hands back the /// frames (or reports why not and returns `None`); the job follows on the /// same thread. The page opens clean, and a timer drains the events. +// Each argument is a different part of the job: where it runs, what it is +// called, where it goes, how its frames arrive, and who is told. +#[allow(clippy::too_many_arguments)] fn start( window: &AppWindow, ctl: &Rc, gpu: dr_gpu::GpuContext, names: Vec, destination: MergeDestination, + library: Option, fetch: Fetch, - on_done: &Rc, + on_done: &Rc) + 'static>, ) where Fetch: FnOnce(&Sender, &Cancel) -> Option> + Send + 'static, { @@ -495,6 +541,7 @@ fn start( cancel, names, activity, + library, }); *ctl.report.borrow_mut() = None; ctl.projection.set(0); @@ -547,7 +594,11 @@ fn chip_projection(i: i32) -> Option { } /// Take everything the job has said and reflect it on the page. -fn drain(window: &AppWindow, ctl: &Rc, on_done: &Rc) { +fn drain( + window: &AppWindow, + ctl: &Rc, + on_done: &Rc)>, +) { let (events, gone) = { let job = ctl.job.borrow(); let Some(job) = job.as_ref() else { return }; @@ -606,6 +657,7 @@ fn drain(window: &AppWindow, ctl: &Rc, on_done: &Rc { let name = path .file_name() @@ -616,7 +668,7 @@ fn drain(window: &AppWindow, ctl: &Rc, on_done: &Rc, on_done: &Rc { window.set_merge_running(false);